mirror of
https://github.com/stablyai/orca.git
synced 2026-09-29 16:02:50 +00:00
fix: validate history before claiming adopted sessions
This commit is contained in:
@@ -90,6 +90,24 @@ export async function importLegacyTranscriptIntoJournal(input: {
|
||||
fence: number
|
||||
options?: LegacyImportOptions
|
||||
}): Promise<LegacyImportResult> {
|
||||
const prepared = await prepareLegacyTranscriptImport(input)
|
||||
if (!prepared.ok) {
|
||||
return prepared
|
||||
}
|
||||
// An empty import must preserve any existing repair anchor and disclosure.
|
||||
if (prepared.items.length === 0) {
|
||||
const current = input.journal.cursor()
|
||||
return { ok: true, epoch: current.epoch, cursor: current, imported: 0, replaced: false }
|
||||
}
|
||||
const cursor = await input.journal.replaceEpochItems('legacy_import', input.fence, prepared.items)
|
||||
return { ok: true, epoch: cursor.epoch, cursor, imported: prepared.items.length, replaced: true }
|
||||
}
|
||||
|
||||
export async function prepareLegacyTranscriptImport(input: {
|
||||
agent: AgentType
|
||||
sessionId: string
|
||||
options?: LegacyImportOptions
|
||||
}): Promise<{ ok: true; items: JournalReplacementItem[] } | { ok: false; error: string }> {
|
||||
const options = input.options ?? {}
|
||||
const limits = options.limits ?? DEFAULT_JOURNAL_PAYLOAD_LIMITS
|
||||
const transcriptAgent = resolveNativeChatTranscriptAgent(input.agent)
|
||||
@@ -145,22 +163,7 @@ export async function importLegacyTranscriptIntoJournal(input: {
|
||||
observedAt: message.timestamp ?? undefined
|
||||
})
|
||||
}
|
||||
// A transcript that decodes to nothing reconstructs nothing, and an empty
|
||||
// replacement is not a harmless no-op: it would delete the repair's anchor and
|
||||
// its disclosure, leaving nothing to ask for the history again. The epoch
|
||||
// stands so a later read can still rebuild it.
|
||||
if (replacement.length === 0) {
|
||||
const current = input.journal.cursor()
|
||||
return { ok: true, epoch: current.epoch, cursor: current, imported: 0, replaced: false }
|
||||
}
|
||||
const cursor = await input.journal.replaceEpochItems('legacy_import', input.fence, replacement)
|
||||
return {
|
||||
ok: true,
|
||||
epoch: cursor.epoch,
|
||||
cursor,
|
||||
imported: decoded.messages.length,
|
||||
replaced: true
|
||||
}
|
||||
return { ok: true, items: replacement }
|
||||
}
|
||||
|
||||
const TRANSCRIPT_DECODERS = {
|
||||
|
||||
+67
-48
@@ -1,12 +1,6 @@
|
||||
// Adopting a history row imports its transcript into the new session's journal.
|
||||
//
|
||||
// The import runs AFTER the provider is acquired, because for a create there is no earlier moment
|
||||
// — the journal does not exist until the child does. That ordering is what makes the failure cases
|
||||
// here load-bearing: by then the provider has already resumed and holds the conversation in
|
||||
// context, so an attach that succeeds with an empty journal would show the user a blank chat beside
|
||||
// an agent that can already answer from history.
|
||||
// Source validation must finish before a new session claims the provider conversation.
|
||||
|
||||
import { mkdtemp, rm, writeFile } from 'node:fs/promises'
|
||||
import { mkdtemp, rm, writeFile, truncate } from 'node:fs/promises'
|
||||
import { tmpdir } from 'node:os'
|
||||
import { join } from 'node:path'
|
||||
import { afterEach, describe, expect, it, vi } from 'vitest'
|
||||
@@ -18,7 +12,9 @@ import {
|
||||
type AgentSessionAttachParams
|
||||
} from './structured-agent-session-attach'
|
||||
import { performAttach, type AttachFlowInput } from './structured-agent-session-attach-flow'
|
||||
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'
|
||||
|
||||
const NOW = 1_800_000_000_000
|
||||
const SESSION = 'codex_adopting_session'
|
||||
@@ -181,51 +177,74 @@ describe('adopting a provider conversation on create', () => {
|
||||
expect(sessionAdapter.acquire).toHaveBeenCalledTimes(1)
|
||||
})
|
||||
|
||||
it('fails the attach when the adopted transcript cannot be read', async () => {
|
||||
root = await mkdtemp(join(tmpdir(), 'orca-adopt-missing-'))
|
||||
const sessionAdapter = adapter()
|
||||
it.each(['missing', 'oversized', 'empty', 'invalid', 'source-less'] as const)(
|
||||
'refuses %s source before claiming a conversation',
|
||||
async (kind) => {
|
||||
root = await mkdtemp(join(tmpdir(), 'orca-adopt-preflight-'))
|
||||
const transcriptPath = join(root, 'rollout.jsonl')
|
||||
if (kind === 'oversized') {
|
||||
await writeCodexRollout(transcriptPath, 'original turn')
|
||||
await truncate(transcriptPath, 16 * 1024 * 1024 + 1)
|
||||
} else if (kind === 'empty' || kind === 'invalid') {
|
||||
await writeFile(transcriptPath, kind === 'empty' ? '' : 'not json\n')
|
||||
}
|
||||
const sessionAdapter = adapter()
|
||||
const onAttached = vi.fn()
|
||||
const result = await attach(
|
||||
kind === 'source-less' ? undefined : transcriptPath,
|
||||
sessionAdapter,
|
||||
onAttached
|
||||
)
|
||||
expect(result).toMatchObject({
|
||||
ok: false,
|
||||
refusal: { code: 'agent_session_identity_required' }
|
||||
})
|
||||
expect(sessionAdapter.acquire).not.toHaveBeenCalled()
|
||||
expect(sessionAdapter.releaseAcquisition).not.toHaveBeenCalled()
|
||||
expect(onAttached).not.toHaveBeenCalled()
|
||||
expect(store?.getRecord(SESSION)).toBeNull()
|
||||
expect(store?.listOperationRows()).toEqual([])
|
||||
if (kind === 'oversized') {
|
||||
expect(JSON.stringify(result)).toContain('import bound')
|
||||
}
|
||||
}
|
||||
)
|
||||
|
||||
it('still releases acquisition and closes the provisional journal on an import write failure', async () => {
|
||||
root = await mkdtemp(join(tmpdir(), 'orca-adopt-write-failure-'))
|
||||
const transcriptPath = join(root, 'rollout.jsonl')
|
||||
await writeCodexRollout(transcriptPath, 'valid source')
|
||||
vi.spyOn(AgentSessionJournal.prototype, 'replaceEpochItems').mockRejectedValueOnce(
|
||||
new Error('disk write failed')
|
||||
)
|
||||
const close = vi.spyOn(agentSessionJournalCloseRetries, 'closeOrRetain')
|
||||
|
||||
// A post-acquisition failure throws rather than answering a refusal — the same path a journal
|
||||
// failure already takes — so the caller learns the outcome is unknown, not that nothing ran.
|
||||
await expect(attach(join(root, 'does-not-exist.jsonl'), sessionAdapter)).rejects.toThrow(
|
||||
/ENOENT|no such file/
|
||||
)
|
||||
// The provider had already resumed, so its acquisition is released rather than left holding a
|
||||
// conversation no surface will ever show.
|
||||
expect(sessionAdapter.releaseAcquisition).toHaveBeenCalled()
|
||||
const sessionAdapter = adapter()
|
||||
await expect(attach(transcriptPath, sessionAdapter)).rejects.toThrow('disk write failed')
|
||||
expect(sessionAdapter.acquire).toHaveBeenCalledTimes(1)
|
||||
expect(sessionAdapter.releaseAcquisition).toHaveBeenCalledTimes(1)
|
||||
expect(close).toHaveBeenCalledTimes(1)
|
||||
close.mockRestore()
|
||||
})
|
||||
|
||||
it('refuses a source-less replay identity when the journal was never imported', async () => {
|
||||
root = await mkdtemp(join(tmpdir(), 'orca-adopt-unimported-replay-'))
|
||||
it('prepares a valid source once before acquisition and imports those exact items', async () => {
|
||||
root = await mkdtemp(join(tmpdir(), 'orca-adopt-once-'))
|
||||
const transcriptPath = join(root, 'rollout.jsonl')
|
||||
await writeCodexRollout(transcriptPath, 'prepared before acquiring')
|
||||
const prepare = vi.spyOn(legacyImport, 'prepareLegacyTranscriptImport')
|
||||
const sessionAdapter = adapter()
|
||||
|
||||
await expect(attach(undefined, sessionAdapter)).rejects.toThrow(
|
||||
'agent_session_identity_required'
|
||||
const acquire = sessionAdapter.acquire
|
||||
sessionAdapter.acquire = vi.fn(async (input) => {
|
||||
expect(prepare).toHaveBeenCalledTimes(1)
|
||||
await rm(transcriptPath)
|
||||
return acquire(input)
|
||||
})
|
||||
const result = await attach(transcriptPath, sessionAdapter, async ({ journal }) =>
|
||||
journal.close()
|
||||
)
|
||||
expect(sessionAdapter.releaseAcquisition).toHaveBeenCalled()
|
||||
})
|
||||
|
||||
it('fails the attach when the adopted transcript decodes to no messages', async () => {
|
||||
root = await mkdtemp(join(tmpdir(), 'orca-adopt-empty-'))
|
||||
const transcriptPath = join(root, 'empty.jsonl')
|
||||
// Well-formed but conversation-free: the row promised turns and the provider resumed them, so
|
||||
// an empty journal here is a disagreement, not an empty chat.
|
||||
await writeFile(
|
||||
transcriptPath,
|
||||
`${JSON.stringify({
|
||||
type: 'session_meta',
|
||||
payload: { id: THREAD, timestamp: '2026-09-06T18:00:00.000Z', cwd: '/workspace' }
|
||||
})}\n`,
|
||||
'utf8'
|
||||
)
|
||||
const sessionAdapter = adapter()
|
||||
|
||||
await expect(attach(transcriptPath, sessionAdapter)).rejects.toThrow(
|
||||
'agent_session_identity_required'
|
||||
)
|
||||
expect(sessionAdapter.releaseAcquisition).toHaveBeenCalled()
|
||||
expect(result.ok).toBe(true)
|
||||
if (!result.ok) {
|
||||
throw new Error('attach failed')
|
||||
}
|
||||
expect(JSON.stringify(result.value.page.items)).toContain('prepared before acquiring')
|
||||
expect(prepare).toHaveBeenCalledTimes(1)
|
||||
})
|
||||
})
|
||||
|
||||
@@ -1,18 +1,91 @@
|
||||
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 { importLegacyTranscriptIntoJournal } from '../agent-session-journal/journal-legacy-import'
|
||||
import { agentSessionJournalCloseRetries } from '../agent-session-journal/journal-close-retry'
|
||||
import type { JournalReplacementItem } from '../agent-session-journal/journal-epoch-replacement'
|
||||
import {
|
||||
importLegacyTranscriptIntoJournal,
|
||||
prepareLegacyTranscriptImport
|
||||
} from '../agent-session-journal/journal-legacy-import'
|
||||
|
||||
export async function prepareAdoptedTranscript(
|
||||
params: AgentSessionAttachParams
|
||||
): Promise<
|
||||
| { ok: true; items: JournalReplacementItem[] | null }
|
||||
| { ok: false; refusal: AgentSessionWireRefusal }
|
||||
> {
|
||||
try {
|
||||
return { ok: true, items: await readAdoptedTranscript(params) }
|
||||
} catch (error) {
|
||||
return {
|
||||
ok: false,
|
||||
refusal: {
|
||||
code: 'agent_session_identity_required',
|
||||
message: error instanceof Error ? error.message : String(error)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Validate source input before a new record can claim the provider conversation.
|
||||
async function readAdoptedTranscript(
|
||||
params: AgentSessionAttachParams
|
||||
): Promise<JournalReplacementItem[] | null> {
|
||||
const adopt = params.adopt
|
||||
if (!adopt) {
|
||||
return null
|
||||
}
|
||||
if (!adopt.transcriptPath) {
|
||||
throw new Error('agent_session_identity_required')
|
||||
}
|
||||
const prepared = await prepareLegacyTranscriptImport({
|
||||
agent: params.agent,
|
||||
sessionId:
|
||||
adopt.providerHandle.kind === 'claude'
|
||||
? adopt.providerHandle.sessionId
|
||||
: adopt.providerHandle.threadId,
|
||||
options: { filePath: adopt.transcriptPath }
|
||||
})
|
||||
if (!prepared.ok) {
|
||||
throw new Error(prepared.error)
|
||||
}
|
||||
if (prepared.items.length === 0) {
|
||||
throw new Error('agent_session_identity_required')
|
||||
}
|
||||
return prepared.items
|
||||
}
|
||||
|
||||
// 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
|
||||
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,
|
||||
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) {
|
||||
return
|
||||
}
|
||||
if (prepared) {
|
||||
await attached.journal.replaceEpochItems('legacy_import', record.lease.runtimeFence, prepared)
|
||||
return
|
||||
}
|
||||
if (!adopt.transcriptPath) {
|
||||
throw new Error('agent_session_identity_required')
|
||||
}
|
||||
|
||||
@@ -36,8 +36,10 @@ import type { StructuredAgentSessionEventSink } from './structured-agent-session
|
||||
import { readNativeSessionOptions } from './structured-agent-session-option-restoration'
|
||||
import { resolveAgentSessionReplayOutcome } from './structured-agent-session-replay-outcome'
|
||||
import { readAgentSessionHydrationPage } from './agent-session-history-page'
|
||||
import { agentSessionJournalCloseRetries } from '../agent-session-journal/journal-close-retry'
|
||||
import { importAdoptedTranscript } from './structured-agent-session-adopted-import'
|
||||
import {
|
||||
importAdoptedTranscript,
|
||||
prepareAdoptedTranscript
|
||||
} from './structured-agent-session-adopted-import'
|
||||
|
||||
export type AttachFlowInput = {
|
||||
store: AgentSessionRecordStore
|
||||
@@ -80,6 +82,12 @@ export async function performAttach(
|
||||
let acquisitionGeneration: string | null = null
|
||||
let reservedRecord: AgentSessionRecord | null = null
|
||||
let replayed = false
|
||||
const preparedTranscript = store.getRecord(sessionId)
|
||||
? { ok: true as const, items: null }
|
||||
: await prepareAdoptedTranscript(params)
|
||||
if (!preparedTranscript.ok) {
|
||||
return preparedTranscript
|
||||
}
|
||||
try {
|
||||
const reserved = await store.reserveOwner(
|
||||
reserveRequestFor({
|
||||
@@ -183,13 +191,7 @@ export async function performAttach(
|
||||
journalRoot: input.journalRoot,
|
||||
adapter: input.adapter
|
||||
})
|
||||
try {
|
||||
await importAdoptedTranscript(params, attached, record)
|
||||
} catch (error) {
|
||||
// Publication has not taken ownership of this provisional journal yet.
|
||||
await agentSessionJournalCloseRetries.closeOrRetain(attached.journal)
|
||||
throw error
|
||||
}
|
||||
await importAdoptedTranscript(params, attached, record, preparedTranscript.items)
|
||||
await input.onAttached(attached, acquisitionGeneration)
|
||||
await store.recordOperationOutcome({
|
||||
callerKey: input.callerKey,
|
||||
|
||||
Reference in New Issue
Block a user