Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@ import { performAttach, type AttachFlowInput } from './structured-agent-session-
import { AgentSessionJournal } from '../agent-session-journal/journal-store'
import { agentSessionJournalCloseRetries } from '../agent-session-journal/journal-close-retry'
import * as legacyImport from '../agent-session-journal/journal-legacy-import'
import { restoreStructuredAgentSessionRead } from './structured-agent-session-read-restore'

const NOW = 1_800_000_000_000
const SESSION = 'codex_adopting_session'
Expand Down Expand Up @@ -210,7 +211,7 @@ describe('adopting a provider conversation on create', () => {
}
)

it('still releases acquisition and closes the provisional journal on an import write failure', async () => {
it('fails an import write before any spawn and closes the founding journal', async () => {
root = await mkdtemp(join(tmpdir(), 'orca-adopt-write-failure-'))
const transcriptPath = join(root, 'rollout.jsonl')
await writeCodexRollout(transcriptPath, 'valid source')
Expand All @@ -220,9 +221,29 @@ describe('adopting a provider conversation on create', () => {
const close = vi.spyOn(agentSessionJournalCloseRetries, 'closeOrRetain')
const sessionAdapter = adapter()
await expect(attach(transcriptPath, sessionAdapter)).rejects.toThrow('disk write failed')
expect(sessionAdapter.acquire).toHaveBeenCalledTimes(1)
expect(sessionAdapter.releaseAcquisition).toHaveBeenCalledTimes(1)
// The conversation is founded before a child exists, so nothing was spawned to release.
expect(sessionAdapter.acquire).not.toHaveBeenCalled()
expect(sessionAdapter.releaseAcquisition).not.toHaveBeenCalled()
expect(close).toHaveBeenCalledTimes(1)
expect(store?.getRecord(SESSION)?.lease.claimStatus).toBe('released')
})

it('keeps the adopted history readable when the first start fails', async () => {
root = await mkdtemp(join(tmpdir(), 'orca-adopt-failed-start-'))
const transcriptPath = join(root, 'rollout.jsonl')
await writeCodexRollout(transcriptPath, 'history before the failed start')
const sessionAdapter = adapter()
vi.mocked(sessionAdapter.acquire).mockRejectedValueOnce(new Error('codex exited (code 1)'))

await expect(attach(transcriptPath, sessionAdapter)).resolves.toMatchObject({ ok: false })

// A send restarts from the record alone, with no adopt source, so the history must already
// be in the conversation the create founded.
const restored = await restoreStructuredAgentSessionRead(store!, root, SESSION)
expect(JSON.stringify(restored?.journal.snapshot().items)).toContain(
'history before the failed start'
)
await restored?.journal.close()
})

it('prepares a valid source once before acquisition and imports those exact items', async () => {
Expand Down
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
import type { AgentSessionWireRefusal } from '../../../shared/agent-session-wire'
import type { AgentSessionRecord } from '../../../shared/agent-session-record'
import type { AgentSessionAttachParams, AttachedJournal } from './structured-agent-session-attach'
import { agentSessionJournalCloseRetries } from '../agent-session-journal/journal-close-retry'
import type { AgentSessionAttachParams } from './structured-agent-session-attach'
import type { AgentSessionJournal } from '../agent-session-journal/journal-store'
import type { JournalReplacementItem } from '../agent-session-journal/journal-epoch-replacement'
import {
importLegacyTranscriptIntoJournal,
Expand Down Expand Up @@ -58,39 +58,24 @@ async function readAdoptedTranscript(
// Import before publication so the first visible chat agrees with the provider's resumed context.
export async function importAdoptedTranscript(
params: AgentSessionAttachParams,
attached: AttachedJournal,
record: AgentSessionRecord,
prepared: JournalReplacementItem[] | null
): Promise<void> {
try {
await applyAdoptedTranscript(params, attached, record, prepared)
} catch (error) {
// Publication has not taken ownership of this provisional journal yet.
await agentSessionJournalCloseRetries.closeOrRetain(attached.journal)
throw error
}
}

async function applyAdoptedTranscript(
params: AgentSessionAttachParams,
attached: AttachedJournal,
journal: AgentSessionJournal,
record: AgentSessionRecord,
prepared: JournalReplacementItem[] | null
): Promise<void> {
const adopt = params.adopt
// A new journal contains only its epoch row; replay must preserve subsequent durable writes.
if (!adopt || attached.journal.cursor().sequence > 1) {
if (!adopt || journal.cursor().sequence > 1) {
return
}
if (prepared) {
await attached.journal.replaceEpochItems('legacy_import', record.lease.runtimeFence, prepared)
await journal.replaceEpochItems('legacy_import', record.lease.runtimeFence, prepared)
return
}
if (!adopt.transcriptPath) {
throw new Error('agent_session_identity_required')
}
const imported = await importLegacyTranscriptIntoJournal({
journal: attached.journal,
journal,
agent: params.agent,
sessionId:
adopt.providerHandle.kind === 'claude'
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -3,9 +3,10 @@ import {
failedAcquisitionRefusal,
failedAcquisitionSettlement
} from './structured-agent-session-failed-create-refusal'
import type {
StructuredAgentSessionAdapter,
StructuredAgentSessionProviderChildPhase
import {
AgentSessionPreSpawnError,
type StructuredAgentSessionAdapter,
type StructuredAgentSessionProviderChildPhase
} from './structured-agent-session-adapter'
// The host supplies owner authority; this flow reserves, proves, and publishes the session.

Expand All @@ -32,15 +33,14 @@ import type { StructuredAgentSessionEventSink } from './structured-agent-session
import { resolveAgentSessionReplayOutcome } from './structured-agent-session-replay-outcome'
import { readAgentSessionHydrationPage } from './agent-session-history-page'
import { acquireOwner } from './structured-agent-session-acquisition'
import {
importAdoptedTranscript,
prepareAdoptedTranscript
} from './structured-agent-session-adopted-import'
import { prepareAdoptedTranscript } from './structured-agent-session-adopted-import'
import { foundAgentSessionConversation } from './structured-agent-session-conversation-founding'
import {
withAgentSessionCreatePhase,
type AgentSessionCreatePhaseRecorder
} from '../../observability/agent-session-instrumentation'
import type { ProviderHistoryWindow } from '../agent-session-journal/journal-submission-reconciler'
import type { JournalReplacementItem } from '../agent-session-journal/journal-epoch-replacement'

export type AttachFlowInput = {
store: AgentSessionRecordStore
Expand Down Expand Up @@ -159,6 +159,7 @@ export async function performAttach(
ownerAlreadyAdmitted: agentSessionLeaseAdmitsWriter(record.lease)
})
if (!agentSessionLeaseAdmitsWriter(record.lease)) {
await foundConversationBeforeSpawn(input, record, preparedTranscript.items)
const acquired = await withAgentSessionCreatePhase('acquire_owner', input.recordPhase, () =>
acquireOwner(input, record)
)
Expand Down Expand Up @@ -210,7 +211,6 @@ export async function performAttach(
adapter: input.adapter,
providerHistoryWindow
})
await importAdoptedTranscript(params, attached, record, preparedTranscript.items)
await input.onAttached(attached, acquisitionGeneration, acquiredOwner, providerChildPhase)
await store.recordOperationOutcome({
callerKey: input.callerKey,
Expand All @@ -237,6 +237,24 @@ export async function performAttach(
}
}

/** Nothing has spawned yet, so a failure here settles the reservation as processless. */
async function foundConversationBeforeSpawn(
input: AttachFlowInput,
record: AgentSessionRecord,
adoptedItems: JournalReplacementItem[] | null
): Promise<void> {
try {
await foundAgentSessionConversation({
record,
params: input.params,
journalRoot: input.journalRoot,
adoptedItems
})
} catch (error) {
throw new AgentSessionPreSpawnError(error)
}
}

async function readProviderHistoryWindow(input: {
adapter: StructuredAgentSessionAdapter
identity: AgentSessionJournalIdentity
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,44 @@
// A conversation exists from its reservation, not from its first live provider child.
//
// The journal used to be opened only once acquisition succeeded, so a create whose child died at
// startup left a record and nothing to read: no chat to publish, and nothing a later send could be
// admitted into. Founding it here makes that create an ordinary readable session with a released
// lease, which a send restarts like any other.

import { existsSync } from 'node:fs'
import type { AgentSessionRecord } from '../../../shared/agent-session-record'
import { agentSessionJournalCloseRetries } from '../agent-session-journal/journal-close-retry'
import type { JournalReplacementItem } from '../agent-session-journal/journal-epoch-replacement'
import { journalDatabaseFile, journalDirectoryFor } from '../agent-session-journal/journal-paths'
import { openAgentSessionJournal } from '../agent-session-journal/journal-store-factory'
import { importAdoptedTranscript } from './structured-agent-session-adopted-import'
import {
journalIdentityFor,
type AgentSessionAttachParams
} from './structured-agent-session-attach'

export async function foundAgentSessionConversation(input: {
record: AgentSessionRecord
params: AgentSessionAttachParams
journalRoot: string
adoptedItems: JournalReplacementItem[] | null
}): Promise<void> {
const journalDir = journalDirectoryFor(input.journalRoot, {
workspaceId: input.params.location.workspaceId,
sessionId: input.record.sessionId
})
// An adopted conversation whose import never landed is still owed it; the import skips a
// journal that already holds more than its epoch.
if (existsSync(journalDatabaseFile(journalDir)) && !input.params.adopt) {
return
}
const journal = await openAgentSessionJournal({
identity: journalIdentityFor(input.record, input.params),
journalDir
})
try {
await importAdoptedTranscript(input.params, journal, input.record, input.adoptedItems)
} finally {
await agentSessionJournalCloseRetries.closeOrRetain(journal)
}
}
Loading
Loading