From c0b07d3a71d3c2ec9cae7fca1b933da1d206bc05 Mon Sep 17 00:00:00 2001 From: Brennan Benson <79079362+brennanb2025@users.noreply.github.com> Date: Mon, 5 Oct 2026 22:20:59 -0700 Subject: [PATCH] Admit one provider event's journal writes as one queued operation, decided when it runs (#25141) * Move the turn message ordinals and the turn-row revision to the neutral timeline folder Pure moves so a shared timeline assembler can use them: Codex's message ordinal counter becomes ProviderTurnMessageOrdinals and Claude's turn-row revision becomes the provider-neutral agent-journal turn-row revision. Only names and import paths change. * Admit one provider event's writes as one transition, and let rows be found again after a restart - A sink transition is admitted whole or not at all; its steps run back to back at their turn in the journal's write queue, and each resolver reads the fold with every earlier write landed. A resolver may also say where the row belongs (turn scope, provider reference), and the writer always hears how the transition landed. A resolved lifecycle batch chooses its settlement mutations from the fold at execution. - New optional row field providerItemRef: the provider's own reference for the item a row is, written only where the row's identity cannot spell it (Codex keys messages by their place in the turn and renumbers its item ids on resume). Set by the creating write, kept by revisions, indexed by the journal fold, never read by clients. A downgrade test shows an older host and client render such rows unchanged. - Provider timeline identity schemes (shared legacy arm, Codex) and the join index that resolves a provider item to its row from memory or the fold: ordinals and request incarnations are read back from the rows, so a restart or an evicted entry finds the original row instead of placing a new one. * Recover message ordinals from the journal's highest place, and forget joins read from a replaced epoch A fresh join index continued a turn's messages at the first free place, so a journal holding only a later ordinal (an imported or removed earlier row) had its sequence back-filled. The place is now one past the highest ordinal any row or echoed send holds there, read through a pure scheme reader. The join caches also drop what they read when the journal's epoch is replaced. * Spell the subagent thread's message slot without spreading an identity union * Keep the journal store under its line limit after the main merge * Drop the provider item reference, join index and identity schemes from the transition PR Nothing in production reaches the state they defended (an assembler that lost its memory while its child keeps streaming the same turn), and the stored Codex id was positional. The legacy identity scheme moves to the assembler PR with its first caller; the Codex scheme and any persisted reference wait for Codex to move onto the assembler. The Codex ordinal counter goes back to codex/, since no neutral code imports it. * Write a resolved settlement in one transaction through enqueueRows A settlement too large for one row now commits all its rows or none, through the journal's existing all-or-nothing write, instead of a row-by-row writer. Every row is built before any commits, so a settlement naming one item twice is refused before anything is written. * Drop the transition's landing report; keep the turn-row write fire-and-forget Nothing reads which steps of a transition wrote. A failed step fails the sink, leaving the steps before it written; the header says so, and tests cover it plus a settlement whose second row fails inside the transaction. writeAgentJournalTurnRow returns nothing again, as on main. * Run a transition's steps as a prefix; drop the paced flag and resolved options A failed step no longer lets the steps after it write: each step checks, at its own turn in the journal's queue, whether the write handed over just ahead of it completed, using the queue's count of completed write bodies (a promise would report the failure only after the next step ran). The sink fails only once every step has had its turn. The item step's `paced` size bypass and the resolver's replacement `options` are removed; nothing planned uses them. * Run a transition's steps in one queued write that loops over them The steps of one event now share one turn in the journal's write queue: a loop writes each in its own transaction through the row writer's synchronous writeRows (split out of enqueueRows) and stops at the first throw. Prefix semantics and "nothing lands between the steps" now hold by construction, so the completed-write counter on the queue, the step gate and the allSettled barrier are gone; the queue is back to main's bytes. * Use current provider handles in transition tests --- src/main/claude/claude-context-facts.ts | 11 +- src/main/claude/claude-open-turn.ts | 4 +- .../journal-lifecycle-batch-appender.ts | 41 +- .../journal-row-writer.ts | 55 +- .../journal-step-writer.ts | 58 ++ .../journal-store-collaborators.ts | 17 +- .../journal-store-contracts.ts | 8 + .../agent-session-journal/journal-store.ts | 20 +- .../agent-journal-turn-row-revision.ts} | 97 ++-- .../agent-session-delta-coalescer.ts | 25 + ...ructured-agent-session-event-sink-queue.ts | 3 + .../structured-agent-session-event-sink.ts | 7 +- ...tructured-agent-session-transition.test.ts | 520 ++++++++++++++++++ .../structured-agent-session-transition.ts | 126 +++++ 14 files changed, 909 insertions(+), 83 deletions(-) create mode 100644 src/main/native-chat/agent-session-journal/journal-step-writer.ts rename src/main/{claude/claude-turn-row-revision.ts => native-chat/agent-session-timeline/agent-journal-turn-row-revision.ts} (68%) create mode 100644 src/main/native-chat/agent-session-wire/structured-agent-session-transition.test.ts create mode 100644 src/main/native-chat/agent-session-wire/structured-agent-session-transition.ts diff --git a/src/main/claude/claude-context-facts.ts b/src/main/claude/claude-context-facts.ts index 358fd134c87..3a1326e0936 100644 --- a/src/main/claude/claude-context-facts.ts +++ b/src/main/claude/claude-context-facts.ts @@ -25,7 +25,10 @@ import type { ClaudeOpenTurn } from './claude-open-turn' import { claudeRecord, claudeText } from './claude-structured-item-translation' import type { ClaudeTurnEnd } from './claude-turn-lifecycle-item' import { isRootClaudeFrame } from './claude-turn-opening' -import { writeClaudeTurnRow, type ClaudeTurnRowTarget } from './claude-turn-row-revision' +import { + writeAgentJournalTurnRow, + type AgentJournalTurnRowTarget +} from '../native-chat/agent-session-timeline/agent-journal-turn-row-revision' const CONVERSATION_FRAME_TYPES = new Set(['assistant', 'user', 'stream_event']) @@ -188,7 +191,7 @@ export class ClaudeContextFacts { this.windowHint = report.windowTokens const contextUsage: AgentSessionContextUsage = part === 'report' ? { used: { kind: 'report', ...report }, window } : { window } - writeClaudeTurnRow( + writeAgentJournalTurnRow( this.sink, target === null ? { newest: true } : { identity: target }, { contextUsage }, @@ -211,9 +214,9 @@ export class ClaudeContextFacts { windowIfNoneHeld?: AgentSessionContextWindow ): void { const identity = this.turn.identity - const target: ClaudeTurnRowTarget = identity ? { identity } : { newest: true } + const target: AgentJournalTurnRowTarget = identity ? { identity } : { newest: true } // A context fact often lands with no later frame to publish it, so it publishes itself. - writeClaudeTurnRow( + writeAgentJournalTurnRow( this.sink, target, { contextUsage, ...(windowIfNoneHeld ? { windowIfNoneHeld } : {}) }, diff --git a/src/main/claude/claude-open-turn.ts b/src/main/claude/claude-open-turn.ts index 1cf0a8d048c..32149bf21f3 100644 --- a/src/main/claude/claude-open-turn.ts +++ b/src/main/claude/claude-open-turn.ts @@ -18,7 +18,7 @@ import { type ClaudeCurrentTurn, type ClaudeTurnEnd } from './claude-turn-lifecycle-item' -import { writeClaudeTurnRow } from './claude-turn-row-revision' +import { writeAgentJournalTurnRow } from '../native-chat/agent-session-timeline/agent-journal-turn-row-revision' import type { ClaudeCommandTurn } from './claude-command-turn' import { createClaudeTurnOpener, type ClaudeTurnSource } from './claude-turn-opening' @@ -198,7 +198,7 @@ export class ClaudeOpenTurn { contextUsage?: AgentSessionContextUsage ): void { const item = claudeTurnLifecycleItem(turn, end) - writeClaudeTurnRow( + writeAgentJournalTurnRow( this.deps.sink, { identity: item.identity }, { lifecycle: item.body, ...(contextUsage ? { contextUsage } : {}) }, diff --git a/src/main/native-chat/agent-session-journal/journal-lifecycle-batch-appender.ts b/src/main/native-chat/agent-session-journal/journal-lifecycle-batch-appender.ts index 68ba188e402..3cac8d03c21 100644 --- a/src/main/native-chat/agent-session-journal/journal-lifecycle-batch-appender.ts +++ b/src/main/native-chat/agent-session-journal/journal-lifecycle-batch-appender.ts @@ -1,7 +1,15 @@ +import { agentJournalItemKey } from '../../../shared/agent-session-journal-item-key' import type { AgentJournalCursor } from '../../../shared/agent-session-journal-types' import type { JournalReducerState } from './journal-reducer' -import { journalLifecycleBatchRowBuilder } from './journal-row-builders' -import type { JournalLifecycleBatchInput } from './journal-store-contracts' +import { partitionJournalLifecycleMutations } from './journal-lifecycle-batch-partition' +import { + journalLifecycleBatchRowBuilder, + type JournalLifecycleMutationInput +} from './journal-row-builders' +import type { + JournalLifecycleBatchInput, + JournalResolvedLifecycleBatchInput +} from './journal-store-contracts' import type { JournalRow } from './journal-row-schema' import { journalQueuedRejectionRowBuilders } from './journal-pending-submission-recovery' @@ -70,6 +78,35 @@ export class JournalLifecycleBatchAppender { }) } + /** The rows a resolved settlement writes, planned at its own turn in the queue: its mutations, + * chosen then, in as many consecutive rows as they need, minus any already applied. */ + planResolved( + input: JournalResolvedLifecycleBatchInput + ): ((seq: number, ts: number) => JournalRow)[] { + const mutations = input.resolve() + // Every chunk is built before any commits, so a second chunk naming the same item would + // reuse the first chunk's revision. + this.assertDistinctItems(mutations) + return partitionJournalLifecycleMutations(input.settlementId, mutations) + .filter((chunk) => !this.wasApplied(chunk.settlementId)) + .map((chunk) => + journalLifecycleBatchRowBuilder(this.deps.state, chunk.settlementId, chunk.mutations, input) + ) + } + + private assertDistinctItems(mutations: readonly JournalLifecycleMutationInput[]): void { + const { aliases } = this.deps.state() + const seen = new Set() + for (const mutation of mutations) { + const itemId = agentJournalItemKey(mutation.identity) + const resolved = aliases.get(itemId) ?? itemId + if (seen.has(resolved)) { + throw new Error('journal_resolved_lifecycle_batch_names_item_twice') + } + seen.add(resolved) + } + } + private wasApplied(settlementId: string): boolean { return this.deps.state().appliedSettlementIds.has(settlementId) } diff --git a/src/main/native-chat/agent-session-journal/journal-row-writer.ts b/src/main/native-chat/agent-session-journal/journal-row-writer.ts index dde1a66741d..60db6e0cbd1 100644 --- a/src/main/native-chat/agent-session-journal/journal-row-writer.ts +++ b/src/main/native-chat/agent-session-journal/journal-row-writer.ts @@ -76,33 +76,36 @@ export class JournalRowWriter { enqueueRows( plan: () => readonly ((seq: number, ts: number) => JournalRow)[] ): Promise { - return this.deps.serialize(() => { - assertJournalWritable(this.deps.readOnly(), this.deps.sessionId) - const first = this.deps.nextSequence() - const ts = this.deps.now() - const rows = plan().map((build, index) => build(first + index, ts)) - if (rows.length === 0) { - return rows - } - for (const row of rows) { - assertJournalFence(row.fence, this.deps.highestFence()) - } - try { - this.deps.database().transaction((db) => { - for (const row of rows) { - insertJournalRow(db, this.deps.sessionId, row) - this.runBookkeeping(db, row) - } - }) - } catch (error) { - this.deps.rolledBack?.() - throw error - } - for (const row of rows) { - this.deps.commit(row) - } + return this.deps.serialize(() => this.writeRows(plan)) + } + + /** `enqueueRows`' write, for a caller already running at its own turn in the queue. */ + writeRows(plan: () => readonly ((seq: number, ts: number) => JournalRow)[]): JournalRow[] { + assertJournalWritable(this.deps.readOnly(), this.deps.sessionId) + const first = this.deps.nextSequence() + const ts = this.deps.now() + const rows = plan().map((build, index) => build(first + index, ts)) + if (rows.length === 0) { return rows - }) + } + for (const row of rows) { + assertJournalFence(row.fence, this.deps.highestFence()) + } + try { + this.deps.database().transaction((db) => { + for (const row of rows) { + insertJournalRow(db, this.deps.sessionId, row) + this.runBookkeeping(db, row) + } + }) + } catch (error) { + this.deps.rolledBack?.() + throw error + } + for (const row of rows) { + this.deps.commit(row) + } + return rows } /** Assign the next sequence, make the row durable, and fold it through the SAME reducer diff --git a/src/main/native-chat/agent-session-journal/journal-step-writer.ts b/src/main/native-chat/agent-session-journal/journal-step-writer.ts new file mode 100644 index 00000000000..c0b62ce50c4 --- /dev/null +++ b/src/main/native-chat/agent-session-journal/journal-step-writer.ts @@ -0,0 +1,58 @@ +// Several writes run as ONE turn in the chat's write queue, one after another. +// +// Each step is planned from the fold with every step before it landed, and commits in its own +// transaction. The steps share one queue body, so no other write lands between them, and the +// first step that throws ends the run: the steps before it stay written and none after it runs. + +import type { + JournalItemAppendOptions, + JournalResolvedLifecycleBatchInput +} from './journal-store-contracts' +import type { JournalResolvedItem } from './journal-item-appender' +import type { JournalLifecycleBatchAppender } from './journal-lifecycle-batch-appender' +import type { JournalReducerState } from './journal-reducer' +import { journalItemRowBuilder } from './journal-row-builders' +import type { JournalRow } from './journal-row-schema' +import type { JournalRowWriter } from './journal-row-writer' +import type { JournalWriteBody } from './journal-write-queue' + +export type JournalStep = + | { + kind: 'item' + /** Read at the step's turn; null writes nothing. */ + resolve: () => JournalResolvedItem | null + options: JournalItemAppendOptions + } + | { kind: 'settlement'; batch: JournalResolvedLifecycleBatchInput } + +export class JournalStepWriter { + constructor( + private readonly deps: { + serialize: (run: JournalWriteBody) => Promise + state: () => JournalReducerState + writeRows: JournalRowWriter['writeRows'] + planSettlement: JournalLifecycleBatchAppender['planResolved'] + } + ) {} + + /** Whether each step wrote; rejects with the error of the step that threw. */ + append(steps: readonly JournalStep[]): Promise { + return this.deps.serialize(() => { + const wrote: boolean[] = [] + for (const step of steps) { + wrote.push(this.deps.writeRows(() => this.plan(step)).length > 0) + } + return wrote + }) + } + + private plan(step: JournalStep): ((seq: number, ts: number) => JournalRow)[] { + if (step.kind === 'settlement') { + return this.deps.planSettlement(step.batch) + } + const resolved = step.resolve() + return resolved === null + ? [] + : [journalItemRowBuilder(this.deps.state, resolved.identity, resolved.body, step.options)] + } +} diff --git a/src/main/native-chat/agent-session-journal/journal-store-collaborators.ts b/src/main/native-chat/agent-session-journal/journal-store-collaborators.ts index dc270bc97a6..1995b2e9875 100644 --- a/src/main/native-chat/agent-session-journal/journal-store-collaborators.ts +++ b/src/main/native-chat/agent-session-journal/journal-store-collaborators.ts @@ -17,6 +17,7 @@ import { JournalStopMarks } from './journal-stop-marks' import { journalQueuePauseRestatement } from './queued-message-pause' import type { JournalReducerState } from './journal-reducer' import { JournalRowWriter } from './journal-row-writer' +import { JournalStepWriter } from './journal-step-writer' import { restoreJournalStore } from './journal-store-restore' import type { JournalRow } from './journal-row-schema' import type { AgentSessionJournal } from './journal-store' @@ -52,6 +53,7 @@ export type JournalStoreCollaborators = { epochController: JournalEpochController itemAppender: JournalItemAppender lifecycleBatchAppender: JournalLifecycleBatchAppender + stepWriter: JournalStepWriter queuedMessages: JournalQueuedMessages stopMarks: JournalStopMarks /** Restores the store's state from disk. Owned here because it needs the same @@ -100,6 +102,12 @@ export function createJournalStoreCollaborators(host: JournalStoreHost): Journal inTransaction: (db, row) => queuedMessages.onRowInTransaction(db, row), rolledBack: () => queuedMessages.invalidate() }) + const lifecycleBatchAppender = new JournalLifecycleBatchAppender({ + state: host.state, + cursor: host.cursor, + enqueue: host.enqueue, + enqueueRows: (plan) => rowWriter.enqueueRows(plan) + }) return { epochController, queuedMessages, @@ -115,11 +123,12 @@ export function createJournalStoreCollaborators(host: JournalStoreHost): Journal state: host.state, enqueue: host.enqueue }), - lifecycleBatchAppender: new JournalLifecycleBatchAppender({ + lifecycleBatchAppender, + stepWriter: new JournalStepWriter({ + serialize: host.serialize, state: host.state, - cursor: host.cursor, - enqueue: host.enqueue, - enqueueRows: (plan) => rowWriter.enqueueRows(plan) + writeRows: (plan) => rowWriter.writeRows(plan), + planSettlement: (batch) => lifecycleBatchAppender.planResolved(batch) }) } } diff --git a/src/main/native-chat/agent-session-journal/journal-store-contracts.ts b/src/main/native-chat/agent-session-journal/journal-store-contracts.ts index 4232e805bbe..c2547c658c8 100644 --- a/src/main/native-chat/agent-session-journal/journal-store-contracts.ts +++ b/src/main/native-chat/agent-session-journal/journal-store-contracts.ts @@ -80,6 +80,14 @@ export type JournalLifecycleBatchInput = { rejectsQueued?: AgentJournalDispatchRejection } +export type JournalResolvedLifecycleBatchInput = Omit< + JournalLifecycleBatchInput, + 'mutations' | 'rejectsQueued' +> & { + /** Read from the fold with every earlier write landed; may return none. */ + resolve: () => readonly JournalLifecycleMutationInput[] +} + export type JournalSubmissionInput = { clientMessageId: string payloadFingerprint: string diff --git a/src/main/native-chat/agent-session-journal/journal-store.ts b/src/main/native-chat/agent-session-journal/journal-store.ts index 90d5a826946..626b77cf614 100644 --- a/src/main/native-chat/agent-session-journal/journal-store.ts +++ b/src/main/native-chat/agent-session-journal/journal-store.ts @@ -8,6 +8,7 @@ import type { AgentJournalCursor, AgentJournalItemBody, AgentJournalItemIdentity, + AgentJournalRenderItem, AgentJournalSnapshot, AgentJournalSubmission, AgentJournalThreadGoal, @@ -67,8 +68,9 @@ import type { JournalOperationReceipt, JournalRowWriter } from './journal-row-wr import type { JournalEpochController } from './journal-epoch-controller' import { JournalWriteQueue } from './journal-write-queue' import { createJournalStoreCollaborators } from './journal-store-collaborators' -import type { JournalItemAppender, JournalResolvedItem } from './journal-item-appender' +import type { JournalItemAppender } from './journal-item-appender' import type { JournalLifecycleBatchAppender } from './journal-lifecycle-batch-appender' +import type { JournalStepWriter } from './journal-step-writer' import type { JournalStopMarks } from './journal-stop-marks' export { AgentSessionJournalError } from './journal-write-guards' @@ -87,6 +89,7 @@ export class AgentSessionJournal { private readonly epochController: JournalEpochController private readonly itemAppender: JournalItemAppender private readonly lifecycleBatchAppender: JournalLifecycleBatchAppender + private readonly stepWriter: JournalStepWriter private readonly restore: () => Promise /** Draft rows queued while the agent works; never reducer input or owed work. */ readonly queuedMessages: JournalQueuedMessages @@ -129,6 +132,7 @@ export class AgentSessionJournal { this.epochController = collaborators.epochController this.itemAppender = collaborators.itemAppender this.lifecycleBatchAppender = collaborators.lifecycleBatchAppender + this.stepWriter = collaborators.stepWriter this.queuedMessages = collaborators.queuedMessages this.stopMarks = collaborators.stopMarks this.restore = collaborators.restore @@ -204,6 +208,9 @@ export class AgentSessionJournal { itemBody = (itemId: string): AgentJournalItemBody | null => this.state.items.get(itemId)?.body ?? null + /** One reduced item with its attribution, for a writer that needs the turn a row joined. */ + item = (itemId: string): AgentJournalRenderItem | null => this.state.items.get(itemId) ?? null + /** Visits reduced items with the producer that wrote each, for a producer re-deriving what an * earlier run of this session left. */ visitItemsWithLinkage = (visit: JournalItemLinkageVisitor): void => { @@ -276,12 +283,8 @@ export class AgentSessionJournal { } /** An upsert whose row is chosen from the fold at its own turn in the queue; null writes nothing. */ - appendResolvedItem( - resolve: () => JournalResolvedItem | null, - options: JournalItemAppendOptions - ): Promise { - return this.itemAppender.appendResolved(resolve, options) - } + appendResolvedItem: JournalItemAppender['appendResolved'] = (resolve, options) => + this.itemAppender.appendResolved(resolve, options) appendTombstone( identity: AgentJournalItemIdentity, @@ -307,6 +310,9 @@ export class AgentSessionJournal { return this.lifecycleBatchAppender.append(input) } + /** Several writes as one turn in the queue; see `JournalStepWriter`. */ + appendSteps: JournalStepWriter['append'] = (steps) => this.stepWriter.append(steps) + /** * Write-ahead submission row. It is durable before the caller dispatches * anything, and it doubles as the optimistic user bubble so an accepted echo diff --git a/src/main/claude/claude-turn-row-revision.ts b/src/main/native-chat/agent-session-timeline/agent-journal-turn-row-revision.ts similarity index 68% rename from src/main/claude/claude-turn-row-revision.ts rename to src/main/native-chat/agent-session-timeline/agent-journal-turn-row-revision.ts index 1b684ba89aa..dfede8c46f7 100644 --- a/src/main/claude/claude-turn-row-revision.ts +++ b/src/main/native-chat/agent-session-timeline/agent-journal-turn-row-revision.ts @@ -1,4 +1,4 @@ -// Every write to a Claude turn row is a revision of the row as the journal holds +// Every write to a turn row is a revision of the row as the journal holds // it at execution: a writer overrides only the fields it owns and keeps the // rest, whoever wrote them. Nothing about ended turns is kept in memory, so a // restart or reattach revises the same rows a live translator would. @@ -9,34 +9,36 @@ import { MAX_CONTEXT_MODEL_ID_CHARS, type AgentSessionContextUsage, type AgentSessionContextWindow -} from '../../shared/agent-session-context-usage' +} from '../../../shared/agent-session-context-usage' import { agentJournalItemKey, parseAgentJournalItemKey -} from '../../shared/agent-session-journal-item-key' +} from '../../../shared/agent-session-journal-item-key' import { AGENT_JOURNAL_THREAD_SCOPE, type AgentJournalItemIdentity, type AgentJournalTurnItem -} from '../../shared/agent-session-journal-types' -import { readAgentJournalTurn } from '../../shared/agent-session-turn-record' -import { estimateStructuredAgentSessionItemBytes } from '../native-chat/agent-session-wire/structured-agent-session-event-sink-estimate' +} from '../../../shared/agent-session-journal-types' +import { readAgentJournalTurn } from '../../../shared/agent-session-turn-record' +import { estimateStructuredAgentSessionItemBytes } from '../agent-session-wire/structured-agent-session-event-sink-estimate' import type { StructuredAgentSessionEventSink, StructuredAgentSessionRevisionJournal, StructuredAgentSessionRevisionOptions -} from '../native-chat/agent-session-wire/structured-agent-session-event-sink' +} from '../agent-session-wire/structured-agent-session-event-sink' /** The row a write revises: a known one, or the newest turn in the journal. */ -export type ClaudeTurnRowTarget = { identity: AgentJournalItemIdentity } | { newest: true } +export type AgentJournalTurnRowTarget = { identity: AgentJournalItemIdentity } | { newest: true } -export type ClaudeTurnRowWrite = { +export type AgentJournalTurnRowWrite = { /** The lifecycle fields this turn has now; the lifecycle writer owns all of them. */ lifecycle?: AgentJournalTurnItem /** Context parts, each replacing its namesake on the row. */ contextUsage?: AgentSessionContextUsage /** A window written only while no row in the journal holds one. */ windowIfNoneHeld?: AgentSessionContextWindow + /** Lands only while the row is absent or still running, so an ended turn is never rewritten. */ + onlyWhileRunning?: true } function jsonBytes(value: unknown): number { @@ -84,7 +86,7 @@ const TURN_ROW_BYTES_WITHOUT_CONTEXT = 8 * 1024 /** `lifecycle` over the row's current body, keeping every field it does not own. */ function reviseTurnBody( current: AgentJournalTurnItem | null, - write: ClaudeTurnRowWrite + write: AgentJournalTurnRowWrite ): AgentJournalTurnItem | null { const { lifecycle, contextUsage } = write const base = @@ -113,8 +115,8 @@ function withoutLifecycle(turn: AgentJournalTurnItem) { /** The write with its fallback window resolved against every row, since any of them may hold the newest. */ function withFallbackWindow( journal: StructuredAgentSessionRevisionJournal, - { windowIfNoneHeld, ...write }: ClaudeTurnRowWrite -): ClaudeTurnRowWrite { + { windowIfNoneHeld, ...write }: AgentJournalTurnRowWrite +): AgentJournalTurnRowWrite { if (!windowIfNoneHeld || write.contextUsage?.window) { return write } @@ -129,7 +131,7 @@ function withFallbackWindow( function findTurnRow( journal: StructuredAgentSessionRevisionJournal, - target: ClaudeTurnRowTarget + target: AgentJournalTurnRowTarget ): { itemId: string; body: AgentJournalTurnItem } | null { if ('identity' in target) { const itemId = agentJournalItemKey(target.identity) @@ -147,22 +149,59 @@ function findTurnRow( } /** How a turn-row write reaches live subscribers. */ -export type ClaudeTurnRowDelivery = { +export type AgentJournalTurnRowDelivery = { /** False only for a writer that publishes each write itself right after queueing it. */ publish: boolean options?: Omit } +/** Bytes a turn-row write may resolve to. */ +export function agentJournalTurnRowReservedBytes( + target: AgentJournalTurnRowTarget, + write: AgentJournalTurnRowWrite +): number { + const known = 'identity' in target ? target.identity : null + return ( + (known && write.lifecycle + ? estimateStructuredAgentSessionItemBytes(known, write.lifecycle) + : TURN_ROW_BYTES_WITHOUT_CONTEXT) + contextBytesBound(write.contextUsage) + ) +} + +/** The row a turn-row write lands as, read from the journal at execution; null writes nothing. */ +export function resolveAgentJournalTurnRowWrite( + journal: StructuredAgentSessionRevisionJournal, + target: AgentJournalTurnRowTarget, + write: AgentJournalTurnRowWrite, + reservedBytes: number +): { identity: AgentJournalItemIdentity; body: AgentJournalTurnItem } | null { + const known = 'identity' in target ? target.identity : null + const row = findTurnRow(journal, target) + if (write.onlyWhileRunning && row && row.body.state !== 'running') { + return null + } + const identity = known ?? (row ? parseAgentJournalItemKey(row.itemId) : null) + const body = reviseTurnBody(row?.body ?? null, withFallbackWindow(journal, write)) + if (!identity || !body) { + return null + } + if (estimateStructuredAgentSessionItemBytes(identity, body) <= reservedBytes) { + return { identity, body } + } + // Only a row grown by fields this build does not know gets here; the lifecycle still lands. + return write.lifecycle ? { identity, body: write.lifecycle } : null +} + /** - * Queue one revision of a Claude turn row. A lifecycle write creates the row + * Queue one revision of a turn row. A lifecycle write creates the row * when it is absent; a context write only ever revises one that exists. * Without a journal-reading sink only the lifecycle is written, as it was built. */ -export function writeClaudeTurnRow( +export function writeAgentJournalTurnRow( sink: StructuredAgentSessionEventSink, - target: ClaudeTurnRowTarget, - write: ClaudeTurnRowWrite, - { publish, options: delivery = {} }: ClaudeTurnRowDelivery + target: AgentJournalTurnRowTarget, + write: AgentJournalTurnRowWrite, + { publish, options: delivery = {} }: AgentJournalTurnRowDelivery ): void { // A turn record belongs to no turn. const options = { ...delivery, turnScope: AGENT_JOURNAL_THREAD_SCOPE } @@ -176,27 +215,11 @@ export function writeClaudeTurnRow( } return } - const known = 'identity' in target ? target.identity : null - const reservedBytes = - (known && write.lifecycle - ? estimateStructuredAgentSessionItemBytes(known, write.lifecycle) - : TURN_ROW_BYTES_WITHOUT_CONTEXT) + contextBytesBound(write.contextUsage) + const reservedBytes = agentJournalTurnRowReservedBytes(target, write) revise.call( sink, reservedBytes, - (journal) => { - const row = findTurnRow(journal, target) - const identity = known ?? (row ? parseAgentJournalItemKey(row.itemId) : null) - const body = reviseTurnBody(row?.body ?? null, withFallbackWindow(journal, write)) - if (!identity || !body) { - return null - } - if (estimateStructuredAgentSessionItemBytes(identity, body) <= reservedBytes) { - return { identity, body } - } - // Only a row grown by fields this build does not know gets here; the lifecycle still lands. - return write.lifecycle ? { identity, body: write.lifecycle } : null - }, + (journal) => resolveAgentJournalTurnRowWrite(journal, target, write, reservedBytes), options ) } diff --git a/src/main/native-chat/agent-session-wire/agent-session-delta-coalescer.ts b/src/main/native-chat/agent-session-wire/agent-session-delta-coalescer.ts index b1d4c483d9b..eff850ae8e0 100644 --- a/src/main/native-chat/agent-session-wire/agent-session-delta-coalescer.ts +++ b/src/main/native-chat/agent-session-wire/agent-session-delta-coalescer.ts @@ -56,6 +56,10 @@ export type AgentSessionDeltaCoalescer = { dispose: () => void /** Bounded last-known state for terminalizing a rejected completion. */ snapshot: (key: string) => AgentSessionDeltaSnapshot | null + /** Streams with text not yet emitted, in arrival order, for a caller that writes them itself. */ + dirty: () => { key: string; snapshot: AgentSessionDeltaSnapshot }[] + /** The caller wrote this stream's current text itself; nothing is owed until it grows. */ + markFlushed: (key: string) => void } function defaultSchedule(run: () => void, ms: number): () => void { @@ -215,6 +219,27 @@ export function createAgentSessionDeltaCoalescer( truncated: stream.truncated } : null + }, + dirty: () => + [...streams].flatMap(([key, stream]) => + stream.dirty + ? [ + { + key, + snapshot: { + text: stream.chunks.join(''), + observedBytes: stream.observedBytes, + truncated: stream.truncated + } + } + ] + : [] + ), + markFlushed: (key) => { + const stream = streams.get(key) + if (stream) { + stream.dirty = false + } } } } diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-event-sink-queue.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-event-sink-queue.ts index afcecadc744..16f636af093 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-event-sink-queue.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-event-sink-queue.ts @@ -8,6 +8,7 @@ import type { StructuredAgentSessionSinkState, StructuredAgentSessionSinkWatermarks } from './structured-agent-session-event-sink' +import type { StructuredAgentSessionTransitionJournal } from './structured-agent-session-transition' export type StructuredAgentSessionSinkOperation = { sequence: number @@ -85,6 +86,8 @@ export class StructuredAgentSessionSinkQueue { journalLinkage = (): StructuredAgentSessionLinkageJournal | null => this.target?.journal ?? null + journalItems = (): StructuredAgentSessionTransitionJournal | null => this.target?.journal ?? null + journalStopDecidesTurn = (turnId: string, endedAt: number): boolean => this.target?.journal.stopMarks.personStopDecides(turnId, endedAt) ?? false diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-event-sink.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-event-sink.ts index 12966b73a97..4a09b9fcacc 100644 --- a/src/main/native-chat/agent-session-wire/structured-agent-session-event-sink.ts +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-event-sink.ts @@ -12,6 +12,10 @@ import { estimateStructuredAgentSessionItemBytes } from './structured-agent-sess import { StructuredAgentSessionSinkQueue } from './structured-agent-session-event-sink-queue' import { structuredAgentSessionJournalAppendOptions } from './structured-agent-session-journal-append-options' import { createStructuredAgentSessionResolvedAppend } from './structured-agent-session-resolved-append' +import { + createStructuredAgentSessionTransitionMembers, + type StructuredAgentSessionTransitionSink +} from './structured-agent-session-transition' import type { StructuredAgentSessionLogger } from './structured-agent-session-logger' export type StructuredAgentSessionSinkAdmission = @@ -80,7 +84,7 @@ export type StructuredAgentSessionRevisionOptions = StructuredAgentSessionItemAp /** Compatibility alias for lifecycle callers that already use this resolver. */ export type StructuredAgentSessionLifecycleIdentityResolver = StructuredAgentSessionIdentityResolver -export type StructuredAgentSessionEventSink = { +export type StructuredAgentSessionEventSink = StructuredAgentSessionTransitionSink & { appendItem( identity: AgentJournalItemIdentity, body: AgentJournalItemBody, @@ -285,6 +289,7 @@ export function createDeferredStructuredAgentSessionEventSink(deps: { }, tryAppendItem: appendItem, ...resolvedAppend, + ...createStructuredAgentSessionTransitionMembers(queue), journalEpoch: queue.journalEpoch, journalLinkage: queue.journalLinkage, journalStopDecidesTurn: queue.journalStopDecidesTurn, diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-transition.test.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-transition.test.ts new file mode 100644 index 00000000000..d9a152d0e6e --- /dev/null +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-transition.test.ts @@ -0,0 +1,520 @@ +import { mkdtemp, rm } from 'node:fs/promises' +import { tmpdir } from 'node:os' +import { join } from 'node:path' +import { afterEach, describe, expect, it, vi } from 'vitest' +import { agentJournalItemKey } from '../../../shared/agent-session-journal-item-key' +import { + AGENT_JOURNAL_THREAD_SCOPE, + type AgentJournalItemIdentity, + type AgentJournalToolCallItem +} from '../../../shared/agent-session-journal-types' +import { + createTrackedJournalOpener, + openTestJournalHostDatabase +} from '../agent-session-journal/journal-host-database-test-support' +import type { JournalLifecycleMutationInput } from '../agent-session-journal/journal-row-builders' +import { openJournalOwingImport } from '../agent-session-journal/journal-owed-import-test-support' +import type { AgentSessionJournal } from '../agent-session-journal/journal-store' +import { createAgentSessionDeltaCoalescer } from './agent-session-delta-coalescer' +import { + createDeferredStructuredAgentSessionEventSink, + type StructuredAgentSessionSinkWatermarks +} from './structured-agent-session-event-sink' +import { testEventSinkLogging } from './structured-agent-session-logger-test-support' +import type { + StructuredAgentSessionTransition, + StructuredAgentSessionTransitionStep +} from './structured-agent-session-transition' + +const SESSION = 'session-transition' +const journals = createTrackedJournalOpener() +const roots: string[] = [] + +afterEach(async () => { + await journals.closeAll() + await Promise.all(roots.splice(0).map((root) => rm(root, { recursive: true, force: true }))) +}) + +function identity(recordId: string): AgentJournalItemIdentity { + return { provider: 'legacy', agent: 'grok', sessionId: SESSION, recordId } +} + +function tool(name: string, state: AgentJournalToolCallItem['state']): AgentJournalToolCallItem { + return { kind: 'tool-call', name, input: { name }, state } +} + +function failedTool(id: AgentJournalItemIdentity): JournalLifecycleMutationInput { + return { + kind: 'item', + identity: id, + body: tool('read', 'failed'), + turnScope: AGENT_JOURNAL_THREAD_SCOPE + } +} + +function itemStep( + resolve: Extract['resolve'], + reservedBytes = 4096 +): StructuredAgentSessionTransitionStep { + return { + kind: 'item', + reservedBytes, + resolve, + options: { turnScope: AGENT_JOURNAL_THREAD_SCOPE } + } +} + +async function rig(watermarks: Partial = {}) { + const root = await mkdtemp(join(tmpdir(), 'orca-transition-')) + roots.push(root) + const journal: AgentSessionJournal = await journals.open({ + identity: { + sessionId: SESSION, + workspaceId: 'workspace-1', + hostId: 'local', + agent: 'grok', + providerHandle: { transport: 'acp', agent: 'grok', nativeId: 'provider-session-1' } + }, + stateDirectory: root, + now: () => 1_000 + }) + const publishes: number[] = [] + const deferred = createDeferredStructuredAgentSessionEventSink({ + ...testEventSinkLogging(SESSION), + watermarks + }) + const bind = () => deferred.bind({ journal, fence: 1, publish: () => publishes.push(1) }) + return { root, journal, deferred, sink: deferred.sink, publishes, bind } +} + +/** A bound sink over a chat whose copy into the host's database is still owed, so every write, + * the transition's included, waits in the journal's queue until the first one pays it. */ +async function owedRig() { + const root = await mkdtemp(join(tmpdir(), 'orca-transition-owed-')) + roots.push(root) + const { journal } = await openJournalOwingImport({ + stateDirectory: root, + identity: { + sessionId: SESSION, + workspaceId: 'workspace-1', + hostId: 'local', + agent: 'grok', + providerHandle: { transport: 'acp', agent: 'grok', nativeId: 'provider-session-1' } + }, + now: () => 1_000 + }) + const failures: unknown[] = [] + const deferred = createDeferredStructuredAgentSessionEventSink({ + ...testEventSinkLogging(SESSION), + onFailed: (error) => failures.push(error) + }) + deferred.bind({ journal, fence: 1, publish: () => undefined }) + const ownKeys = () => + journal + .snapshot() + .items.map((item) => item.itemId) + .filter((itemId) => itemId.includes(SESSION)) + return { journal, deferred, sink: deferred.sink, failures, ownKeys, close: () => journal.close() } +} + +describe('structured agent-session transitions', () => { + it('lands its steps back to back, each resolved after the one before it', async () => { + const { journal, deferred, sink, publishes, bind } = await rig() + bind() + const first = identity('first') + const transition: StructuredAgentSessionTransition = { + lifecycle: false, + publish: true, + steps: [ + itemStep(() => ({ identity: first, body: tool('read', 'running') })), + // Reads the row the step before it wrote. + itemStep((view) => { + const before = view.itemBody(agentJournalItemKey(first)) + return before?.kind === 'tool-call' + ? { identity: identity('second'), body: tool(`after-${before.name}`, 'running') } + : null + }) + ] + } + expect(sink.tryAppendTransition?.(transition)).toEqual({ accepted: true }) + sink.tryAppendItem?.(identity('later'), tool('later', 'running'), { + turnScope: AGENT_JOURNAL_THREAD_SCOPE + }) + await deferred.drained() + + const items = journal.snapshot().items + const start = items[0]?.sequence ?? 0 + expect(items.map((item) => [item.itemId, item.sequence - start])).toEqual([ + [agentJournalItemKey(first), 0], + [agentJournalItemKey(identity('second')), 1], + [agentJournalItemKey(identity('later')), 2] + ]) + expect(items[1]?.body).toMatchObject({ name: 'after-read' }) + expect(publishes).toHaveLength(1) + }) + + it('admitted whole is not executed whole: a failed step keeps the ones before it and fails the sink', async () => { + const { journal, deferred, sink, publishes, bind } = await rig() + bind() + sink.tryAppendTransition?.({ + lifecycle: false, + publish: true, + steps: [ + itemStep(() => ({ identity: identity('kept'), body: tool('read', 'running') })), + itemStep(() => { + throw new Error('resolver failed') + }) + ] + }) + + await expect(deferred.drained()).resolves.toMatchObject({ ok: false }) + expect(journal.snapshot().items.map((item) => item.itemId)).toEqual([ + agentJournalItemKey(identity('kept')) + ]) + expect(publishes).toEqual([]) + expect( + sink.tryAppendTransition?.({ + lifecycle: false, + publish: false, + steps: [itemStep(() => null)] + }) + ).toEqual({ accepted: false, reason: 'failed' }) + }) + + it('writes nothing after a step that overflows its reservation', async () => { + const { journal, deferred, sink, publishes, bind } = await rig() + bind() + const after = vi.fn(() => ({ identity: identity('third'), body: tool('read', 'running') })) + sink.tryAppendTransition?.({ + lifecycle: false, + publish: true, + steps: [ + itemStep(() => ({ identity: identity('first'), body: tool('read', 'running') })), + itemStep( + () => ({ identity: identity('second'), body: tool('x'.repeat(10_000), 'running') }), + 64 + ), + itemStep(after) + ] + }) + + await expect(deferred.drained()).resolves.toMatchObject({ ok: false }) + expect(journal.snapshot().items.map((item) => item.itemId)).toEqual([ + agentJournalItemKey(identity('first')) + ]) + expect(after).not.toHaveBeenCalled() + expect(publishes).toEqual([]) + }) + + it('opens no next-turn work after a settlement the journal refused', async () => { + const { journal, deferred, sink, bind } = await rig() + bind() + const old = identity('old-tool') + sink.tryAppendItem?.(old, tool('read', 'running'), { turnScope: AGENT_JOURNAL_THREAD_SCOPE }) + await deferred.drained() + sink.tryAppendTransition?.({ + lifecycle: true, + publish: true, + steps: [ + itemStep(() => ({ identity: identity('before'), body: tool('read', 'completed') })), + // Names one item twice, which the journal refuses. + { + kind: 'settlement', + settlementId: 'settle', + reservedBytes: 1, + resolve: () => [failedTool(old), failedTool(old)] + }, + itemStep(() => ({ identity: identity('next-turn-work'), body: tool('read', 'running') })) + ] + }) + + await expect(deferred.drained()).resolves.toMatchObject({ ok: false }) + expect(journal.snapshot().items.map((item) => item.itemId)).toEqual([ + agentJournalItemKey(old), + agentJournalItemKey(identity('before')) + ]) + expect(journal.item(agentJournalItemKey(old))?.body).toMatchObject({ state: 'running' }) + }) + + it('writes nothing after a step whose row the database refused', async () => { + const { root, journal, deferred, sink, bind } = await rig() + bind() + // Refuses only the second step's row, so the third would land in its place. + openTestJournalHostDatabase(root).db.exec(`CREATE TEMP TRIGGER fail_second_step +BEFORE INSERT ON main.journal_rows WHEN instr(NEW.row_json, 'refused-step') > 0 +BEGIN SELECT RAISE(ABORT, 'second step refused'); END`) + const step = (recordId: string) => + itemStep(() => ({ identity: identity(recordId), body: tool(recordId, 'running') })) + sink.tryAppendTransition?.({ + lifecycle: false, + publish: true, + steps: [step('first'), step('refused-step'), step('third')] + }) + + await expect(deferred.drained()).resolves.toMatchObject({ ok: false }) + expect(journal.snapshot().items.map((item) => item.itemId)).toEqual([ + agentJournalItemKey(identity('first')) + ]) + }) + + it('refuses a transition whole, so none of its steps ever lands', async () => { + const { journal, deferred, sink, bind } = await rig({ maxQueuedOperations: 1 }) + const step = (recordId: string) => + itemStep(() => ({ identity: identity(recordId), body: tool(recordId, 'running') })) + const admitted = { lifecycle: false, publish: false, steps: [step('a')] } + const refused = { lifecycle: false, publish: false, steps: [step('b'), step('c')] } + + expect(sink.tryAppendTransition?.(admitted)).toEqual({ accepted: true }) + expect(sink.tryAppendTransition?.(refused)).toEqual({ + accepted: false, + reason: 'backpressure' + }) + bind() + await deferred.drained() + + expect(journal.snapshot().items.map((item) => item.itemId)).toEqual([ + agentJournalItemKey(identity('a')) + ]) + }) + + it('settles from the rows as they stand, in consecutive rows when one cannot hold them', async () => { + const { journal, deferred, sink, publishes, bind } = await rig() + bind() + const running = Array.from({ length: 250 }, (_, index) => identity(`tool-${index}`)) + for (const id of running) { + sink.tryAppendItem?.(id, tool('read', 'running'), { turnScope: AGENT_JOURNAL_THREAD_SCOPE }) + } + sink.tryAppendItem?.(identity('done'), tool('read', 'completed'), { + turnScope: AGENT_JOURNAL_THREAD_SCOPE + }) + expect( + sink.tryAppendTransition?.({ + lifecycle: true, + publish: true, + steps: [ + { + kind: 'settlement', + settlementId: 'settle-all', + reservedBytes: 1, + // The settled row is a candidate too; the fold, not the caller, rules it out. + resolve: (view) => + [...running, identity('done')].flatMap((id) => { + const body = view.itemBody(agentJournalItemKey(id)) + return body?.kind === 'tool-call' && body.state === 'running' + ? [ + { + kind: 'item' as const, + identity: id, + body: { ...body, state: 'failed' as const }, + turnScope: AGENT_JOURNAL_THREAD_SCOPE + } + ] + : [] + }) + } + ] + }) + ).toEqual({ accepted: true }) + sink.tryAppendItem?.(identity('after'), tool('after', 'running'), { + turnScope: AGENT_JOURNAL_THREAD_SCOPE + }) + await deferred.drained() + + const items = journal.snapshot().items + const sequenceOf = (recordId: string) => + items.find((item) => item.itemId === agentJournalItemKey(identity(recordId)))?.sequence ?? 0 + expect( + items.filter((item) => item.body.kind === 'tool-call' && item.body.state === 'failed') + ).toHaveLength(250) + expect( + items.find((item) => item.itemId === agentJournalItemKey(identity('done')))?.body + ).toMatchObject({ state: 'completed' }) + // Two batch rows, back to back, between the last append before it and the first after it. + expect(sequenceOf('after') - sequenceOf('done')).toBe(3) + expect(publishes).toHaveLength(1) + }) + + it('leaves no row of a settlement durable when a later one of its rows fails to write', async () => { + const { root, journal, deferred, sink, bind } = await rig() + bind() + const running = Array.from({ length: 250 }, (_, index) => identity(`tool-${index}`)) + for (const id of running) { + sink.tryAppendItem?.(id, tool('read', 'running'), { turnScope: AGENT_JOURNAL_THREAD_SCOPE }) + } + await deferred.drained() + const before = journal.cursor().sequence + // The settlement needs two rows; the second one's insert aborts inside the transaction. + openTestJournalHostDatabase(root).db.exec(`CREATE TEMP TRIGGER fail_second_chunk +BEFORE INSERT ON main.journal_rows WHEN NEW.seq = ${before + 2} +BEGIN SELECT RAISE(ABORT, 'second chunk refused'); END`) + sink.tryAppendTransition?.({ + lifecycle: true, + publish: true, + steps: [ + { + kind: 'settlement', + settlementId: 'settle-all', + reservedBytes: 1, + resolve: () => running.map(failedTool) + } + ] + }) + + await expect(deferred.drained()).resolves.toMatchObject({ ok: false }) + expect(journal.cursor().sequence).toBe(before) + expect( + journal + .snapshot() + .items.filter((item) => item.body.kind === 'tool-call' && item.body.state === 'failed') + ).toEqual([]) + }) + + it('writes and announces nothing when a step resolves to nothing', async () => { + const { journal, deferred, sink, publishes, bind } = await rig() + bind() + sink.tryAppendTransition?.({ + lifecycle: true, + publish: true, + steps: [ + itemStep(() => null), + { kind: 'settlement', settlementId: 'none', reservedBytes: 1, resolve: () => [] } + ] + }) + await deferred.drained() + + expect(journal.snapshot().items).toEqual([]) + expect(publishes).toEqual([]) + }) + + it('lands a write a step issues while it runs after the whole transition', async () => { + const { journal, deferred, sink, bind } = await rig() + bind() + sink.tryAppendTransition?.({ + lifecycle: false, + publish: false, + steps: [ + itemStep(() => { + // Another writer, reached from inside the step: it waits for the transition's turn to end. + sink.tryAppendItem?.(identity('nested'), tool('nested', 'running'), { + turnScope: AGENT_JOURNAL_THREAD_SCOPE + }) + return { identity: identity('first'), body: tool('first', 'running') } + }), + itemStep(() => ({ identity: identity('second'), body: tool('second', 'running') })) + ] + }) + await deferred.drained() + + expect(journal.snapshot().items.map((item) => item.itemId)).toEqual( + ['first', 'second', 'nested'].map((recordId) => agentJournalItemKey(identity(recordId))) + ) + }) + + it('runs its steps back to back after an owed import, ahead of a write issued after it', async () => { + const { journal, deferred, sink, ownKeys, close } = await owedRig() + try { + sink.tryAppendTransition?.({ + lifecycle: false, + publish: false, + steps: [ + itemStep(() => ({ identity: identity('first'), body: tool('read', 'running') })), + itemStep((view) => + view.itemBody(agentJournalItemKey(identity('first'))) + ? { identity: identity('second'), body: tool('read', 'running') } + : null + ) + ] + }) + sink.tryAppendItem?.(identity('later'), tool('later', 'running'), { + turnScope: AGENT_JOURNAL_THREAD_SCOPE + }) + expect(journal.importPending).toBe(true) + await expect(deferred.drained()).resolves.toEqual({ ok: true }) + + expect(journal.importPending).toBe(false) + expect(ownKeys()).toEqual( + ['first', 'second', 'later'].map((recordId) => agentJournalItemKey(identity(recordId))) + ) + } finally { + await close() + } + }) + + it('stops at a failed middle step after an owed import; a write accepted before the failure still lands', async () => { + const { deferred, sink, failures, ownKeys, close } = await owedRig() + try { + const third = vi.fn(() => ({ identity: identity('third'), body: tool('read', 'running') })) + sink.tryAppendTransition?.({ + lifecycle: false, + publish: false, + steps: [ + itemStep(() => ({ identity: identity('first'), body: tool('read', 'running') })), + itemStep( + () => ({ identity: identity('second'), body: tool('x'.repeat(10_000), 'running') }), + 64 + ), + itemStep(third) + ] + }) + // Accepted before the failure is known: the sink's existing rule lets it land. + sink.tryAppendItem?.(identity('later'), tool('later', 'running'), { + turnScope: AGENT_JOURNAL_THREAD_SCOPE + }) + await expect(deferred.drained()).resolves.toMatchObject({ ok: false }) + + expect(ownKeys()).toEqual( + ['first', 'later'].map((recordId) => agentJournalItemKey(identity(recordId))) + ) + expect(third).not.toHaveBeenCalled() + expect(failures).toHaveLength(1) + } finally { + await close() + } + }) +}) + +describe('resolved lifecycle batches', () => { + it('refuses, before writing anything, a settlement that names one item twice', async () => { + const { journal } = await rig() + const before = journal.cursor().sequence + const twice = identity('twice') + + await expect( + journal.appendSteps([ + { + kind: 'settlement', + batch: { + settlementId: 'twice', + fence: 1, + resolve: () => [failedTool(identity('once')), failedTool(twice), failedTool(twice)] + } + } + ]) + ).rejects.toThrow('journal_resolved_lifecycle_batch_names_item_twice') + expect(journal.cursor().sequence).toBe(before) + }) +}) + +describe('coalescer text a caller writes itself', () => { + it('reports unwritten streams and owes nothing once the caller marks them written', () => { + const emitted: string[] = [] + const coalescer = createAgentSessionDeltaCoalescer({ + emit: (_key, text) => { + emitted.push(text) + }, + schedule: () => () => {} + }) + coalescer.append('a', 'Hel') + coalescer.append('b', 'Wor') + coalescer.append('a', 'lo') + + expect(coalescer.dirty().map(({ key, snapshot }) => [key, snapshot.text])).toEqual([ + ['a', 'Hello'], + ['b', 'Wor'] + ]) + coalescer.markFlushed('a') + coalescer.flushAll() + expect(emitted).toEqual(['Wor']) + expect(coalescer.dirty()).toEqual([]) + }) +}) diff --git a/src/main/native-chat/agent-session-wire/structured-agent-session-transition.ts b/src/main/native-chat/agent-session-wire/structured-agent-session-transition.ts new file mode 100644 index 00000000000..f13b897f12b --- /dev/null +++ b/src/main/native-chat/agent-session-wire/structured-agent-session-transition.ts @@ -0,0 +1,126 @@ +// One provider event's journal writes, admitted as a single sink operation. +// +// A producer that keeps state about what it wrote must change that state only for writes the sink +// took; when an event needs several rows, a refusal of the third after the first two were taken +// would leave the producer and the journal disagreeing. A transition is admitted whole or not at +// all. At execution its steps run as one turn in the journal's write queue, each resolved against +// the fold with every earlier write landed (the steps before it included), so what a step writes +// is decided by the journal, not by memory. Admitted whole, executed as a prefix: once a step +// fails, the steps after it never run, the steps before it stay written, and the sink fails. + +import type { + AgentJournalItemBody, + AgentJournalItemIdentity +} from '../../../shared/agent-session-journal-types' +import type { JournalLifecycleMutationInput } from '../agent-session-journal/journal-row-builders' +import type { AgentSessionJournal } from '../agent-session-journal/journal-store' +import type { JournalStep } from '../agent-session-journal/journal-step-writer' +import { estimateStructuredAgentSessionItemBytes } from './structured-agent-session-event-sink-estimate' +import type { + StructuredAgentSessionItemAppendOptions, + StructuredAgentSessionSinkAdmission +} from './structured-agent-session-event-sink' +import type { StructuredAgentSessionSinkQueue } from './structured-agent-session-event-sink-queue' +import { structuredAgentSessionJournalAppendOptions } from './structured-agent-session-journal-append-options' + +/** What a transition step reads: rows by key, every row, and the turns they joined. */ +export type StructuredAgentSessionTransitionJournal = Pick< + AgentSessionJournal, + 'epoch' | 'visitItems' | 'itemBody' | 'item' | 'visitItemsWithLinkage' +> + +export type StructuredAgentSessionTransitionStep = + | { + kind: 'item' + /** Bounds what `resolve` may write; a larger write fails the sink. */ + reservedBytes: number + /** The row and its whole body; null writes nothing. */ + resolve: (journal: StructuredAgentSessionTransitionJournal) => { + identity: AgentJournalItemIdentity + body: AgentJournalItemBody + } | null + options: StructuredAgentSessionItemAppendOptions + } + | { + kind: 'settlement' + /** Unique per settlement: the journal applies one id once. */ + settlementId: string + /** Paces the queue only; the mutations are the journal's to choose. */ + reservedBytes: number + /** Read at execution; none writes nothing. */ + resolve: ( + journal: StructuredAgentSessionTransitionJournal + ) => readonly JournalLifecycleMutationInput[] + } + +export type StructuredAgentSessionTransition = { + steps: readonly StructuredAgentSessionTransitionStep[] + /** Rides the sink's lifecycle budget: it ends or settles something. */ + lifecycle: boolean + /** Announce the writes once they land, when any step wrote. */ + publish: boolean +} + +const STEP_OVERFLOW = 'structured agent-session transition step exceeded its reserved size' + +/** The sink members a transition writer uses. */ +export type StructuredAgentSessionTransitionSink = { + /** Queues one event's writes as a single admitted operation. */ + tryAppendTransition?( + transition: StructuredAgentSessionTransition + ): StructuredAgentSessionSinkAdmission + /** The bound journal's rows as they stand now; null until bound. */ + journalItems?(): StructuredAgentSessionTransitionJournal | null +} + +export function createStructuredAgentSessionTransitionMembers( + queue: StructuredAgentSessionSinkQueue +): Required { + return { tryAppendTransition: transitionAppend(queue), journalItems: queue.journalItems } +} + +function transitionAppend( + queue: StructuredAgentSessionSinkQueue +): (transition: StructuredAgentSessionTransition) => StructuredAgentSessionSinkAdmission { + return (transition) => + queue.submit({ + bytes: + transition.steps.reduce((total, step) => total + step.reservedBytes, 0) + + (transition.publish ? 1 : 0), + lifecycle: transition.lifecycle, + run: async (bound) => { + const { journal, fence } = bound + const wrote = await journal.appendSteps( + transition.steps.map((step): JournalStep => + step.kind === 'item' + ? { + kind: 'item', + resolve: () => { + const resolved = step.resolve(journal) + if ( + resolved && + estimateStructuredAgentSessionItemBytes(resolved.identity, resolved.body) > + step.reservedBytes + ) { + throw new Error(STEP_OVERFLOW) + } + return resolved + }, + options: structuredAgentSessionJournalAppendOptions(fence, step.options) + } + : { + kind: 'settlement', + batch: { + settlementId: step.settlementId, + fence, + resolve: () => step.resolve(journal) + } + } + ) + ) + if (transition.publish && wrote.includes(true)) { + bound.publish() + } + } + }) +}