| /* |
| * Licensed to the Apache Software Foundation (ASF) under one |
| * or more contributor license agreements. See the NOTICE file |
| * distributed with this work for additional information |
| * regarding copyright ownership. The ASF licenses this file |
| * to you under the Apache License, Version 2.0 (the |
| * "License"); you may not use this file except in compliance |
| * with the License. You may obtain a copy of the License at |
| * |
| * http://www.apache.org/licenses/LICENSE-2.0 |
| * |
| * Unless required by applicable law or agreed to in writing, |
| * software distributed under the License is distributed on an |
| * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY |
| * KIND, either express or implied. See the License for the |
| * specific language governing permissions and limitations |
| * under the License. |
| */ |
| |
| import { MODEL_FAILURE_MESSAGE_MAX_BYTES } from '@maka/core/model-failure'; |
| import { truncateUtf8 } from '@maka/core/diagnostic-log'; |
| import type { RuntimeInvocationRecord } from '@maka/core/runtime-invocation'; |
| import type { |
| AssistantStepContentKind, |
| StoredMessage, |
| TurnStatus, |
| WorkHubCoordinationActionMessage, |
| } from '@maka/core/session'; |
| import type { RuntimeEvent, RuntimeEventStatus } from '@maka/core/runtime-event'; |
| import type { ToolActivityKind, ToolResultContent } from '@maka/core/events'; |
| import { markPersisted } from '@maka/core/persisted-value'; |
| import { |
| SANDBOX_BOUNDARY_REQUEST_STATUSES, |
| validateSandboxBoundaryExpansion, |
| } from '@maka/core/sandbox-boundary'; |
| |
| import { TOOL_ACTIVITY_KINDS, normalizeMessageContent } from '@maka/core/events'; |
| |
| import { |
| isPartialRuntimeEvent, |
| isTerminalRuntimeEvent, |
| isTerminalRuntimeEventStatus, |
| } from '@maka/core/runtime-event'; |
| |
| import { decodePersistedToolResultContent } from '@maka/core/tool-result-record-schema'; |
| |
| /** The statuses a settled boundary decision can carry — every status but `pending`. */ |
| type SettledSandboxBoundaryStatus = Exclude< |
| (typeof SANDBOX_BOUNDARY_REQUEST_STATUSES)[number], |
| 'pending' |
| >; |
| const SETTLED_SANDBOX_BOUNDARY_STATUSES: readonly SettledSandboxBoundaryStatus[] = |
| SANDBOX_BOUNDARY_REQUEST_STATUSES.filter( |
| (status): status is SettledSandboxBoundaryStatus => status !== 'pending', |
| ); |
| import type { CanonicalPermissionOutcomeRecord } from './interaction-authority.js'; |
| import { isArchivedToolResultPlaceholder } from './tool-result-archive.js'; |
| import { truncateToolOutput } from './tool-output.js'; |
| |
| export type RuntimeEventReadModelDiagnosticCode = |
| | 'partial_skipped' |
| | 'unsupported_event' |
| | 'unclaimed_control_fact' |
| | 'incomplete_event' |
| | 'archived_tool_result_placeholder' |
| | 'generated_id' |
| | 'tool_use_id_mismatch' |
| | 'unexpected_projected_message'; |
| |
| /** |
| * Whether a diagnostic means the projection may have lost user-visible content. |
| * |
| * `hard` — a row a reader would have seen may be missing, so the projection is |
| * not a faithful view of the session and must not be served in place of one. |
| * `soft` — the fact is reported without withholding the view; it does not by |
| * itself mean a row is missing. |
| * |
| * The table is keyed by code so a new diagnostic cannot exist without deciding |
| * which side of that line it falls on. |
| */ |
| const RUNTIME_EVENT_READ_MODEL_DIAGNOSTIC_SEVERITY: Record< |
| RuntimeEventReadModelDiagnosticCode, |
| 'hard' | 'soft' |
| > = { |
| partial_skipped: 'soft', |
| unsupported_event: 'hard', |
| unclaimed_control_fact: 'soft', |
| incomplete_event: 'hard', |
| archived_tool_result_placeholder: 'soft', |
| generated_id: 'soft', |
| tool_use_id_mismatch: 'hard', |
| unexpected_projected_message: 'soft', |
| }; |
| |
| export function isHardRuntimeEventReadModelDiagnostic(diagnostic: { |
| code: RuntimeEventReadModelDiagnosticCode; |
| }): boolean { |
| return RUNTIME_EVENT_READ_MODEL_DIAGNOSTIC_SEVERITY[diagnostic.code] === 'hard'; |
| } |
| |
| export function isContinuationStartRuntimeEvent(event: RuntimeEvent): boolean { |
| return ( |
| event.actions?.stateDelta?.continuationStart === true || |
| event.actions?.continuationStart !== undefined |
| ); |
| } |
| |
| export function projectRuntimeEventCoordinationReceipt( |
| event: RuntimeEvent, |
| ): WorkHubCoordinationActionMessage | undefined { |
| if (!event.actions?.coordination) return undefined; |
| return { |
| type: 'workhub_coordination', |
| kind: 'action_receipt', |
| schemaVersion: 1, |
| id: event.id, |
| turnId: event.turnId, |
| ts: event.ts, |
| receipt: event.actions.coordination, |
| }; |
| } |
| |
| /** |
| * Whether the event can affect the StoredMessage projection or the state needed |
| * to construct one. Pure control-plane facts are intentionally absent so a |
| * transcript reader can stream past them without retaining the whole ledger. |
| */ |
| export function affectsRuntimeEventStoredMessageProjection(event: RuntimeEvent): boolean { |
| return ( |
| event.actions?.coordination !== undefined || |
| event.content !== undefined || |
| isTerminalRuntimeEvent(event) || |
| event.actions?.permissionRequest !== undefined || |
| event.actions?.permissionDecision !== undefined || |
| event.actions?.permissionAnswerAccepted !== undefined || |
| event.actions?.tokenUsage !== undefined |
| ); |
| } |
| |
| /** |
| * Codes that mean the projection did not claim an event, at either severity. |
| * |
| * Severity decides whether a session still opens; this decides whether the |
| * projection has a coverage gap. The projection-coverage contract asserts on |
| * this set, so softening an event's severity never softens the contract. |
| */ |
| const UNCLAIMED_RUNTIME_EVENT_DIAGNOSTIC_CODES: readonly RuntimeEventReadModelDiagnosticCode[] = [ |
| 'unsupported_event', |
| 'unclaimed_control_fact', |
| ]; |
| |
| export function isUnclaimedRuntimeEventDiagnostic(diagnostic: { |
| code: RuntimeEventReadModelDiagnosticCode; |
| }): boolean { |
| return UNCLAIMED_RUNTIME_EVENT_DIAGNOSTIC_CODES.includes(diagnostic.code); |
| } |
| |
| export interface RuntimeEventReadModelDiagnostic { |
| code: RuntimeEventReadModelDiagnosticCode; |
| eventId?: string; |
| runId?: string; |
| turnId?: string; |
| message: string; |
| detail?: unknown; |
| } |
| |
| export interface RuntimeEventReadModelProjection { |
| messages: StoredMessage[]; |
| diagnostics: RuntimeEventReadModelDiagnostic[]; |
| /** The id of the event each message was projected from, by position. */ |
| sourceEventIds: string[]; |
| } |
| |
| export interface ProjectRuntimeEventsToStoredMessagesOptions { |
| invocations: |
| | readonly RuntimeInvocationRecord[] |
| | Readonly<Record<string, RuntimeInvocationRecord>>; |
| canonicalPermissionOutcomes?: ReadonlyMap<string, CanonicalPermissionOutcomeRecord>; |
| active?: boolean; |
| projectToolResult?: (event: RuntimeEvent, decoded: ToolResultContent) => ToolResultContent; |
| onMessage?: (message: StoredMessage, sourceEventId: string) => void; |
| } |
| |
| export interface RuntimeEventStoredMessageProjector { |
| push(event: RuntimeEvent): void; |
| finish(): RuntimeEventReadModelProjection; |
| readonly permissionRequestIds: readonly string[]; |
| } |
| |
| /** |
| * Keep a completed local terminal result useful but bounded in transcript views. |
| * The durable RuntimeEvent remains untouched and the marker names its retained |
| * result, so readers can fetch the full output without repeating the command. |
| */ |
| export function projectTranscriptToolResult( |
| event: RuntimeEvent, |
| content: ToolResultContent, |
| ): ToolResultContent { |
| if ( |
| content.kind !== 'terminal' || |
| content.output.mode !== 'pipes' || |
| !hasLocalTerminalModelProjection(event) |
| ) { |
| return content; |
| } |
| const recoveryHint = `Read ${JSON.stringify({ |
| path: `maka://runtime/tool-results/${encodeURIComponent(event.id)}`, |
| })} for the retained output; follow next to continue.`; |
| const options = { maxBytes: 1024, maxLines: 20, direction: 'tail' as const, recoveryHint }; |
| const stdout = truncateToolOutput(content.output.stdout, options); |
| const stderr = truncateToolOutput(content.output.stderr, options); |
| if (!stdout.truncated && !stderr.truncated) return content; |
| return { |
| ...content, |
| output: { |
| ...content.output, |
| stdout: stdout.content, |
| stderr: stderr.content, |
| stdoutTruncated: content.output.stdoutTruncated || stdout.truncated, |
| stderrTruncated: content.output.stderrTruncated || stderr.truncated, |
| }, |
| }; |
| } |
| |
| function hasLocalTerminalModelProjection(event: RuntimeEvent): boolean { |
| const content = event.content; |
| if (content?.kind !== 'function_response' || content.providerExecuted) return false; |
| const projection = content.modelProjection; |
| return ( |
| projection?.kind === 'json' && |
| projection.value !== null && |
| typeof projection.value === 'object' && |
| 'kind' in projection.value && |
| projection.value.kind === 'terminal' |
| ); |
| } |
| |
| export interface ArchivedToolResultReadModelStatus { |
| runtimeEventId: string; |
| status: Extract<ToolResultContent, { kind: 'archived_tool_result' }>['status']; |
| } |
| |
| export interface RuntimeEventTerminalFact { |
| runId: string; |
| turnId: string; |
| runStatus: 'completed' | 'failed' | 'cancelled'; |
| turnStatus: 'completed' | 'failed' | 'aborted'; |
| terminalEvent: RuntimeEvent; |
| failureClass?: string; |
| abortSource?: string; |
| diagnostics: RuntimeEventReadModelDiagnostic[]; |
| } |
| |
| export interface RuntimeEventTerminalFactResult { |
| fact?: RuntimeEventTerminalFact; |
| diagnostics: RuntimeEventReadModelDiagnostic[]; |
| } |
| |
| interface PermissionRequestProjectionMetadata { |
| requestId: string; |
| toolUseId: string; |
| toolName: string; |
| sessionId: string; |
| runId: string; |
| turnId: string; |
| hint?: string; |
| } |
| |
| interface ProjectionState { |
| invocations: Map<string, RuntimeInvocationRecord>; |
| diagnostics: RuntimeEventReadModelDiagnostic[]; |
| toolNameByUseId: Map<string, string>; |
| permissionRequestById: Map<string, PermissionRequestProjectionMetadata>; |
| /** |
| * Thinking awaiting its assistant text row, keyed by the step message id |
| * (function of the event's providerEventId / storedMessageId — the same id the |
| * step's assistant row gets). Per-step turns have several entries per turn, so |
| * keying by message id (not turn) attaches each step's reasoning to its own row. |
| */ |
| thinkingByMessageId: Map<string, PendingThinking[]>; |
| contentOrderByMessageId: Map<string, AssistantStepContentKind[]>; |
| projectToolResult?: (event: RuntimeEvent, decoded: ToolResultContent) => ToolResultContent; |
| sourceOrder: number; |
| } |
| |
| interface PendingThinking { |
| event: RuntimeEvent; |
| messageId: string; |
| text: string; |
| signature?: string; |
| providerOptions?: Record<string, unknown>; |
| sourceOrder: number; |
| } |
| |
| export function projectRuntimeEventsToStoredMessages( |
| events: readonly RuntimeEvent[], |
| options: ProjectRuntimeEventsToStoredMessagesOptions, |
| ): RuntimeEventReadModelProjection { |
| const projector = createRuntimeEventStoredMessageProjector(options); |
| for (const event of events) projector.push(event); |
| return projector.finish(); |
| } |
| |
| export function createRuntimeEventStoredMessageProjector( |
| options: ProjectRuntimeEventsToStoredMessagesOptions, |
| ): RuntimeEventStoredMessageProjector { |
| const state: ProjectionState = { |
| invocations: normalizeInvocations(options.invocations), |
| diagnostics: [], |
| toolNameByUseId: new Map(), |
| permissionRequestById: new Map(), |
| thinkingByMessageId: new Map(), |
| contentOrderByMessageId: new Map(), |
| ...(options.projectToolResult ? { projectToolResult: options.projectToolResult } : {}), |
| sourceOrder: -1, |
| }; |
| const messages: StoredMessage[] = []; |
| /** |
| * Which event each message came out of, by position. |
| * |
| * A message belongs to the event being read when it was appended: nothing |
| * rewrites an earlier message, so the rows that appear while one event is |
| * handled are exactly that event's rows. A durable reader numbers its pages |
| * from this, which is why it is recorded here rather than rediscovered. |
| */ |
| const messageSources: Array<{ |
| eventId: string; |
| sourceOrder: number; |
| sourcePosition: number; |
| emittedOrder: number; |
| }> = []; |
| const diagnosticSources: Array<{ |
| sourceOrder: number; |
| sourcePosition: number; |
| emittedOrder: number; |
| }> = []; |
| const deferredPermissionAcceptances: Array<{ |
| event: RuntimeEvent; |
| ledgerRequest: PermissionRequestProjectionMetadata | undefined; |
| sourceOrder: number; |
| messagePosition: number; |
| diagnosticPosition: number; |
| }> = []; |
| const permissionRequestIds = new Set<string>(); |
| let nextSourceOrder = 0; |
| let finished: RuntimeEventReadModelProjection | undefined; |
| const attributeEmitted = ( |
| event: RuntimeEvent, |
| sourceOrder: number, |
| firstSourcePosition: number, |
| ): number => { |
| let sourcePosition = firstSourcePosition; |
| while (messageSources.length < messages.length) { |
| const message = messages[messageSources.length]!; |
| messageSources.push({ |
| eventId: event.id, |
| sourceOrder, |
| sourcePosition, |
| emittedOrder: messageSources.length, |
| }); |
| sourcePosition += 1; |
| options.onMessage?.(message, event.id); |
| } |
| return sourcePosition; |
| }; |
| const attributeDiagnostics = (sourceOrder: number, firstSourcePosition: number): number => { |
| let sourcePosition = firstSourcePosition; |
| while (diagnosticSources.length < state.diagnostics.length) { |
| diagnosticSources.push({ |
| sourceOrder, |
| sourcePosition, |
| emittedOrder: diagnosticSources.length, |
| }); |
| sourcePosition += 1; |
| } |
| return sourcePosition; |
| }; |
| |
| const projectEvent = (event: RuntimeEvent, sourceOrder: number): void => { |
| let nextMessagePosition = 0; |
| let nextDiagnosticPosition = 0; |
| state.sourceOrder = sourceOrder; |
| recordStepContentOrder(event, state); |
| if (isPartialRuntimeEvent(event)) { |
| diagnostic(state, event, 'partial_skipped', 'partial RuntimeEvent skipped'); |
| attributeDiagnostics(sourceOrder, nextDiagnosticPosition); |
| return; |
| } |
| |
| let projected = false; |
| const content = event.content; |
| if (content) { |
| switch (content.kind) { |
| case 'text': |
| projected = projectText(event, state, messages) || projected; |
| break; |
| case 'function_call': |
| projected = projectFunctionCall(event, state, messages) || projected; |
| break; |
| case 'function_response': |
| projected = projectFunctionResponse(event, state, messages) || projected; |
| break; |
| case 'thinking': |
| projected = projectThinking(event, state, messages) || projected; |
| break; |
| case 'system_note': |
| projected = projectSystemNote(event, state, messages) || projected; |
| break; |
| case 'invocation_opened': |
| // The opening fact records route, configuration and lineage once per |
| // invocation. Every reader joins it by invocationId; it has no chat row. |
| projected = true; |
| break; |
| case 'error': |
| if (!isTerminalRuntimeEvent(event)) { |
| diagnostic( |
| state, |
| event, |
| 'unsupported_event', |
| 'non-terminal error content has no safe legacy read-model row', |
| ); |
| } |
| break; |
| } |
| } |
| |
| const coordinationReceipt = projectRuntimeEventCoordinationReceipt(event); |
| if (coordinationReceipt) { |
| messages.push(coordinationReceipt); |
| projected = true; |
| } |
| |
| if (event.actions?.permissionRequest) { |
| const request = event.actions.permissionRequest; |
| state.permissionRequestById.set(request.requestId, { |
| requestId: request.requestId, |
| toolUseId: request.toolUseId, |
| toolName: request.toolName, |
| sessionId: event.sessionId, |
| runId: event.runId, |
| turnId: event.turnId, |
| ...(request.hint !== undefined ? { hint: request.hint } : {}), |
| }); |
| state.toolNameByUseId.set(request.toolUseId, request.toolName); |
| projected = true; |
| } |
| |
| if (event.actions?.userQuestionRequest) { |
| // The matching function_call/function_response own the legacy rows; |
| // this request is live interaction state only. |
| projected = true; |
| } |
| |
| if (event.actions?.userQuestionAnswerAccepted) { |
| // InteractionStore owns the canonical answer. This Run-local audit fact |
| // intentionally has no legacy chat row. |
| projected = true; |
| } |
| |
| if (event.actions?.formRequest) { |
| // The matching function_call/function_response own the legacy rows; |
| // this request is live interaction state only. |
| projected = true; |
| } |
| |
| if (event.actions?.formAnswerAccepted) { |
| // InteractionStore owns the canonical result. This Run-local audit fact |
| // intentionally has no legacy chat row. |
| projected = true; |
| } |
| |
| if (event.actions?.permissionAnswerAccepted) { |
| nextMessagePosition = attributeEmitted(event, sourceOrder, nextMessagePosition); |
| nextDiagnosticPosition = attributeDiagnostics(sourceOrder, nextDiagnosticPosition); |
| const requestId = event.actions.permissionAnswerAccepted.requestId; |
| permissionRequestIds.add(requestId); |
| deferredPermissionAcceptances.push({ |
| event: permissionAcceptanceProjectionEvent(event), |
| ledgerRequest: state.permissionRequestById.get(requestId), |
| sourceOrder, |
| messagePosition: nextMessagePosition++, |
| diagnosticPosition: nextDiagnosticPosition++, |
| }); |
| projected = true; |
| } |
| |
| if (event.actions?.permissionClosureAccepted) { |
| // The canonical closure is already represented by this identity-only |
| // RuntimeEvent; unlike an answer it has no legacy permission-decision row. |
| projected = true; |
| } |
| |
| if (event.actions?.toolDispatch) { |
| // Dispatch is a canonical recovery fact with no legacy chat row. It is |
| // consumed by RecoveryResolver, but must remain invisible to messages. |
| projected = true; |
| } |
| |
| if (event.actions?.toolRecovery) { |
| // Recovery observations and decisions are canonical audit facts. The |
| // matching function_call/function_response own any provider-visible row. |
| projected = true; |
| } |
| |
| if (event.actions?.workspaceFact) { |
| // Workspace epoch/version facts belong to the store-owned control-plane |
| // stream. They are canonical recovery inputs, never chat messages. |
| projected = true; |
| } |
| |
| if (event.actions?.managedMutationTerminal) { |
| // The matching function_response owns the provider-visible row. This |
| // action only proves that the managed reservation reached a no-effect |
| // terminal through its dedicated atomic writer. |
| projected = true; |
| } |
| |
| if (event.actions?.artifactDelta) { |
| // Artifact counters are storage bookkeeping. The tool result that owns the |
| // artifact owns its row; this delta has none of its own. |
| projected = true; |
| } |
| |
| if (event.actions?.transferToAgent !== undefined) { |
| // A hand-off is control routing. The receiving agent's own events own |
| // every provider-visible row the transfer leads to. |
| projected = true; |
| } |
| |
| if (event.actions?.handoffPause) { |
| // Physical pause is not a logical Turn outcome or a chat message. |
| projected = true; |
| } |
| |
| if (event.actions?.runtimeProtocol) { |
| // The protocol marker records which runtime contracts were live from a |
| // run's first event. RecoveryResolver reads it; it has no chat row. |
| projected = true; |
| } |
| |
| if (isContinuationStartRuntimeEvent(event)) { |
| // Continuation start is a canonical lineage/recovery fact with no |
| // legacy chat row. Its following model events own the visible output. |
| projected = true; |
| } |
| |
| if (isSandboxBoundaryStateDelta(event)) { |
| // The session sandbox boundary owns enforcement and its own durable |
| // revisions. These are canonical control/audit facts, and the tool call |
| // and response around them own every provider-visible row. |
| projected = true; |
| } |
| |
| if (isPlanProposalStateDelta(event)) { |
| // Plan proposals render from PlanStore as approval cards. This event is |
| // still a canonical runtime fact, but intentionally has no legacy chat row. |
| projected = true; |
| } |
| |
| if (event.actions?.permissionDecision) { |
| projected = projectPermissionDecision(event, state, messages) || projected; |
| } |
| |
| if (event.actions?.tokenUsage) { |
| projected = projectTokenUsage(event, state, messages) || projected; |
| } |
| |
| if (isTerminalRuntimeEvent(event) && !event.actions?.handoffPause) { |
| projected = projectTerminalTurnState(event, state, messages) || projected; |
| } |
| |
| if (!projected) { |
| // Content is the only payload an unclaimed shape could still have owed a |
| // row, so its absence is what makes degrading safe here — not a promise |
| // that actions never produce rows (permissionDecision, tokenUsage and the |
| // terminal fact all do). What holds that up is claim coverage: every |
| // action field a reader can meet is claimed above, proven by the |
| // projection-coverage contract, so nothing with a row reaches this branch. |
| if (event.content === undefined) { |
| diagnostic( |
| state, |
| event, |
| 'unclaimed_control_fact', |
| 'control-only RuntimeEvent is not claimed by the legacy read-model projection', |
| ); |
| } else { |
| diagnostic( |
| state, |
| event, |
| 'unsupported_event', |
| 'RuntimeEvent shape is not supported by the legacy read-model projection', |
| ); |
| } |
| } |
| attributeEmitted(event, sourceOrder, nextMessagePosition); |
| attributeDiagnostics(sourceOrder, nextDiagnosticPosition); |
| }; |
| |
| const push = (event: RuntimeEvent): void => { |
| if (finished) throw new Error('RuntimeEvent StoredMessage projector is already finished'); |
| const sourceOrder = nextSourceOrder++; |
| projectEvent(options.active ? settledPresentationEvent(event) : event, sourceOrder); |
| }; |
| |
| const finish = (): RuntimeEventReadModelProjection => { |
| if (finished) return finished; |
| |
| for (const acceptance of deferredPermissionAcceptances) { |
| const before = messages.length; |
| projectCanonicalPermissionOutcome( |
| acceptance.event, |
| state, |
| messages, |
| options.canonicalPermissionOutcomes, |
| acceptance.ledgerRequest, |
| ); |
| if (messages.length > before) { |
| attributeEmitted(acceptance.event, acceptance.sourceOrder, acceptance.messagePosition); |
| } |
| attributeDiagnostics(acceptance.sourceOrder, acceptance.diagnosticPosition); |
| } |
| |
| if (options.active) { |
| for (const pendingItems of [...state.thinkingByMessageId.values()]) { |
| const pending = pendingItems.at(-1); |
| if (pending) projectEvent(emptyAssistantText(pending.event), pending.sourceOrder); |
| } |
| } |
| |
| const orderedDiagnostics = state.diagnostics |
| .map((diagnostic, index) => ({ diagnostic, source: diagnosticSources[index]! })) |
| .sort( |
| (left, right) => |
| left.source.sourceOrder - right.source.sourceOrder || |
| left.source.sourcePosition - right.source.sourcePosition || |
| left.source.emittedOrder - right.source.emittedOrder, |
| ); |
| state.diagnostics.splice( |
| 0, |
| state.diagnostics.length, |
| ...orderedDiagnostics.map(({ diagnostic }) => diagnostic), |
| ); |
| |
| if (!options.active) { |
| for (const pendingItems of state.thinkingByMessageId.values()) { |
| for (const pending of pendingItems) { |
| diagnostic( |
| state, |
| pending.event, |
| 'unsupported_event', |
| 'thinking content has no assistant text row with a matching message id', |
| ); |
| } |
| } |
| } |
| |
| const ordered = messages |
| .map((message, index) => ({ message, source: messageSources[index]! })) |
| .sort( |
| (left, right) => |
| left.source.sourceOrder - right.source.sourceOrder || |
| left.source.sourcePosition - right.source.sourcePosition || |
| left.source.emittedOrder - right.source.emittedOrder, |
| ); |
| finished = { |
| messages: ordered.map(({ message }) => message), |
| diagnostics: state.diagnostics, |
| sourceEventIds: ordered.map(({ source }) => source.eventId), |
| }; |
| deferredPermissionAcceptances.length = 0; |
| messageSources.length = 0; |
| diagnosticSources.length = 0; |
| state.invocations.clear(); |
| state.toolNameByUseId.clear(); |
| state.permissionRequestById.clear(); |
| state.thinkingByMessageId.clear(); |
| state.contentOrderByMessageId.clear(); |
| state.projectToolResult = undefined; |
| return finished; |
| }; |
| |
| return { |
| push, |
| finish, |
| get permissionRequestIds(): readonly string[] { |
| return [...permissionRequestIds]; |
| }, |
| }; |
| } |
| |
| /** |
| * A running invocation's events as the transcript should show them right now. |
| * |
| * Two things separate a live run from a finished one. Its last text or thinking |
| * event is still arriving, so it is presented as settled rather than withheld; |
| * and a step that has only thought so far has no assistant row to hang that |
| * thinking on, so an empty one is opened for it. Neither changes the ledger: |
| * both are how the same events read before the run ends. |
| */ |
| export function activePresentationRuntimeEvents(events: readonly RuntimeEvent[]): RuntimeEvent[] { |
| const textMessages = new Set<string>(); |
| const lastThinkingByMessage = new Map<string, RuntimeEvent>(); |
| |
| for (const event of events) { |
| const content = event.content; |
| if (event.role !== 'model' || (content?.kind !== 'text' && content?.kind !== 'thinking')) { |
| continue; |
| } |
| const messageKey = activeMessageKey(event); |
| if (content.kind === 'text') textMessages.add(messageKey); |
| else lastThinkingByMessage.set(messageKey, event); |
| } |
| |
| const syntheticAfter = new Map<RuntimeEvent, RuntimeEvent[]>(); |
| for (const [messageKey, thinking] of lastThinkingByMessage) { |
| if (textMessages.has(messageKey)) continue; |
| const existing = syntheticAfter.get(thinking) ?? []; |
| existing.push(emptyAssistantText(thinking)); |
| syntheticAfter.set(thinking, existing); |
| } |
| |
| const presented: RuntimeEvent[] = []; |
| for (const event of events) { |
| presented.push(settledPresentationEvent(event)); |
| presented.push(...(syntheticAfter.get(event) ?? [])); |
| } |
| return presented; |
| } |
| |
| function activeMessageKey(event: RuntimeEvent): string { |
| const messageId = event.refs?.providerEventId ?? event.refs?.storedMessageId ?? event.id; |
| return `${event.runId}\0${messageId}`; |
| } |
| |
| function settledPresentationEvent(event: RuntimeEvent): RuntimeEvent { |
| const content = event.content; |
| return event.partial && |
| event.role === 'model' && |
| (content?.kind === 'text' || content?.kind === 'thinking') |
| ? { ...event, partial: false } |
| : event; |
| } |
| |
| function emptyAssistantText(thinking: RuntimeEvent): RuntimeEvent { |
| return { |
| ...thinking, |
| id: `${thinking.id}:active-transcript-empty-text`, |
| partial: false, |
| content: { kind: 'text', text: '' }, |
| }; |
| } |
| |
| export function projectRuntimeEventsToStoredMessagesWithArchiveStatuses( |
| events: readonly RuntimeEvent[], |
| options: ProjectRuntimeEventsToStoredMessagesOptions & { |
| archiveStatuses: |
| | readonly ArchivedToolResultReadModelStatus[] |
| | Readonly<Record<string, ArchivedToolResultReadModelStatus['status']>>; |
| }, |
| ): RuntimeEventReadModelProjection { |
| return projectRuntimeEventsToStoredMessages( |
| applyArchivedToolResultReadModelStatuses(events, options.archiveStatuses), |
| options, |
| ); |
| } |
| |
| export function applyArchivedToolResultReadModelStatuses( |
| events: readonly RuntimeEvent[], |
| archiveStatuses: |
| | readonly ArchivedToolResultReadModelStatus[] |
| | Readonly<Record<string, ArchivedToolResultReadModelStatus['status']>>, |
| ): RuntimeEvent[] { |
| const statuses = normalizeArchiveStatuses(archiveStatuses); |
| if (statuses.size === 0) return [...events]; |
| return events.map((event) => { |
| const status = statuses.get(event.id); |
| if (!status || event.content?.kind !== 'function_response') return event; |
| if (!isArchivedToolResultPlaceholder(event.content.result)) return event; |
| const placeholder = event.content.result; |
| return { |
| ...event, |
| content: { |
| ...event.content, |
| result: { |
| kind: 'archived_tool_result', |
| status, |
| runtimeEventId: placeholder.runtimeEventId, |
| toolCallId: placeholder.toolCallId, |
| toolName: placeholder.toolName, |
| artifactId: placeholder.artifactId, |
| ...(placeholder.rewriteVersion === 2 ? { resourceRef: placeholder.resourceRef } : {}), |
| bodySha256: placeholder.bodySha256, |
| originalEstimatedTokens: placeholder.originalEstimatedTokens, |
| originalBytes: placeholder.originalBytes, |
| rewriteVersion: placeholder.rewriteVersion, |
| reason: placeholder.reason, |
| } satisfies ToolResultContent, |
| }, |
| }; |
| }); |
| } |
| |
| export function classifyRuntimeEventTerminalFact( |
| invocation: Pick<RuntimeInvocationRecord, 'sessionId' | 'runId' | 'turnId'>, |
| events: readonly RuntimeEvent[], |
| ): RuntimeEventTerminalFactResult { |
| const diagnostics: RuntimeEventReadModelDiagnostic[] = []; |
| if (events.length === 0) { |
| diagnostics.push( |
| readModelDiagnostic('incomplete_event', 'runtime ledger has no readable RuntimeEvents', { |
| runId: invocation.runId, |
| turnId: invocation.turnId, |
| }), |
| ); |
| return { diagnostics }; |
| } |
| |
| const terminalSignals = events.filter( |
| (event) => |
| !isPartialRuntimeEvent(event) && |
| event.sessionId === invocation.sessionId && |
| event.runId === invocation.runId && |
| event.turnId === invocation.turnId && |
| isTerminalRuntimeEvent(event), |
| ); |
| |
| if (terminalSignals.length === 0) { |
| diagnostics.push( |
| readModelDiagnostic( |
| 'incomplete_event', |
| 'runtime ledger has no matching terminal RuntimeEvent', |
| { runId: invocation.runId, turnId: invocation.turnId }, |
| ), |
| ); |
| return { diagnostics }; |
| } |
| if (terminalSignals.length > 1) { |
| diagnostics.push( |
| readModelDiagnostic( |
| 'incomplete_event', |
| 'runtime ledger has multiple matching terminal RuntimeEvents', |
| { |
| runId: invocation.runId, |
| turnId: invocation.turnId, |
| eventIds: terminalSignals.map((event) => event.id), |
| }, |
| ), |
| ); |
| return { diagnostics }; |
| } |
| |
| const terminalEvent = terminalSignals[0]!; |
| if (!isTerminalRuntimeEventStatus(terminalEvent.status)) { |
| diagnostics.push( |
| readModelDiagnostic( |
| 'incomplete_event', |
| 'terminal RuntimeEvent requires a terminal status for recovery', |
| terminalEvent, |
| ), |
| ); |
| return { diagnostics }; |
| } |
| |
| if (terminalEvent.status === 'completed') { |
| const fact: RuntimeEventTerminalFact = { |
| runId: invocation.runId, |
| turnId: invocation.turnId, |
| runStatus: 'completed', |
| turnStatus: 'completed', |
| terminalEvent, |
| diagnostics, |
| }; |
| return { fact, diagnostics }; |
| } |
| |
| // A terminal event is the run's ending, and it is immutable once written, so |
| // an omitted failure class or abort source is a detail nobody can ever supply |
| // afterwards. Withholding the fact over it would only leave the reader with a |
| // run that ended and no way to say so; the omission is worth a diagnostic, not |
| // a refusal. |
| if (terminalEvent.status === 'failed') { |
| const failureClass = failureClassFromRuntimeEvent(terminalEvent); |
| if (!failureClass) { |
| diagnostics.push( |
| readModelDiagnostic( |
| 'incomplete_event', |
| 'failed terminal RuntimeEvent states no failure class', |
| terminalEvent, |
| ), |
| ); |
| } |
| const fact: RuntimeEventTerminalFact = { |
| runId: invocation.runId, |
| turnId: invocation.turnId, |
| runStatus: 'failed', |
| turnStatus: 'failed', |
| terminalEvent, |
| failureClass: failureClass ?? 'unknown', |
| diagnostics, |
| }; |
| return { fact, diagnostics }; |
| } |
| |
| const abortSource = abortSourceFromRuntime(terminalEvent); |
| if (!abortSource) { |
| diagnostics.push( |
| readModelDiagnostic( |
| 'incomplete_event', |
| 'aborted terminal RuntimeEvent states no abort source', |
| terminalEvent, |
| ), |
| ); |
| } |
| const fact: RuntimeEventTerminalFact = { |
| runId: invocation.runId, |
| turnId: invocation.turnId, |
| runStatus: 'cancelled', |
| turnStatus: 'aborted', |
| terminalEvent, |
| abortSource: abortSource ?? 'unknown', |
| diagnostics, |
| }; |
| return { fact, diagnostics }; |
| } |
| |
| function projectText( |
| event: RuntimeEvent, |
| state: ProjectionState, |
| messages: StoredMessage[], |
| ): boolean { |
| if (event.content?.kind !== 'text') return false; |
| if (event.role === 'user') { |
| const message = projectRuntimeEventUserMessage(event, stableMessageId(event, state, 'user')); |
| if (!message) return false; |
| messages.push(message); |
| return true; |
| } |
| |
| if (event.role === 'model') { |
| const invocation = state.invocations.get(event.runId); |
| if (!invocation?.opening.route.modelId) { |
| diagnostic( |
| state, |
| event, |
| 'incomplete_event', |
| 'model text RuntimeEvent requires the opening fact of its invocation', |
| ); |
| return false; |
| } |
| const assistantId = stableMessageId(event, state, 'assistant'); |
| const contentOrder = nonCanonicalContentOrder(state.contentOrderByMessageId.get(assistantId)); |
| messages.push({ |
| type: 'assistant', |
| id: assistantId, |
| turnId: event.turnId, |
| ts: event.ts, |
| text: event.content.text, |
| ...(event.content.providerOptions !== undefined |
| ? { providerOptions: structuredClone(event.content.providerOptions) } |
| : {}), |
| ...(contentOrder ? { contentOrder } : {}), |
| modelId: invocation.opening.route.modelId, |
| }); |
| attachPendingThinking(event, state, messages, assistantId); |
| return true; |
| } |
| |
| diagnostic( |
| state, |
| event, |
| 'unsupported_event', |
| `text content with role ${event.role} is not projected`, |
| ); |
| return false; |
| } |
| |
| export function projectRuntimeEventUserMessage( |
| event: RuntimeEvent, |
| messageId: string, |
| ): Extract<StoredMessage, { type: 'user' }> | undefined { |
| if (event.role !== 'user' || event.content?.kind !== 'text') return undefined; |
| return { |
| type: 'user', |
| id: messageId, |
| turnId: event.turnId, |
| ts: event.ts, |
| ...normalizeMessageContent(event.content), |
| ...(event.content.origin !== undefined ? { origin: event.content.origin } : {}), |
| ...(event.content.steering === true ? { steeringEventId: event.id } : {}), |
| }; |
| } |
| |
| export function nonCanonicalContentOrder( |
| order: readonly AssistantStepContentKind[] | undefined, |
| ): AssistantStepContentKind[] | undefined { |
| if (!order?.length) return undefined; |
| const present = new Set(order); |
| const canonical = (['thinking', 'text', 'tools'] as const).filter((kind) => present.has(kind)); |
| return order.every((kind, index) => kind === canonical[index]) ? undefined : [...order]; |
| } |
| |
| function recordStepContentOrder(event: RuntimeEvent, state: ProjectionState): void { |
| const content = event.content; |
| let messageId: string | undefined; |
| let kind: AssistantStepContentKind | undefined; |
| if (event.role === 'model' && content?.kind === 'text') { |
| messageId = event.refs?.providerEventId ?? event.refs?.storedMessageId ?? event.id; |
| kind = 'text'; |
| } else if (event.role === 'model' && content?.kind === 'thinking') { |
| messageId = event.refs?.providerEventId ?? event.refs?.storedMessageId ?? event.id; |
| kind = 'thinking'; |
| } else if (event.role === 'model' && content?.kind === 'function_call' && event.refs?.stepId) { |
| messageId = event.refs.stepId; |
| kind = 'tools'; |
| } |
| if (!messageId || !kind) return; |
| const order = state.contentOrderByMessageId.get(messageId) ?? []; |
| if (!order.includes(kind)) state.contentOrderByMessageId.set(messageId, [...order, kind]); |
| } |
| |
| function normalizeArchiveStatuses( |
| archiveStatuses: |
| | readonly ArchivedToolResultReadModelStatus[] |
| | Readonly<Record<string, ArchivedToolResultReadModelStatus['status']>>, |
| ): Map<string, ArchivedToolResultReadModelStatus['status']> { |
| const map = new Map<string, ArchivedToolResultReadModelStatus['status']>(); |
| if (Array.isArray(archiveStatuses)) { |
| for (const item of archiveStatuses) { |
| map.set(item.runtimeEventId, item.status); |
| } |
| return map; |
| } |
| for (const [runtimeEventId, status] of Object.entries(archiveStatuses)) { |
| map.set(runtimeEventId, status); |
| } |
| return map; |
| } |
| |
| function projectThinking( |
| event: RuntimeEvent, |
| state: ProjectionState, |
| messages: StoredMessage[], |
| ): boolean { |
| if (event.content?.kind !== 'thinking') return false; |
| const messageId = thinkingMessageId(event); |
| const pending: PendingThinking = { |
| event, |
| messageId, |
| text: event.content.text, |
| sourceOrder: state.sourceOrder, |
| ...(event.content.signature !== undefined ? { signature: event.content.signature } : {}), |
| ...(event.content.providerOptions !== undefined |
| ? { providerOptions: structuredClone(event.content.providerOptions) } |
| : {}), |
| }; |
| // The step's assistant text row lands after its thinking in ledger order, so |
| // attach eagerly if it already exists (older ordering), else park by message id |
| // for projectText's attachPendingThinking to claim. |
| if (attachThinkingToAssistant(event, pending, messages)) return true; |
| const pendingItems = state.thinkingByMessageId.get(messageId) ?? []; |
| pendingItems.push(pending); |
| state.thinkingByMessageId.set(messageId, pendingItems); |
| return true; |
| } |
| |
| function projectFunctionCall( |
| event: RuntimeEvent, |
| state: ProjectionState, |
| messages: StoredMessage[], |
| ): boolean { |
| if (event.content?.kind !== 'function_call') return false; |
| const toolUseId = toolUseIdFor(event); |
| if (!toolUseId) { |
| diagnostic( |
| state, |
| event, |
| 'incomplete_event', |
| 'function_call RuntimeEvent requires content.id or refs.toolCallId', |
| ); |
| return false; |
| } |
| if (event.content.id !== toolUseId) { |
| diagnostic( |
| state, |
| event, |
| 'tool_use_id_mismatch', |
| 'function_call content.id differs from refs.toolCallId', |
| { |
| contentId: event.content.id, |
| refToolCallId: event.refs?.toolCallId, |
| }, |
| ); |
| } |
| state.toolNameByUseId.set(toolUseId, event.content.name); |
| messages.push({ |
| type: 'tool_call', |
| id: toolUseId, |
| turnId: event.turnId, |
| ts: event.ts, |
| toolName: event.content.name, |
| ...(toolActivityKindStateDelta(event) !== undefined |
| ? { activityKind: toolActivityKindStateDelta(event) } |
| : {}), |
| ...(stringStateDelta(event, 'displayName') !== undefined |
| ? { displayName: stringStateDelta(event, 'displayName') } |
| : {}), |
| ...(stringStateDelta(event, 'intent') !== undefined |
| ? { intent: stringStateDelta(event, 'intent') } |
| : {}), |
| // Carry the step pairing through the projection: without it, sessions |
| // rebuilt from the runtime event log lose the tool↔step association and |
| // the UI timeline falls back to legacy tools-before-text ordering. |
| ...(event.refs?.stepId ? { stepId: event.refs.stepId } : {}), |
| ...toolActivityIdentity(event), |
| args: event.content.args, |
| ...(event.content.providerOptions !== undefined |
| ? { providerOptions: structuredClone(event.content.providerOptions) } |
| : {}), |
| ...(event.content.providerExecuted !== undefined |
| ? { providerExecuted: event.content.providerExecuted } |
| : {}), |
| }); |
| return true; |
| } |
| |
| function projectFunctionResponse( |
| event: RuntimeEvent, |
| state: ProjectionState, |
| messages: StoredMessage[], |
| ): boolean { |
| if (event.content?.kind !== 'function_response') return false; |
| const toolUseId = toolUseIdFor(event); |
| if (!toolUseId) { |
| diagnostic( |
| state, |
| event, |
| 'incomplete_event', |
| 'function_response RuntimeEvent requires content.id or refs.toolCallId', |
| ); |
| return false; |
| } |
| if (event.content.id !== toolUseId) { |
| diagnostic( |
| state, |
| event, |
| 'tool_use_id_mismatch', |
| 'function_response content.id differs from refs.toolCallId', |
| { |
| contentId: event.content.id, |
| refToolCallId: event.refs?.toolCallId, |
| }, |
| ); |
| } |
| const legacyPlanResult = isLegacyPlanToolResult(event.content.result) |
| ? { kind: 'json' as const, value: event.content.result } |
| : undefined; |
| const compatibleResult = legacyPlanResult ?? event.content.result; |
| const archivedPlaceholder = isArchivedToolResultPlaceholder(compatibleResult) |
| ? compatibleResult |
| : undefined; |
| let normalizedResult: ToolResultContent | undefined; |
| if (!archivedPlaceholder) { |
| try { |
| normalizedResult = decodePersistedToolResultContent( |
| markPersisted<ToolResultContent>(compatibleResult), |
| ); |
| } catch (error) { |
| diagnostic( |
| state, |
| event, |
| 'incomplete_event', |
| error instanceof Error && error.message === 'Invalid shell tool result content' |
| ? 'function_response contains an invalid shell tool result' |
| : 'function_response result is not a supported ToolResultContent', |
| ); |
| return false; |
| } |
| } |
| if (archivedPlaceholder) { |
| diagnostic( |
| state, |
| event, |
| 'archived_tool_result_placeholder', |
| 'function_response result is archived and not loaded in read model', |
| { |
| artifactId: archivedPlaceholder.artifactId, |
| runtimeEventId: archivedPlaceholder.runtimeEventId, |
| toolCallId: archivedPlaceholder.toolCallId, |
| toolName: archivedPlaceholder.toolName, |
| reason: archivedPlaceholder.reason, |
| rewriteVersion: archivedPlaceholder.rewriteVersion, |
| }, |
| ); |
| } |
| if (event.content.name) state.toolNameByUseId.set(toolUseId, event.content.name); |
| const decodedResult: ToolResultContent = archivedPlaceholder |
| ? { |
| kind: 'archived_tool_result', |
| status: 'not_loaded', |
| runtimeEventId: archivedPlaceholder.runtimeEventId, |
| toolCallId: archivedPlaceholder.toolCallId, |
| toolName: archivedPlaceholder.toolName, |
| artifactId: archivedPlaceholder.artifactId, |
| bodySha256: archivedPlaceholder.bodySha256, |
| ...(archivedPlaceholder.rewriteVersion === 2 |
| ? { resourceRef: archivedPlaceholder.resourceRef } |
| : {}), |
| originalEstimatedTokens: archivedPlaceholder.originalEstimatedTokens, |
| originalBytes: archivedPlaceholder.originalBytes, |
| rewriteVersion: archivedPlaceholder.rewriteVersion, |
| reason: archivedPlaceholder.reason, |
| } |
| : normalizedResult!; |
| const resultContent = state.projectToolResult |
| ? state.projectToolResult(event, decodedResult) |
| : decodedResult; |
| messages.push({ |
| type: 'tool_result', |
| id: stableMessageId(event, state, 'tool_result'), |
| turnId: event.turnId, |
| ts: event.ts, |
| toolUseId, |
| isError: event.content.isError === true, |
| content: resultContent, |
| ...(event.content.providerExecuted !== undefined |
| ? { providerExecuted: event.content.providerExecuted } |
| : {}), |
| ...(event.content.providerExecuted && event.content.providerOutput !== undefined |
| ? { providerOutput: structuredClone(event.content.providerOutput) } |
| : {}), |
| ...(numberStateDelta(event, 'durationMs') !== undefined |
| ? { durationMs: numberStateDelta(event, 'durationMs') } |
| : {}), |
| ...toolActivityIdentity(event), |
| }); |
| return true; |
| } |
| |
| function toolActivityIdentity(event: RuntimeEvent): { |
| origin?: 'provider' | 'code_mode'; |
| modelVisibility?: 'visible' | 'hidden'; |
| parentToolCallId?: string; |
| parentOperationId?: string; |
| } { |
| return { |
| ...(event.origin !== undefined ? { origin: event.origin } : {}), |
| ...(event.modelVisibility !== undefined ? { modelVisibility: event.modelVisibility } : {}), |
| ...(event.refs?.parentToolCallId !== undefined |
| ? { parentToolCallId: event.refs.parentToolCallId } |
| : {}), |
| ...(event.refs?.parentOperationId !== undefined |
| ? { parentOperationId: event.refs.parentOperationId } |
| : {}), |
| }; |
| } |
| |
| function projectPermissionDecision( |
| event: RuntimeEvent, |
| state: ProjectionState, |
| messages: StoredMessage[], |
| ): boolean { |
| const decision = event.actions?.permissionDecision; |
| if (!decision) return false; |
| const request = state.permissionRequestById.get(decision.requestId); |
| const toolUseId = event.refs?.toolCallId ?? request?.toolUseId; |
| if (!toolUseId) { |
| diagnostic( |
| state, |
| event, |
| 'incomplete_event', |
| 'permission decision requires refs.toolCallId or a paired permission request', |
| ); |
| return false; |
| } |
| if (request && request.toolUseId !== toolUseId) { |
| diagnostic( |
| state, |
| event, |
| 'tool_use_id_mismatch', |
| 'permission decision toolUseId does not match its paired permission request', |
| ); |
| return false; |
| } |
| const toolStateName = state.toolNameByUseId.get(toolUseId); |
| const toolName = decision.toolName ?? request?.toolName ?? toolStateName; |
| if (!toolName) { |
| diagnostic( |
| state, |
| event, |
| 'incomplete_event', |
| 'permission decision requires durable toolName or a paired permission request or tool call', |
| ); |
| return false; |
| } |
| if ( |
| (request?.toolName !== undefined && request.toolName !== toolName) || |
| (toolStateName !== undefined && toolStateName !== toolName) |
| ) { |
| diagnostic( |
| state, |
| event, |
| 'incomplete_event', |
| 'permission decision toolName does not match its paired request or tool call', |
| ); |
| return false; |
| } |
| // The prompt's own wording when the request survived, and the decision's copy |
| // of it when the decision is all that is left. |
| const hint = request?.hint ?? decision.hint; |
| messages.push({ |
| type: 'permission_decision', |
| id: decision.requestId, |
| turnId: event.turnId, |
| ts: event.ts, |
| toolUseId, |
| toolName, |
| decision: decision.decision, |
| ...(decision.rememberForTurn !== undefined |
| ? { rememberForTurn: decision.rememberForTurn } |
| : {}), |
| ...(decision.reviewer !== undefined ? { reviewer: decision.reviewer } : {}), |
| ...(decision.rationale !== undefined ? { rationale: decision.rationale } : {}), |
| ...(decision.riskLevel !== undefined ? { riskLevel: decision.riskLevel } : {}), |
| ...(hint !== undefined ? { hint } : {}), |
| }); |
| return true; |
| } |
| |
| function projectCanonicalPermissionOutcome( |
| event: RuntimeEvent, |
| state: ProjectionState, |
| messages: StoredMessage[], |
| outcomes: ReadonlyMap<string, CanonicalPermissionOutcomeRecord> | undefined, |
| ledgerRequest: PermissionRequestProjectionMetadata | undefined, |
| ): void { |
| const accepted = event.actions?.permissionAnswerAccepted; |
| if (!accepted) return; |
| const canonical = outcomes?.get(accepted.requestId); |
| const toolUseId = event.refs?.toolCallId; |
| if (!canonical || !toolUseId) { |
| diagnostic( |
| state, |
| event, |
| 'incomplete_event', |
| 'permission answer acceptance requires a canonical Interaction outcome', |
| { requestId: accepted.requestId }, |
| ); |
| return; |
| } |
| const outcome = canonical.outcome; |
| if ( |
| canonical.sessionId !== event.sessionId || |
| canonical.runId !== event.runId || |
| canonical.turnId !== event.turnId || |
| canonical.requestId !== accepted.requestId || |
| canonical.request.toolUseId !== toolUseId || |
| (ledgerRequest !== undefined && |
| (ledgerRequest.sessionId !== event.sessionId || |
| ledgerRequest.runId !== event.runId || |
| ledgerRequest.turnId !== event.turnId || |
| ledgerRequest.toolUseId !== toolUseId || |
| ledgerRequest.toolName !== canonical.request.prompt.toolName)) |
| ) { |
| diagnostic( |
| state, |
| event, |
| 'incomplete_event', |
| 'permission answer canonical outcome identity does not match its acceptance', |
| { requestId: accepted.requestId }, |
| ); |
| return; |
| } |
| messages.push({ |
| type: 'permission_decision', |
| id: accepted.requestId, |
| turnId: event.turnId, |
| ts: outcome.committedAt, |
| toolUseId, |
| toolName: canonical.request.prompt.toolName, |
| decision: outcome.decision, |
| rememberForTurn: outcome.rememberForTurn, |
| reviewer: outcome.reviewer, |
| ...(outcome.rationale !== undefined ? { rationale: outcome.rationale } : {}), |
| ...(outcome.riskLevel !== undefined ? { riskLevel: outcome.riskLevel } : {}), |
| ...(ledgerRequest?.hint !== undefined ? { hint: ledgerRequest.hint } : {}), |
| }); |
| } |
| |
| function permissionAcceptanceProjectionEvent(event: RuntimeEvent): RuntimeEvent { |
| return { |
| id: event.id, |
| sessionId: event.sessionId, |
| invocationId: event.invocationId, |
| runId: event.runId, |
| turnId: event.turnId, |
| ts: event.ts, |
| partial: event.partial, |
| role: event.role, |
| author: event.author, |
| actions: { permissionAnswerAccepted: event.actions!.permissionAnswerAccepted! }, |
| ...(event.refs?.toolCallId ? { refs: { toolCallId: event.refs.toolCallId } } : {}), |
| }; |
| } |
| |
| function projectTokenUsage( |
| event: RuntimeEvent, |
| state: ProjectionState, |
| messages: StoredMessage[], |
| ): boolean { |
| const usage = event.actions?.tokenUsage; |
| if (!usage) return false; |
| messages.push({ |
| type: 'token_usage', |
| id: stableMessageId(event, state, 'token_usage'), |
| turnId: event.turnId, |
| ts: event.ts, |
| input: usage.input, |
| output: usage.output, |
| ...(usage.cacheHitInput !== undefined ? { cacheHitInput: usage.cacheHitInput } : {}), |
| ...(usage.cacheMissInput !== undefined ? { cacheMissInput: usage.cacheMissInput } : {}), |
| ...(usage.cacheMissInputSource !== undefined |
| ? { cacheMissInputSource: usage.cacheMissInputSource } |
| : {}), |
| ...(usage.cacheWriteInput !== undefined ? { cacheWriteInput: usage.cacheWriteInput } : {}), |
| ...(usage.reasoning !== undefined ? { reasoning: usage.reasoning } : {}), |
| ...(usage.total !== undefined ? { total: usage.total } : {}), |
| ...(usage.rawFinishReason !== undefined ? { rawFinishReason: usage.rawFinishReason } : {}), |
| ...(usage.runtimeSteps !== undefined ? { runtimeSteps: usage.runtimeSteps } : {}), |
| ...(usage.cacheRead !== undefined ? { cacheRead: usage.cacheRead } : {}), |
| ...(usage.cacheCreation !== undefined ? { cacheCreation: usage.cacheCreation } : {}), |
| ...(usage.costUsd !== undefined ? { costUsd: usage.costUsd } : {}), |
| ...(usage.systemPromptHash !== undefined ? { systemPromptHash: usage.systemPromptHash } : {}), |
| ...(usage.contextRemaining !== undefined ? { contextRemaining: usage.contextRemaining } : {}), |
| ...(usage.prefixHash !== undefined ? { prefixHash: usage.prefixHash } : {}), |
| ...(usage.prefixChangeReason !== undefined |
| ? { prefixChangeReason: usage.prefixChangeReason } |
| : {}), |
| ...(usage.requestShapeHash !== undefined ? { requestShapeHash: usage.requestShapeHash } : {}), |
| ...(usage.requestShapeChangeReason !== undefined |
| ? { requestShapeChangeReason: usage.requestShapeChangeReason } |
| : {}), |
| ...(usage.promptSegments !== undefined ? { promptSegments: usage.promptSegments } : {}), |
| ...(usage.contextBudget !== undefined ? { contextBudget: usage.contextBudget } : {}), |
| ...(usage.lastRequestAnchor !== undefined |
| ? { lastRequestAnchor: usage.lastRequestAnchor } |
| : {}), |
| ...(event.refs?.providerRequestTraceId !== undefined |
| ? { providerRequestTraceId: event.refs.providerRequestTraceId } |
| : {}), |
| }); |
| return true; |
| } |
| |
| function projectTerminalTurnState( |
| event: RuntimeEvent, |
| state: ProjectionState, |
| messages: StoredMessage[], |
| ): boolean { |
| const invocation = state.invocations.get(event.runId); |
| if (!invocation) { |
| diagnostic( |
| state, |
| event, |
| 'incomplete_event', |
| 'terminal RuntimeEvent requires the opening fact of its invocation', |
| ); |
| return false; |
| } |
| const lineage = invocation.opening.lineage; |
| const status = turnStatusFor(event.status); |
| if (!status) { |
| diagnostic( |
| state, |
| event, |
| 'incomplete_event', |
| 'terminal RuntimeEvent status cannot be mapped to a legacy TurnStatus', |
| ); |
| return false; |
| } |
| const abortSource = status === 'aborted' ? abortSourceFromRuntime(event) : undefined; |
| const failureClass = status === 'failed' ? failureClassFromRuntimeEvent(event) : undefined; |
| messages.push({ |
| type: 'turn_state', |
| id: stableMessageId(event, state, 'turn_state'), |
| turnId: event.turnId, |
| ts: event.ts, |
| status, |
| ...(lineage?.parentTurnId ? { parentTurnId: lineage.parentTurnId } : {}), |
| ...(lineage?.retriedFromTurnId ? { retriedFromTurnId: lineage.retriedFromTurnId } : {}), |
| ...(lineage?.regeneratedFromTurnId |
| ? { regeneratedFromTurnId: lineage.regeneratedFromTurnId } |
| : {}), |
| ...(lineage?.branchOfTurnId ? { branchOfTurnId: lineage.branchOfTurnId } : {}), |
| ...(lineage?.parentSessionId ? { parentSessionId: lineage.parentSessionId } : {}), |
| ...(status === 'aborted' ? { abortedAt: event.ts } : {}), |
| ...(abortSource ? { abortSource } : {}), |
| ...(status === 'failed' ? { errorClass: failureClass ?? 'unknown' } : {}), |
| ...(status === 'failed' && event.content?.kind === 'error' && event.content.message |
| ? { |
| failureMessage: truncateUtf8(event.content.message, MODEL_FAILURE_MESSAGE_MAX_BYTES, '…'), |
| } |
| : {}), |
| ...(status === 'failed' && event.content?.kind === 'error' && event.content.retry |
| ? { retry: event.content.retry } |
| : {}), |
| }); |
| if (failureClass === 'tool_step_cap_reached') { |
| messages.push({ |
| type: 'system_note', |
| id: `${event.id}:step-limit-notice`, |
| turnId: event.turnId, |
| ts: event.ts, |
| kind: 'step_limit', |
| }); |
| } |
| // An omitted failure class or abort source is `classifyRuntimeEventTerminalFact`'s |
| // observation to make. Repeating it here would only turn a transcript row that |
| // already reads `unknown` into an unreadable Session. |
| return true; |
| } |
| |
| /** |
| * The note row of an invocation that wrote one. |
| * |
| * There is nothing to reconcile: the event carries the kind and the payload the |
| * row is made of, so the row is the event said back in the transcript's shape. |
| */ |
| function projectSystemNote( |
| event: RuntimeEvent, |
| state: ProjectionState, |
| messages: StoredMessage[], |
| ): boolean { |
| if (event.content?.kind !== 'system_note') return false; |
| messages.push({ |
| type: 'system_note', |
| id: stableMessageId(event, state, 'system_note'), |
| turnId: event.turnId, |
| ts: event.ts, |
| kind: event.content.note, |
| ...(event.content.data !== undefined ? { data: structuredClone(event.content.data) } : {}), |
| }); |
| return true; |
| } |
| |
| function attachPendingThinking( |
| event: RuntimeEvent, |
| state: ProjectionState, |
| messages: StoredMessage[], |
| assistantMessageId: string, |
| ): void { |
| const pendingItems = state.thinkingByMessageId.get(assistantMessageId); |
| if (!pendingItems) return; |
| if (pendingItems.every((pending) => attachThinkingToAssistant(event, pending, messages))) { |
| state.thinkingByMessageId.delete(assistantMessageId); |
| } |
| } |
| |
| function attachThinkingToAssistant( |
| event: RuntimeEvent, |
| pending: PendingThinking, |
| messages: StoredMessage[], |
| ): boolean { |
| // Attach to the assistant row whose id equals the thinking's step message id |
| // (per-step pairing). Scans from the tail so the newest matching row wins. |
| for (let index = messages.length - 1; index >= 0; index -= 1) { |
| const message = messages[index]!; |
| if (message.type !== 'assistant' || message.turnId !== event.turnId) continue; |
| if (message.id !== pending.messageId) continue; |
| const incoming = { |
| text: pending.text, |
| ...(pending.signature !== undefined ? { signature: pending.signature } : {}), |
| ...(pending.providerOptions !== undefined |
| ? { providerOptions: structuredClone(pending.providerOptions) } |
| : {}), |
| }; |
| if (!message.thinking) { |
| message.thinking = incoming; |
| return true; |
| } |
| const parts = message.thinking.parts ?? [ |
| { |
| text: message.thinking.text, |
| ...(message.thinking.signature !== undefined |
| ? { signature: message.thinking.signature } |
| : {}), |
| ...(message.thinking.providerOptions !== undefined |
| ? { providerOptions: structuredClone(message.thinking.providerOptions) } |
| : {}), |
| }, |
| ]; |
| message.thinking = { |
| text: message.thinking.text + pending.text, |
| parts: [...parts, incoming], |
| }; |
| return true; |
| } |
| return false; |
| } |
| |
| function thinkingMessageId(event: RuntimeEvent): string { |
| return event.refs?.providerEventId ?? event.refs?.storedMessageId ?? event.id; |
| } |
| |
| /** |
| * Why this invocation failed, according to its own terminal event. |
| * |
| * `undefined` for an invocation that is still running or did not fail. There is |
| * no second place to look: the event that ends the run also states the class. |
| */ |
| export function runtimeInvocationFailureClass(invocation: { |
| terminalEvent?: RuntimeEvent; |
| }): string | undefined { |
| const terminalEvent = invocation.terminalEvent; |
| if (terminalEvent?.status !== 'failed') return undefined; |
| return failureClassFromRuntimeEvent(terminalEvent); |
| } |
| |
| function abortSourceFromRuntime(event: RuntimeEvent): string | undefined { |
| return ( |
| stringStateDelta(event, 'abortSource') ?? |
| stringStateDelta(event, 'source') ?? |
| stringRecordValue(event.refs, 'abortSource') ?? |
| stringRecordValue(event.refs, 'source') |
| ); |
| } |
| |
| function failureClassFromRuntimeEvent(event: RuntimeEvent): string | undefined { |
| const failureClass = |
| stringStateDelta(event, 'failureClass') ?? |
| stringStateDelta(event, 'errorClass') ?? |
| stringStateDelta(event, 'reason') ?? |
| stringStateDelta(event, 'code') ?? |
| (event.content?.kind === 'error' ? nonEmptyString(event.content.reason) : undefined) ?? |
| (event.content?.kind === 'error' ? nonEmptyString(event.content.code) : undefined); |
| // Retired outcome. The runtime no longer decides locally that a request |
| // cannot be shaped to fit — the provider rejects it and recovery compacts and |
| // retries — so a turn that ends over the window is a context overflow like any |
| // other. Sessions written before that still carry the old name; fold it here, |
| // at the one place the durable ledger is read, so nothing downstream has to |
| // know two names for one outcome. |
| return failureClass === 'context_budget_exhausted' ? 'context_overflow' : failureClass; |
| } |
| |
| function stringRecordValue(value: unknown, key: string): string | undefined { |
| if (!value || typeof value !== 'object') return undefined; |
| const result = (value as Record<string, unknown>)[key]; |
| return typeof result === 'string' && result.length > 0 ? result : undefined; |
| } |
| |
| function nonEmptyString(value: unknown): string | undefined { |
| return typeof value === 'string' && value.length > 0 ? value : undefined; |
| } |
| |
| function stableMessageId( |
| event: RuntimeEvent, |
| state: ProjectionState, |
| kind: StoredMessage['type'], |
| contentId?: string, |
| ): string { |
| const stable = |
| event.refs?.storedMessageId ?? event.refs?.providerEventId ?? contentId ?? event.id; |
| if (stable) return stable; |
| const generated = `rtproj:${event.id}:${kind}`; |
| diagnostic(state, event, 'generated_id', 'projection used a deterministic generated id', { |
| id: generated, |
| }); |
| return generated; |
| } |
| |
| function toolUseIdFor(event: RuntimeEvent): string | undefined { |
| if (event.content?.kind !== 'function_call' && event.content?.kind !== 'function_response') { |
| return event.refs?.toolCallId; |
| } |
| return event.content.id || event.refs?.toolCallId; |
| } |
| |
| function normalizeInvocations( |
| invocations: |
| | readonly RuntimeInvocationRecord[] |
| | Readonly<Record<string, RuntimeInvocationRecord>>, |
| ): Map<string, RuntimeInvocationRecord> { |
| const values = Array.isArray(invocations) |
| ? (invocations as readonly RuntimeInvocationRecord[]) |
| : Object.values(invocations as Readonly<Record<string, RuntimeInvocationRecord>>); |
| return new Map(values.map((invocation) => [invocation.runId, invocation])); |
| } |
| |
| /** The terminal event states the outcome; nothing else is allowed to disagree. */ |
| function turnStatusFor(eventStatus: RuntimeEventStatus | undefined): TurnStatus | undefined { |
| if (eventStatus === 'completed') return 'completed'; |
| if (eventStatus === 'failed') return 'failed'; |
| if (eventStatus === 'aborted' || eventStatus === 'cancelled') return 'aborted'; |
| return undefined; |
| } |
| |
| function stringStateDelta(event: RuntimeEvent, key: string): string | undefined { |
| const value = event.actions?.stateDelta?.[key]; |
| return typeof value === 'string' ? value : undefined; |
| } |
| |
| function toolActivityKindStateDelta(event: RuntimeEvent): ToolActivityKind | undefined { |
| const value = stringStateDelta(event, 'activityKind'); |
| return TOOL_ACTIVITY_KINDS.find((kind) => kind === value); |
| } |
| |
| function numberStateDelta(event: RuntimeEvent, key: string): number | undefined { |
| const value = event.actions?.stateDelta?.[key]; |
| return typeof value === 'number' ? value : undefined; |
| } |
| |
| function isLegacyPlanToolResult(value: unknown): boolean { |
| if (!value || typeof value !== 'object') return false; |
| const kind = (value as { kind?: unknown }).kind; |
| return ( |
| kind === 'plan_submitted' || |
| kind === 'plan_progress_updated' || |
| kind === 'plan_execution_completed' || |
| kind === 'plan_execution_cancelled' |
| ); |
| } |
| |
| function isPlanProposalStateDelta(event: RuntimeEvent): boolean { |
| const stateDelta = event.actions?.stateDelta; |
| return ( |
| event.role === 'system' && |
| event.author === 'agent' && |
| typeof stateDelta?.planId === 'string' && |
| typeof stateDelta.title === 'string' |
| ); |
| } |
| |
| /** |
| * A boundary fact is canonical only in the exact shape the Runtime mapper emits: every |
| * field of the source SessionEvent, the identity it maps to, and the tool call |
| * it settles. A partial match is worse than none — it would claim a corrupt |
| * ledger as sound while still paying the cost of rejecting a malformed one. |
| */ |
| function isSandboxBoundaryStateDelta(event: RuntimeEvent): boolean { |
| const stateDelta = event.actions?.stateDelta; |
| if (!stateDelta) return false; |
| const request = stateDelta.sandboxBoundaryRequest; |
| const decision = stateDelta.sandboxBoundaryDecision; |
| if (request === undefined && decision === undefined) return false; |
| if (event.role !== 'system' || typeof event.refs?.toolCallId !== 'string') return false; |
| if (request !== undefined) { |
| return ( |
| event.author === 'system' && |
| isRecord(request) && |
| typeof request.requestId === 'string' && |
| typeof request.toolUseId === 'string' && |
| typeof request.justification === 'string' && |
| validateSandboxBoundaryExpansion(request.expansion).ok |
| ); |
| } |
| return ( |
| event.author === 'user' && |
| isRecord(decision) && |
| typeof decision.requestId === 'string' && |
| (decision.decision === 'allow' || decision.decision === 'deny') && |
| SETTLED_SANDBOX_BOUNDARY_STATUSES.includes(decision.status as SettledSandboxBoundaryStatus) && |
| typeof decision.revision === 'number' && |
| Number.isFinite(decision.revision) |
| ); |
| } |
| |
| function isRecord(value: unknown): value is Record<string, unknown> { |
| return typeof value === 'object' && value !== null && !Array.isArray(value); |
| } |
| |
| function diagnostic( |
| state: ProjectionState, |
| event: RuntimeEvent, |
| code: RuntimeEventReadModelDiagnosticCode, |
| message: string, |
| detail?: unknown, |
| ): void { |
| state.diagnostics.push({ |
| code, |
| eventId: event.id, |
| runId: event.runId, |
| turnId: event.turnId, |
| message, |
| ...(detail !== undefined ? { detail } : {}), |
| }); |
| } |
| |
| function readModelDiagnostic( |
| code: RuntimeEventReadModelDiagnosticCode, |
| message: string, |
| detail: RuntimeEvent | { runId: string; turnId: string; [key: string]: unknown }, |
| ): RuntimeEventReadModelDiagnostic { |
| if (isRuntimeEventDiagnosticDetail(detail)) { |
| return { |
| code, |
| eventId: detail.id, |
| runId: detail.runId, |
| turnId: detail.turnId, |
| message, |
| }; |
| } |
| return { |
| code, |
| runId: detail.runId, |
| turnId: detail.turnId, |
| message, |
| detail, |
| }; |
| } |
| |
| function isRuntimeEventDiagnosticDetail( |
| detail: RuntimeEvent | { runId: string; turnId: string; [key: string]: unknown }, |
| ): detail is RuntimeEvent { |
| return ( |
| typeof (detail as RuntimeEvent).id === 'string' && |
| typeof (detail as RuntimeEvent).sessionId === 'string' && |
| typeof detail.runId === 'string' && |
| typeof detail.turnId === 'string' |
| ); |
| } |