Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
10 changes: 9 additions & 1 deletion packages/core/src/__tests__/recall.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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));
Expand Down
25 changes: 13 additions & 12 deletions packages/core/src/recall.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand All @@ -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 };
}

Expand Down
14 changes: 14 additions & 0 deletions packages/core/src/runtime-event-store.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<string, RuntimeEvent[]>;
readonly durableEventOrdinalById: ReadonlyMap<string, number>;
}

export interface RuntimeEventStore {
/** Canonical stores fail the active run closed on every durable write error. */
readonly durability?: 'best_effort' | 'canonical';
Expand Down Expand Up @@ -210,6 +218,12 @@ export interface RuntimeEventStore {
unsettledOperationIds: readonly string[],
): Promise<void>;
readRuntimeEvents(sessionId: string, runId: string): Promise<RuntimeEvent[]>;
/**
* 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<RuntimeSessionEventSnapshot>;
/** Session-wide immutable append order. */
readSessionRuntimeEventEntries(
sessionId: string,
Expand Down
168 changes: 168 additions & 0 deletions packages/runtime/src/__tests__/recall-ledger-corpus.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -185,7 +186,174 @@ function anchors(result: Awaited<ReturnType<typeof runRecall>>): 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<typeof JSON.parse>) => {
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'));
Expand Down
47 changes: 35 additions & 12 deletions packages/runtime/src/runtime-read-model.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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';
Expand Down Expand Up @@ -78,11 +81,28 @@ export class RuntimeReadModel {
async getSessionView(sessionId: string): Promise<RuntimeReadModelSessionView> {
const diagnostics: RuntimeEventReadModelDiagnostic[] = [];
const inFlightTurnIds = new Set<string>();
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(
Expand All @@ -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', {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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 });
Expand Down
12 changes: 12 additions & 0 deletions packages/storage/src/__tests__/invocation-opening-backfill.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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'],
Expand Down
Loading