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
This commit is contained in:
Brennan Benson
2026-10-05 22:20:59 -07:00
committed by GitHub
parent dc8b3934b2
commit c0b07d3a71
14 changed files with 909 additions and 83 deletions
+7 -4
View File
@@ -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 } : {}) },
+2 -2
View File
@@ -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 } : {}) },
@@ -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<string>()
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)
}
@@ -76,33 +76,36 @@ export class JournalRowWriter {
enqueueRows(
plan: () => readonly ((seq: number, ts: number) => JournalRow)[]
): Promise<JournalRow[]> {
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
@@ -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: <T>(run: JournalWriteBody<T>) => Promise<T>
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<boolean[]> {
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)]
}
}
@@ -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)
})
}
}
@@ -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
@@ -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<void>
/** 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<JournalAppendResult | null> {
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
@@ -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<StructuredAgentSessionRevisionOptions, 'turnScope'>
}
/** 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
)
}
@@ -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
}
}
}
}
@@ -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
@@ -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,
@@ -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<StructuredAgentSessionTransitionStep, { kind: 'item' }>['resolve'],
reservedBytes = 4096
): StructuredAgentSessionTransitionStep {
return {
kind: 'item',
reservedBytes,
resolve,
options: { turnScope: AGENT_JOURNAL_THREAD_SCOPE }
}
}
async function rig(watermarks: Partial<StructuredAgentSessionSinkWatermarks> = {}) {
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([])
})
})
@@ -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<StructuredAgentSessionTransitionSink> {
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()
}
}
})
}