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() + } + } + }) +}