diff --git a/packages/core/src/__tests__/recall.test.ts b/packages/core/src/__tests__/recall.test.ts index f5456169eb..63e97a12e6 100644 --- a/packages/core/src/__tests__/recall.test.ts +++ b/packages/core/src/__tests__/recall.test.ts @@ -363,10 +363,18 @@ test('a candidate source without a corpus count is not used', async () => { assert.equal(missing.scannedFully, true); assert.equal(asked, 0, 'the candidate source must not even be asked'); + const reads: string[] = []; const declining = await runRecall( { terms: ['上下文'] }, - candidateDeps(data, { countSearchableMessages: async () => null }), + candidateDeps(data, { + countSearchableMessages: async () => null, + readMessages: async (sessionId) => { + reads.push(sessionId); + return data.messages.get(sessionId) ?? null; + }, + }), ); + assert.equal(reads.length, new Set(reads).size, 'count fallback must not reread a Session'); assert.ok(declining.ok); assert.equal(declining.scannedFully, true); const scanned = await runRecall({ terms: ['上下文'] }, scanDeps(data)); diff --git a/packages/core/src/recall.ts b/packages/core/src/recall.ts index d8ef58ad35..29470f5b1e 100644 --- a/packages/core/src/recall.ts +++ b/packages/core/src/recall.ts @@ -1029,6 +1029,18 @@ async function collectHits( } } + // A narrowed scan needs the corpus count to rank the same way as a full + // scan. Resolve it before reading transcripts so a declined count cannot + // discard work and reread every candidate during fallback. + const storedCorpusSize = deps.countSearchableMessages + ? await deps.countSearchableMessages({ sessionIds: input.sessionIds }) + : null; + if (input.abortSignal?.aborted) return null; + if (storedCorpusSize === null && !scannedFully) { + sessionIds = input.sessionIds; + scannedFully = true; + } + const hits: VerifiedHit[] = []; let counted = 0; for (const sessionId of sessionIds) { @@ -1053,18 +1065,7 @@ async function collectHits( } } - // The narrowed path never sees the Sessions it skipped, so its corpus size - // has to come from the store. The full scan takes the same number when the - // store offers one, which is what keeps a score independent of the path. - let corpusSize: number | null = counted; - if (deps.countSearchableMessages) { - corpusSize = await deps.countSearchableMessages({ sessionIds: input.sessionIds }); - if (input.abortSignal?.aborted) return null; - if (corpusSize === null) { - if (!scannedFully) return collectHits(deps, { ...input, forceFullScan: true }); - corpusSize = counted; - } - } + const corpusSize = storedCorpusSize ?? counted; return { hits, corpusSize, scannedFully, transcripts }; } diff --git a/packages/core/src/runtime-event-store.ts b/packages/core/src/runtime-event-store.ts index 7c8f65c5ad..39f4e72469 100644 --- a/packages/core/src/runtime-event-store.ts +++ b/packages/core/src/runtime-event-store.ts @@ -116,6 +116,14 @@ export class DurableStoreWriteError extends Error { } } +/** One consistent read of the inputs needed to project a complete Session. */ +export interface RuntimeSessionEventSnapshot { + readonly invocations: RuntimeInvocationRecord[]; + /** Per-run event order, including merged mutable partial presentation. */ + readonly eventsByRun: ReadonlyMap; + readonly durableEventOrdinalById: ReadonlyMap; +} + export interface RuntimeEventStore { /** Canonical stores fail the active run closed on every durable write error. */ readonly durability?: 'best_effort' | 'canonical'; @@ -210,6 +218,12 @@ export interface RuntimeEventStore { unsettledOperationIds: readonly string[], ): Promise; readRuntimeEvents(sessionId: string, runId: string): Promise; + /** + * Batch full-view reads in one snapshot, sharing decoded events between the + * invocation inventory and run histories. Optional for stores without batch + * support; consumers retain the individual-reader path in that case. + */ + readSessionRuntimeSnapshot?(sessionId: string): Promise; /** Session-wide immutable append order. */ readSessionRuntimeEventEntries( sessionId: string, diff --git a/packages/runtime/src/__tests__/recall-ledger-corpus.test.ts b/packages/runtime/src/__tests__/recall-ledger-corpus.test.ts index db0fff9d51..dfa33493e6 100644 --- a/packages/runtime/src/__tests__/recall-ledger-corpus.test.ts +++ b/packages/runtime/src/__tests__/recall-ledger-corpus.test.ts @@ -29,6 +29,7 @@ import type { CreateSessionInput } from '@maka/core/runtime-inputs'; import type { StoredMessage } from '@maka/core/session'; import { createSessionStore } from '@maka/storage/session-store'; import { createSqliteRuntimeStore } from '@maka/storage/sqlite-runtime-store'; +import { openRuntimeEventReadPersistence } from '@maka/storage/runtime-event-persistence'; import { countRecallSearchableMessages, listRecallCandidateSessions, @@ -185,7 +186,174 @@ function anchors(result: Awaited>): string[] { return result.passages.map((passage) => passage.anchorMessageId).sort(); } +function unbatchedReadModel(workspace: Workspace): RuntimeReadModel { + return new RuntimeReadModel({ + runtimeEventStore: new Proxy(workspace.runtime, { + get(target, key) { + if (key === 'readSessionRuntimeSnapshot') return undefined; + const value = Reflect.get(target, key, target); + return typeof value === 'function' ? value.bind(target) : value; + }, + }), + }); +} + describe('recall over the corpus a live workspace actually has', () => { + test('reading Recall messages uses bounded SQL and decodes each durable event once', async (t) => { + await withWorkspace(async (workspace) => { + const session = await workspace.sessions.create(makeInput('long history')); + const turns = 12; + for (let index = 0; index < turns; index += 1) { + await seedLedgerTurn( + workspace, + session.id, + `run-${index}`, + index < 10 ? 'deploy target' : 'unrelated', + 'ordinary response', + ); + } + const before = await runRecall( + { terms: ['deploy'], limit: 10 }, + narrowedDeps({ ...workspace, readModel: unbatchedReadModel(workspace) }), + ); + assert.ok(before.ok); + assert.equal(before.passages.length, 10); + + const db = (workspace.runtime as unknown as { db: DatabaseSync }).db; + const prepare = db.prepare.bind(db); + let statements = 0; + let parses = 0; + t.mock.method(db, 'prepare', (sql: string) => { + const statement = prepare(sql); + return new Proxy(statement, { + get(target, key) { + const value = Reflect.get(target, key, target); + if (typeof value !== 'function') return value; + return (...args: unknown[]) => { + if (key === 'get' || key === 'all' || key === 'iterate') statements += 1; + return Reflect.apply(value, target, args); + }; + }, + }); + }); + const parse = JSON.parse; + t.mock.method(JSON, 'parse', (...args: Parameters) => { + parses += 1; + return parse(...args); + }); + try { + const messages = await workspace.readModel.getSessionMessages(session.id); + assert.equal(messages.filter((message) => message.type === 'user').length, turns); + } finally { + t.mock.restoreAll(); + } + assert.ok(statements <= 6, `${turns} Turns issued ${statements} SQL reads`); + assert.equal(parses, turns * 4, 'each opening, user, assistant and terminal decoded once'); + assert.deepEqual( + await runRecall({ terms: ['deploy'], limit: 10 }, narrowedDeps(workspace)), + before, + ); + }); + }); + + test('batch projection preserves active partials, child exclusion and Recall navigation', async () => { + await withWorkspace(async (workspace) => { + const session = await workspace.sessions.create(makeInput('mixed history')); + await seedLedgerTurn(workspace, session.id, 'settled', 'deploy target', 'deployment notes'); + for (const runId of ['active-a', 'active-b', 'child']) { + const identity = { + sessionId: session.id, + invocationId: runId, + runId, + turnId: `turn-${runId}`, + }; + await workspace.runtime.appendRuntimeEvent( + session.id, + runId, + testInvocationOpenedEvent({ + ...identity, + openedAt: 200, + ...(runId === 'child' ? { opening: { lineage: { parentRunId: 'active-a' } } } : {}), + }), + ); + await workspace.runtime.appendRuntimeEvent(session.id, runId, { + ...identity, + id: `${runId}-prompt`, + ts: 199, + partial: false, + role: 'user', + author: 'user', + content: { kind: 'text', text: 'deploy target' }, + }); + await workspace.runtime.appendRuntimePartialBatch(session.id, runId, [ + { + ...identity, + id: `${runId}-p1`, + ts: 201, + partial: true, + role: 'model', + author: 'agent', + refs: { providerEventId: 'stream' }, + content: { kind: 'text', text: 'deploy ' }, + }, + { + ...identity, + id: `${runId}-p2`, + ts: 202, + partial: true, + role: 'model', + author: 'agent', + refs: { providerEventId: 'stream' }, + content: { kind: 'text', text: 'in progress' }, + }, + ]); + } + const fallback = unbatchedReadModel(workspace); + const expected = await fallback.getSessionView(session.id); + const actual = await workspace.readModel.getSessionView(session.id); + assert.deepEqual(actual, expected); + assert.ok( + actual.messages.some( + (m) => m.type === 'assistant' && m.text.includes('deploy in progress'), + ), + ); + assert.ok(actual.events.every((event) => event.runId !== 'child')); + for (const terms of [['deploy'], ['progress'], ['nothing-here']]) { + assert.deepEqual( + await runRecall({ terms, limit: 10 }, narrowedDeps(workspace)), + await runRecall( + { terms, limit: 10 }, + narrowedDeps({ ...workspace, readModel: fallback }), + ), + ); + } + const reader = await openRuntimeEventReadPersistence({ workspaceRoot: workspace.root }); + try { + assert.ok(reader.runtimeEventStore.readSessionRuntimeSnapshot); + assert.deepEqual( + await reader.runtimeEventStore.readSessionRuntimeSnapshot(session.id), + await workspace.runtime.readSessionRuntimeSnapshot(session.id), + ); + } finally { + reader.close(); + } + }); + }); + + test('a failed batch snapshot is reported instead of retried as separate reads', async (t) => { + await withWorkspace(async (workspace) => { + t.mock.method(workspace.runtime, 'readSessionRuntimeSnapshot', async () => { + throw new Error('snapshot unavailable'); + }); + const list = t.mock.method(workspace.runtime, 'listSessionInvocations'); + await assert.rejects( + workspace.readModel.getSessionMessages('session'), + /snapshot read failed/, + ); + assert.equal(list.mock.callCount(), 0); + }); + }); + test('a Session born on the ledger is reachable, and narrowing matches the full scan', async () => { await withWorkspace(async (workspace) => { const live = await workspace.sessions.create(makeInput('live')); diff --git a/packages/runtime/src/runtime-read-model.ts b/packages/runtime/src/runtime-read-model.ts index 221f260b02..b2096f051c 100644 --- a/packages/runtime/src/runtime-read-model.ts +++ b/packages/runtime/src/runtime-read-model.ts @@ -18,7 +18,10 @@ */ import type { RuntimeEvent } from '@maka/core/runtime-event'; -import type { RuntimeEventStore } from '@maka/core/runtime-event-store'; +import type { + RuntimeEventStore, + RuntimeSessionEventSnapshot, +} from '@maka/core/runtime-event-store'; import type { RuntimeInvocationRecord } from '@maka/core/runtime-invocation'; import type { StoredMessage, TurnRecord } from '@maka/core/session'; import { deriveTurnRecords } from '@maka/core/session'; @@ -78,11 +81,28 @@ export class RuntimeReadModel { async getSessionView(sessionId: string): Promise { const diagnostics: RuntimeEventReadModelDiagnostic[] = []; const inFlightTurnIds = new Set(); + let snapshot: RuntimeSessionEventSnapshot | undefined; + if (this.deps.runtimeEventStore.readSessionRuntimeSnapshot) { + try { + snapshot = await this.deps.runtimeEventStore.readSessionRuntimeSnapshot(sessionId); + } catch (error) { + throw new RuntimeReadModelError('RuntimeEvent Session snapshot read failed', [ + readModelDiagnostic( + 'unsupported_event', + 'RuntimeEventStore.readSessionRuntimeSnapshot failed', + { + error: errorMessage(error), + }, + ), + ]); + } + } let invocations: RuntimeInvocationRecord[]; try { - invocations = (await this.deps.runtimeEventStore.listSessionInvocations(sessionId)).filter( - (invocation) => isSessionInlineInvocation(invocation.opening), - ); + invocations = ( + snapshot?.invocations ?? + (await this.deps.runtimeEventStore.listSessionInvocations(sessionId)) + ).filter((invocation) => isSessionInlineInvocation(invocation.opening)); } catch (error) { throw new RuntimeReadModelError('RuntimeReadModel could not list Session invocations', [ readModelDiagnostic( @@ -99,20 +119,23 @@ export class RuntimeReadModel { return this.buildView({ invocations, events: [], diagnostics }); } - const durableEventOrdinals = await this.readSessionRuntimeEventOrdinals(sessionId); - const durableEventOrdinalById = new Map( - durableEventOrdinals.map(({ event, ordinal }) => [event.id, ordinal]), - ); + const durableEventOrdinalById = + snapshot?.durableEventOrdinalById ?? + new Map( + (await this.readSessionRuntimeEventOrdinals(sessionId)).map(({ event, ordinal }) => [ + event.id, + ordinal, + ]), + ); const ordered: OrderedRuntimeEvent[] = []; const terminalFacts: RuntimeEventTerminalFact[] = []; for (let runIndex = 0; runIndex < invocations.length; runIndex += 1) { const invocation = invocations[runIndex]!; let runEvents: RuntimeEvent[]; try { - runEvents = await this.deps.runtimeEventStore.readRuntimeEvents( - sessionId, - invocation.runId, - ); + runEvents = snapshot + ? (snapshot.eventsByRun.get(invocation.runId) ?? []) + : await this.deps.runtimeEventStore.readRuntimeEvents(sessionId, invocation.runId); } catch (error) { throw new RuntimeReadModelError('RuntimeEvent ledger read failed', [ readModelDiagnostic('unsupported_event', 'RuntimeEventStore.readRuntimeEvents failed', { diff --git a/packages/storage/src/__tests__/execution-provider-conformance.test.ts b/packages/storage/src/__tests__/execution-provider-conformance.test.ts index 7e5643651a..bea4d967ef 100644 --- a/packages/storage/src/__tests__/execution-provider-conformance.test.ts +++ b/packages/storage/src/__tests__/execution-provider-conformance.test.ts @@ -605,6 +605,21 @@ for (const backend of ['Local', 'Memory'] as const) { // Repair renumbers every ordinal, so the extents move with them. await s.resequenceSessionEventOrdinals(sessionId); const entries = await s.readSessionRuntimeEventEntries(sessionId); + if (backend === 'Local') { + assert.ok(s.readSessionRuntimeSnapshot, 'the production facade must expose batch reads'); + const snapshot = await s.readSessionRuntimeSnapshot(sessionId); + assert.deepEqual(snapshot.invocations, await s.listSessionInvocations(sessionId)); + assert.deepEqual( + snapshot.durableEventOrdinalById, + new Map(entries.map(({ event, ordinal }) => [event.id, ordinal])), + ); + for (const invocation of snapshot.invocations) { + assert.deepEqual( + snapshot.eventsByRun.get(invocation.runId), + await s.readRuntimeEvents(sessionId, invocation.runId), + ); + } + } const ordinalsOf = (turnId: string) => entries.filter((entry) => entry.event.turnId === turnId).map((entry) => entry.ordinal); const [moved] = await s.readTranscriptTurns(sessionId, { turnId: outer.turnId }); diff --git a/packages/storage/src/__tests__/invocation-opening-backfill.test.ts b/packages/storage/src/__tests__/invocation-opening-backfill.test.ts index c20343bec1..4656001845 100644 --- a/packages/storage/src/__tests__/invocation-opening-backfill.test.ts +++ b/packages/storage/src/__tests__/invocation-opening-backfill.test.ts @@ -216,6 +216,18 @@ describe('invocation opening fact backfill', () => { const store = createSqliteRuntimeStore(databasePath); try { const invocations = await store.listSessionInvocations('session-1'); + const snapshot = await store.readSessionRuntimeSnapshot('session-1'); + assert.deepEqual( + snapshot.invocations, + invocations, + 'batch reads preserve migrated openings', + ); + for (const invocation of invocations) { + assert.deepEqual( + snapshot.eventsByRun.get(invocation.runId) ?? [], + await store.readRuntimeEvents('session-1', invocation.runId), + ); + } assert.deepEqual( invocations.map((invocation) => invocation.invocationId), ['run-legacy-route', 'run-scheduled', 'run-with-events'], diff --git a/packages/storage/src/__tests__/sqlite-runtime-store.test.ts b/packages/storage/src/__tests__/sqlite-runtime-store.test.ts index 9d1ad92885..aeeb32fbd0 100644 --- a/packages/storage/src/__tests__/sqlite-runtime-store.test.ts +++ b/packages/storage/src/__tests__/sqlite-runtime-store.test.ts @@ -65,6 +65,98 @@ const PREFIX_PROOF_TEST_BUDGET = { }; describe('SqliteRuntimeStore', () => { + it('reads a complete Session from one snapshot while another connection commits its terminal', async (t) => { + await withStore(async (store, dbPath) => { + const opening = invocationOpeningEvent(1); + await store.appendRuntimeEvent(opening.sessionId, opening.runId, opening); + const terminal: RuntimeEvent = { + ...opening, + id: 'concurrent-terminal', + ts: 20, + content: undefined, + status: 'completed', + actions: { endInvocation: true }, + }; + const writer = new DatabaseSync(dbPath); + const db = (store as unknown as { db: DatabaseSync }).db; + const prepare = db.prepare.bind(db); + let committed = false; + t.mock.method(db, 'prepare', (sql: string) => { + if (!committed && sql.includes("CASE WHEN event_kind = 'invocation_opened'")) { + committed = true; + writer.exec('BEGIN'); + writer + .prepare(`INSERT INTO runtime_events + (event_id, session_id, invocation_id, run_id, turn_id, event_seq, event_kind, payload_json, committed_at) + VALUES (?, ?, ?, ?, ?, 2, 'invocation_end', ?, ?)`) + .run( + terminal.id, + terminal.sessionId, + terminal.invocationId, + terminal.runId, + terminal.turnId, + JSON.stringify(terminal), + terminal.ts, + ); + writer + .prepare( + 'INSERT INTO runtime_session_event_ordinals (session_id, ordinal, event_id) VALUES (?, 2, ?)', + ) + .run(terminal.sessionId, terminal.id); + writer.exec('COMMIT'); + } + return prepare(sql); + }); + try { + const snapshot = await store.readSessionRuntimeSnapshot(opening.sessionId); + assert.ok(committed); + assert.equal(snapshot.invocations[0]?.terminalEvent, undefined); + assert.deepEqual( + snapshot.eventsByRun.get(opening.runId)?.map((event) => event.id), + [opening.id], + ); + assert.deepEqual([...snapshot.durableEventOrdinalById], [[opening.id, 1]]); + const next = await store.readSessionRuntimeSnapshot(opening.sessionId); + assert.equal(next.invocations[0]?.terminalEvent?.id, terminal.id); + assert.deepEqual( + next.eventsByRun.get(opening.runId)?.map((event) => event.id), + [opening.id, terminal.id], + ); + assert.equal(next.durableEventOrdinalById.get(terminal.id), 2); + } finally { + t.mock.restoreAll(); + writer.close(); + } + }); + }); + + it('batch Session reads retain event identity validation', async () => { + await withStore(async (store, dbPath) => { + const opening = invocationOpeningEvent(1); + await store.appendRuntimeEvent(opening.sessionId, opening.runId, opening); + const text = functionCallEvent({ + id: 'corrupt-text', + content: { kind: 'text', text: 'hello' }, + }); + await store.appendRuntimeEvent(text.sessionId, text.runId, text); + const writer = new DatabaseSync(dbPath); + try { + writer + .prepare( + "UPDATE runtime_events SET payload_json = json_set(payload_json, '$.sessionId', 'wrong-session') WHERE event_id = ?", + ) + .run(text.id); + await assert.rejects( + store.readRuntimeEvents(text.sessionId, text.runId), + /identity mismatch/, + ); + await assert.rejects(store.readSessionRuntimeSnapshot(text.sessionId), /identity mismatch/); + } finally { + writer.close(); + } + }); + }); + it('applies versioned migrations and reopens the same database without rewriting schema', async () => { await withStore(async (store, dbPath) => { assert.equal(store.schemaVersion(), SQLITE_RUNTIME_SCHEMA_VERSION); diff --git a/packages/storage/src/execution-stores.ts b/packages/storage/src/execution-stores.ts index a11ea09628..ee57dd52e7 100644 --- a/packages/storage/src/execution-stores.ts +++ b/packages/storage/src/execution-stores.ts @@ -22,6 +22,7 @@ import type { RuntimeEvent, ToolBoundaryProtocol } from '@maka/core/runtime-even import type { RuntimeContinuationAuthorityStore, RuntimeInvocationRecoveryInventoryEntry, + RuntimeSessionEventSnapshot, } from '@maka/core/runtime-event-store'; import type { ImmutableRuntimePrefixProofV1 } from '@maka/core/runtime-boundary'; import type { RuntimeTranscriptQueries } from './runtime-transcript-query.js'; @@ -258,6 +259,7 @@ export interface ExecutionAgentRunReader { } export interface ExecutionRuntimeEventReader { + readSessionRuntimeSnapshot?(sessionId: string): Promise; /** * A Session's run inventory, read from its canonical events. This is the * definition of the inventory, not a cache of it, so nothing writes or @@ -756,6 +758,12 @@ async function createExecutionStoresForWrite( ), readRuntimeEvents: (sessionId, runId) => run(() => runtimeEventStore.readRuntimeEvents(sessionId, runId)), + ...(runtimeEventStore.readSessionRuntimeSnapshot + ? { + readSessionRuntimeSnapshot: (sessionId: string) => + run(() => runtimeEventStore.readSessionRuntimeSnapshot!(sessionId)), + } + : {}), scanRuntimeEvents: (sessionId, runId, budget, visit) => run(() => runtimeEventStore.scanRuntimeEvents(sessionId, runId, budget, visit)), readRuntimeEventsBounded: (sessionId, runId, budget) => @@ -921,6 +929,12 @@ async function openExecutionStoresForRead run(() => runtimeEventStore.readRuntimeEvents(sessionId, runId)), + ...(runtimeEventStore.readSessionRuntimeSnapshot + ? { + readSessionRuntimeSnapshot: (sessionId: string) => + run(() => runtimeEventStore.readSessionRuntimeSnapshot!(sessionId)), + } + : {}), readRuntimeEventsBounded: (sessionId, runId, budget) => run(() => runtimeEventStore.readRuntimeEventsBounded(sessionId, runId, budget)), readImmutableRuntimeEvents: (sessionId, runId) => diff --git a/packages/storage/src/runtime-event-persistence.ts b/packages/storage/src/runtime-event-persistence.ts index 5e28ceb4bb..04c00cdf75 100644 --- a/packages/storage/src/runtime-event-persistence.ts +++ b/packages/storage/src/runtime-event-persistence.ts @@ -19,7 +19,10 @@ import { join } from 'node:path'; import type { RuntimeEvent } from '@maka/core/runtime-event'; -import type { RuntimeInvocationRecoveryInventoryEntry } from '@maka/core/runtime-event-store'; +import type { + RuntimeInvocationRecoveryInventoryEntry, + RuntimeSessionEventSnapshot, +} from '@maka/core/runtime-event-store'; import type { RuntimeInvocationPageInput, RuntimeInvocationPageResult, @@ -47,6 +50,7 @@ export type RuntimeEventReadPersistence = { }; export interface RuntimeEventReadStore { + readSessionRuntimeSnapshot?(sessionId: string): Promise; listSessionInvocations(sessionId: string): Promise; listInvocationRecoveryInventory( sessionIds: readonly string[], @@ -110,6 +114,8 @@ export async function openRuntimeEventReadPersistence(input: { return { kind: 'sqlite', runtimeEventStore: Object.freeze({ + readSessionRuntimeSnapshot: (sessionId: string) => + store.readSessionRuntimeSnapshot(sessionId), listSessionInvocations: (sessionId: string) => store.listSessionInvocations(sessionId), listInvocationRecoveryInventory: (sessionIds: readonly string[]) => store.listInvocationRecoveryInventory(sessionIds), diff --git a/packages/storage/src/sqlite-runtime-store.ts b/packages/storage/src/sqlite-runtime-store.ts index 8a4077bc6b..c5a707b519 100644 --- a/packages/storage/src/sqlite-runtime-store.ts +++ b/packages/storage/src/sqlite-runtime-store.ts @@ -80,6 +80,7 @@ import { type RuntimeInvocationRecoveryInventoryEntry, type RuntimeRecoveryBundleCommit, type RuntimeRecoveryBundleStore, + type RuntimeSessionEventSnapshot, type RuntimeWorkspaceVersionAuthorityStore, } from '@maka/core/runtime-event-store'; import { @@ -710,6 +711,85 @@ export class SqliteRuntimeStore return this.readRuntimeEventsSync(sessionId, runId); } + async readSessionRuntimeSnapshot(sessionId: string): Promise { + assertRuntimeStorageSafeId(sessionId, 'Invalid session id'); + return this.readTransaction(() => { + const decodedEvents = new Map(); + const openings = this.readInvocationOpeningsSync( + sessionId, + { direction: 'asc' }, + decodedEvents, + ); + const eventsByRun = new Map(); + const durableEventOrdinalById = new Map(); + if (openings.length === 0) { + return { invocations: [], eventsByRun, durableEventOrdinalById }; + } + + // Read each payload once. The opening query already decoded its events; + // reuse those same objects. Global event_seq order also gives each run + // its immutable order and each invocation its first terminal fact. + const rows = this.db + .prepare(` + SELECT event_id, session_id, invocation_id, run_id, turn_id, + CASE WHEN event_kind = 'invocation_opened' THEN '' ELSE payload_json END AS payload_json + FROM runtime_events + WHERE session_id = ? + ORDER BY event_seq ASC, event_id ASC + `) + .all(sessionId) as unknown as RuntimeEventStorageRow[]; + const terminalByInvocation = new Map(); + for (const row of rows) { + const event = decodedEvents.get(row.event_id) ?? decodeRuntimeEventStorageRow(row); + const events = eventsByRun.get(event.runId) ?? []; + events.push(event); + eventsByRun.set(event.runId, events); + if (isTerminalRuntimeEvent(event) && !terminalByInvocation.has(event.invocationId)) { + terminalByInvocation.set(event.invocationId, event); + } + } + + // Only identities and ordinals are needed here, not a second copy of + // every payload. Keep the same join and identity checks as the entry reader. + const ordinals = this.db + .prepare(` + SELECT o.ordinal, e.event_id, e.session_id + FROM runtime_session_event_ordinals o + JOIN runtime_events e ON e.event_id = o.event_id + WHERE o.session_id = ? + ORDER BY o.ordinal ASC + `) + .all(sessionId) as Array<{ ordinal: number; event_id: string; session_id: string }>; + for (const row of ordinals) { + if (!Number.isSafeInteger(row.ordinal) || row.ordinal < 1) { + throw new Error(`Invalid RuntimeEvent Session ordinal for ${sessionId}`); + } + if (row.session_id !== sessionId) { + throw new Error(`RuntimeEvent Session ordinal identity mismatch for ${row.event_id}`); + } + durableEventOrdinalById.set(row.event_id, row.ordinal); + } + + const partialsByRun = new Map(); + for (const partial of this.readRuntimePartialSnapshotsSync(sessionId)) { + const partials = partialsByRun.get(partial.event.runId) ?? []; + partials.push(partial); + partialsByRun.set(partial.event.runId, partials); + } + for (const [runId, partials] of partialsByRun) { + eventsByRun.set( + runId, + mergeRuntimePartialSnapshots(eventsByRun.get(runId) ?? [], partials), + ); + } + const invocations = openings.map((opening) => { + const terminalEvent = terminalByInvocation.get(opening.invocationId); + return { ...opening, ...(terminalEvent ? { terminalEvent } : {}) }; + }); + return { invocations, eventsByRun, durableEventOrdinalById }; + }); + } + private transcriptQuery(): RuntimeTranscriptQuery { return new RuntimeTranscriptQuery(this.db, (sessionId, invocationId) => { // By invocation rather than by run: both shelves key their opening on it, @@ -1092,6 +1172,7 @@ export class SqliteRuntimeStore invocationId?: string; runId?: string; }, + decodedEvents?: Map, ): Omit[] { const order = options.direction === 'desc' ? 'DESC' : 'ASC'; const rows = this.db @@ -1175,6 +1256,7 @@ export class SqliteRuntimeStore if (!opening) { throw new Error(`RuntimeEvent ${event.id} is indexed as an opening fact but is not one`); } + decodedEvents?.set(event.id, event); return { sessionId: event.sessionId, invocationId: event.invocationId, @@ -1401,17 +1483,19 @@ export class SqliteRuntimeStore private readRuntimePartialSnapshotsSync( sessionId: string, - runId: string, + runId?: string, ): RuntimePartialSnapshot[] { const partials = this.db .prepare(` SELECT stream_key, session_id, invocation_id, run_id, turn_id, payload_json, text_content, after_event_id FROM runtime_partial_snapshots - WHERE session_id = ? AND run_id = ? + WHERE session_id = ? ${runId === undefined ? '' : 'AND run_id = ?'} ORDER BY updated_at ASC, stream_key ASC `) - .all(sessionId, runId) as unknown as RuntimePartialStorageRow[]; + .all( + ...(runId === undefined ? [sessionId] : [sessionId, runId]), + ) as unknown as RuntimePartialStorageRow[]; const segmentText = new Map(); const segments = this.db .prepare(` @@ -1419,10 +1503,13 @@ export class SqliteRuntimeStore FROM runtime_partial_segments AS segment INNER JOIN runtime_partial_snapshots AS snapshot ON snapshot.stream_key = segment.stream_key - WHERE snapshot.session_id = ? AND snapshot.run_id = ? + WHERE snapshot.session_id = ? ${runId === undefined ? '' : 'AND snapshot.run_id = ?'} ORDER BY segment.stream_key ASC, segment.segment_seq ASC `) - .iterate(sessionId, runId) as Iterable<{ stream_key: string; text_content: string }>; + .iterate(...(runId === undefined ? [sessionId] : [sessionId, runId])) as Iterable<{ + stream_key: string; + text_content: string; + }>; let streamKey: string | undefined; let chunks: string[] = []; let tail: string[] = [];