diff --git a/apps/desktop/src/main/__tests__/runtime-host-desktop-candidate.test.ts b/apps/desktop/src/main/__tests__/runtime-host-desktop-candidate.test.ts index 0be40cd9e4..bd6e2eb33e 100644 --- a/apps/desktop/src/main/__tests__/runtime-host-desktop-candidate.test.ts +++ b/apps/desktop/src/main/__tests__/runtime-host-desktop-candidate.test.ts @@ -1034,6 +1034,68 @@ test('keeps a restored observation retryable until replacement seeding succeeds' assert.equal(pendingAt >= 0 && readyAt > pendingAt, true); await replacementCandidate.close(); }); // A failed replacement leaves the registry available to the next candidate. +test('keeps replacement Host ready when one transcript observation cannot seed', async () => { + const observations = new RuntimeHostSessionObservationRegistry(); + const firstIpc = ipcHarness(); + const firstHost = connectionHarness('transcript-source', { + sessionId: 'session-1', + }); + const firstCandidate = await createDesktopRuntimeHostCandidate( + firstHost.connection, + deps(firstIpc), + observations, + ); + await firstIpc.invoke('sessions:observe', 'session-1', 'observer-1'); + await firstIpc.invoke('sessions:observe', 'session-2', 'observer-2'); + await firstCandidate.close(); + + const events: Array<{ channel: string; payload: unknown }> = []; + const failingHost = connectionHarness('transcript-seed-failure', { + sessionId: 'session-1', + subscribeFailure: new RuntimeHostOperationError( + 'subscription.open', + 'transcript_unavailable', + 'Session transcript is unavailable', + ), + subscribeFailureSessionId: 'session-1', + subscriptionSnapshots: { + 'session-2': continuitySnapshot({ + session: { ...continuitySnapshot().session, sessionId: 'session-2' }, + rootTurn: { + sessionId: 'session-2', + turnId: 'turn-2', + runId: 'run-2', + status: 'running', + }, + }), + }, + }); + const candidate = await createDesktopRuntimeHostCandidate( + failingHost.connection, + { + ...deps(ipcHarness()), + renderer: { + send(channel, _host, payload) { + events.push({ channel, payload }); + }, + }, + }, + observations, + ); + assert.deepEqual(observations.observationSessionIds(), ['session-1', 'session-2']); + await waitFor(() => events.some(({ channel, payload }) => + channel === 'sessions:event:session-1' && + (payload as { type?: string }).type === 'host_observation_error')); + assert.ok(events.some(({ channel, payload }) => + channel === 'sessions:active-interactions-changed' && + (payload as { sessionId?: string }).sessionId === 'session-2')); + assert.ok(!events.some(({ channel, payload }) => + channel === 'sessions:event:session-2' && + (payload as { type?: string }).type === 'host_observation_error')); + await candidate.close(); + await observations.close(); +}); + test('drops a stale shared Session observation when Guest access is gone', async () => { const observations = new RuntimeHostSessionObservationRegistry(); const firstIpc = ipcHarness(); @@ -1259,8 +1321,10 @@ function connectionHarness( sessionId?: string; revisionAbandon?: 'abandoned' | 'retained'; subscriptionSnapshot?: SessionContinuitySnapshot; + subscriptionSnapshots?: Record; assistantStreams?: readonly SessionAssistantStreamIdentity[]; subscribeFailure?: Error; + subscribeFailureSessionId?: string; runtimeResourcePty?: ReturnType; runtimeResourceUpdate?: ShellRunUpdate; sharedSessionAvailable?: boolean; @@ -1402,7 +1466,9 @@ function connectionHarness( throw new Error(`Unexpected operation: ${operation}`); }, openSessionSubscription: async ({ sessionId }: { sessionId: string }) => { - if (options.subscribeFailure) throw options.subscribeFailure; + if (options.subscribeFailure && + (!options.subscribeFailureSessionId || options.subscribeFailureSessionId === sessionId)) + throw options.subscribeFailure; const subscriptionFrames = new AsyncFrameQueue(); activeSubscriptionFrames = subscriptionFrames; // The Host holds a subscription's frames until the subscriber calls @@ -1435,7 +1501,7 @@ function connectionHarness( ptyListeners.add(listener); return () => ptyListeners.delete(listener); }, - snapshot: options.subscriptionSnapshot ?? { + snapshot: options.subscriptionSnapshots?.[sessionId] ?? options.subscriptionSnapshot ?? { projectionRevision: 1, session: { sessionId }, }, diff --git a/apps/desktop/src/main/runtime-host-desktop-candidate.ts b/apps/desktop/src/main/runtime-host-desktop-candidate.ts index 58d1b48066..b7cdb5302b 100644 --- a/apps/desktop/src/main/runtime-host-desktop-candidate.ts +++ b/apps/desktop/src/main/runtime-host-desktop-candidate.ts @@ -420,6 +420,7 @@ async function restoreSessionObservations(input: { sessionIds(): string[]; announcePending(sessionId: string): void; attach(): Promise; + isRecoverable?(sessionId: string): boolean; }): Promise { const requested = input.sessionIds(); for (const sessionId of requested) input.announcePending(sessionId); @@ -428,7 +429,7 @@ async function restoreSessionObservations(input: { const restoredSet = new Set(restored); const registeredSet = new Set(input.sessionIds()); const failed = requested.filter( - (sessionId) => registeredSet.has(sessionId) && !restoredSet.has(sessionId), + (sessionId) => registeredSet.has(sessionId) && !restoredSet.has(sessionId) && !input.isRecoverable?.(sessionId), ); if (failed.length > 0) { throw new Error(`Failed to restore Session observations: ${failed.join(', ')}`); @@ -768,10 +769,12 @@ export async function createDesktopRuntimeHostCandidate( } } observationsAttached = Boolean(sessionObserver); + const transcriptSeedFailures = new Set(); const restoredSessionIds = await restoreSessionObservations({ sessionIds: () => sessionObservations.observationSessionIds(), announcePending: (sessionId) => sendToRenderer(`sessions:event:${sessionId}`, { type: 'host_observation_pending' }), + isRecoverable: (sessionId) => transcriptSeedFailures.has(sessionId), attach: () => sessionObservations.attach( sessionObserver, (target) => ({ @@ -786,6 +789,13 @@ export async function createDesktopRuntimeHostCandidate( off: target.off.bind(target), }), (missingSessionId) => emitSessionsChanged("deleted", missingSessionId), + (failedSessionId) => { + transcriptSeedFailures.add(failedSessionId); + sendToRenderer(`sessions:event:${failedSessionId}`, { + type: 'host_observation_error', + message: 'Session transcript is unavailable', + }); + }, ), }); for (const sessionId of restoredSessionIds) { diff --git a/apps/desktop/src/main/runtime-host-session-observation-registry.ts b/apps/desktop/src/main/runtime-host-session-observation-registry.ts index 6b03f0bdc9..9c3b5dc2c8 100644 --- a/apps/desktop/src/main/runtime-host-session-observation-registry.ts +++ b/apps/desktop/src/main/runtime-host-session-observation-registry.ts @@ -99,6 +99,14 @@ function isMissingRuntimeHostSessionError(error: unknown): boolean { return error.code === "not_found"; } +function isTranscriptSeedFailure(error: unknown): boolean { + return ( + error instanceof RuntimeHostOperationError && + error.operation === 'subscription.open' && + error.code === 'transcript_unavailable' + ); +} + interface SessionObservationRegistration { readonly sessionId: string; readonly messageAdmissions: boolean; @@ -175,6 +183,7 @@ export class RuntimeHostSessionObservationRegistry { source: SessionObservationSource, bindTarget: ObservationTargetBinding = (target) => target, onSessionMissing?: (sessionId: string) => void, + onTranscriptSeedFailure?: (sessionId: string) => void, ): Promise { this.#assertOpen(); if (this.#source && this.#source !== source) { @@ -224,6 +233,9 @@ export class RuntimeHostSessionObservationRegistry { if (registration.lifecycle === "pending") { this.#deleteRegistration(observerId, registration); } + if (isTranscriptSeedFailure(error)) { + onTranscriptSeedFailure?.(registration.sessionId); + } this.#onError(error); } return undefined; diff --git a/packages/cli/src/pi-transcript.ts b/packages/cli/src/pi-transcript.ts index b72bc4e9db..a17096edf1 100644 --- a/packages/cli/src/pi-transcript.ts +++ b/packages/cli/src/pi-transcript.ts @@ -32,7 +32,7 @@ import type { } from '@maka/core/events'; import { deriveTurnRecords, - isRuntimeSystemNoteKind, + isUserVisibleSessionSystemNote, STEP_LIMIT_NOTICE_TEXT, type StoredMessage, type SystemNoteMessage, @@ -1348,7 +1348,7 @@ function tokenDelta(before: number | undefined, after: number | undefined): numb function systemNoteText(message: SystemNoteMessage): string | undefined { // Retired kinds are still decoded off legacy transcript rows, and none of // them ever had a line here worth reading. - if (!isRuntimeSystemNoteKind(message.kind)) return undefined; + if (!isUserVisibleSessionSystemNote(message.kind)) return undefined; switch (message.kind) { case 'context_compacted': return 'Context compacted to keep this task within the model window.'; @@ -1395,6 +1395,8 @@ function systemNoteText(message: SystemNoteMessage): string | undefined { } case 'step_limit': return STEP_LIMIT_NOTICE_TEXT; + case 'transcript_omitted': + return 'Part of this session transcript was omitted because it exceeded the display limit.'; } } diff --git a/packages/core/src/__tests__/runtime-event.test.ts b/packages/core/src/__tests__/runtime-event.test.ts index 48c8e6c9f8..fae4216514 100644 --- a/packages/core/src/__tests__/runtime-event.test.ts +++ b/packages/core/src/__tests__/runtime-event.test.ts @@ -38,7 +38,12 @@ import { type RuntimeEvent, type RuntimeEventActions, } from '../runtime-event.js'; -import { decodeCanonicalMessage, isRuntimeSystemNoteKind } from '../session.js'; +import { + decodeCanonicalMessage, + isRuntimeSystemNoteKind, + isSystemNoteKind, + isUserVisibleSessionSystemNote, +} from '../session.js'; import { decodeTurnOrigin } from '../turn-origin.js'; /** Minimal valid RuntimeEvent; callers spread overrides on top. */ @@ -146,6 +151,12 @@ test('decodes a released provider dropping note that nothing writes any more', ( assert.equal(isRuntimeSystemNoteKind('context_provider_dropping'), false); }); +test('keeps reader omission notices outside the runtime-written note kinds', () => { + assert.equal(isSystemNoteKind('transcript_omitted'), true); + assert.equal(isUserVisibleSessionSystemNote('transcript_omitted'), true); + assert.equal(isRuntimeSystemNoteKind('transcript_omitted'), false); +}); + test('shares one decoder across all TurnOrigin variants', () => { const origins = [ { kind: 'scheduled_task', scheduledTaskId: 'task-1' }, diff --git a/packages/core/src/session.ts b/packages/core/src/session.ts index fd18305681..4d79857722 100644 --- a/packages/core/src/session.ts +++ b/packages/core/src/session.ts @@ -832,10 +832,10 @@ export function userFacingText(message: Pick { + const fixture = await createFailureFixture({ + registerBackend: (backends) => + backends.register('ai-sdk', (context) => new FakeBackend(context)), + }); + try { + await assert.rejects( + fixture.coordinator.startWorkHubResult( + { + kind: 'workhub_result', + eventId: 'omitted-result', + actionId: 'action', + delegationId: 'delegation', + targetSessionId: fixture.sessionId, + targetTurnId: 'delegated-turn', + }, + async () => { + throw new OmittedTurnResultError('Delegated turn has an omitted transcript result'); + }, + ), + OmittedTurnResultError, + ); + assert.equal(fixture.drainRequested(), false); + } finally { + await fixture.coordinator.close(); + await fixture.messages.close(); + await fixture.dispose(); + } +}); + test('resume query preserves Session-before-activation lock ordering', async () => { const activation = new RuntimePolicyActivationGate(); const capabilities = new HostClientCapabilityCoordinator({ diff --git a/packages/runtime-host/src/__tests__/session-continuity-coordinator.test.ts b/packages/runtime-host/src/__tests__/session-continuity-coordinator.test.ts index d1fd3492a0..bcf0af029c 100644 --- a/packages/runtime-host/src/__tests__/session-continuity-coordinator.test.ts +++ b/packages/runtime-host/src/__tests__/session-continuity-coordinator.test.ts @@ -1184,7 +1184,7 @@ test('open returns a bounded immutable durable tail', async () => { coordinator.close(); }); -test('logs a durable bootstrap failure before reporting unavailable persistence', async (t) => { +test('isolates a durable bootstrap failure without draining the Host', async (t) => { const logs: string[] = []; t.mock.method(console, 'error', (...args: unknown[]) => { logs.push(args.map(String).join(' ')); @@ -1193,11 +1193,14 @@ test('logs a durable bootstrap failure before reporting unavailable persistence' `injected durable bootstrap failure: api_key=sk-secretvalue123\n${'细'.repeat(4096)}`, ); const publicationFailures: unknown[] = []; + const baseReader = transcriptReader([]); + let bootstrapReads = 0; const reader: SessionTranscriptReader = { - ...transcriptReader([]), + ...baseReader, readDurableHighWater: async () => 0, - readDurablePage: async () => { - throw failure; + readDurablePage: async (sessionId, request, project) => { + if (bootstrapReads++ === 0) throw failure; + return baseReader.readDurablePage(sessionId, request, project); }, }; const coordinator = new SessionContinuityCoordinator( @@ -1205,7 +1208,6 @@ test('logs a durable bootstrap failure before reporting unavailable persistence' async () => canonical(), new SessionAdmissionGate(), (error) => { - assert.equal(logs.length, 1, 'log the cause before the publication-failure hook'); publicationFailures.push(error); }, reader, @@ -1221,9 +1223,17 @@ test('logs a durable bootstrap failure before reporting unavailable persistence' ); assert.deepEqual(outcome, { ok: false, - error: { code: 'persistence_failed', message: 'Session transcript is unavailable' }, + error: { code: 'transcript_unavailable', message: 'Session transcript is unavailable' }, }); - assert.deepEqual(publicationFailures, [failure]); + assert.deepEqual(publicationFailures, []); + const retry = await coordinator.handlers['subscription.open']( + { + sessionId: SESSION_ID, + transcript: { kind: 'tail', maxBytes: SESSION_TRANSCRIPT_BOOTSTRAP_MAX_BYTES }, + }, + connectionContext('connection-failed-bootstrap'), + ); + assert.equal(retry.ok, true, 'the same Host can still serve a subsequent open'); assert.equal(logs.length, 1); const prefix = '[runtime-host] subscription.open transcript bootstrap failed: '; const diagnostic = logs[0] ?? ''; diff --git a/packages/runtime-host/src/__tests__/session-transcript-reader.test.ts b/packages/runtime-host/src/__tests__/session-transcript-reader.test.ts index 92624a15fd..6811923926 100644 --- a/packages/runtime-host/src/__tests__/session-transcript-reader.test.ts +++ b/packages/runtime-host/src/__tests__/session-transcript-reader.test.ts @@ -21,6 +21,7 @@ import assert from 'node:assert/strict'; import { mkdtemp, rm } from 'node:fs/promises'; import { tmpdir } from 'node:os'; import { join } from 'node:path'; +import { DatabaseSync } from 'node:sqlite'; import test from 'node:test'; import { seedInvocation, @@ -42,7 +43,11 @@ import { readPageSchema } from '@maka/runtime/read-page'; import { openToolResultArchiveEvidenceReader } from '@maka/storage/tool-result-archive-evidence'; import { foldTurnContribution } from '@maka/storage/session-message-projection'; import type { SessionTurnContribution } from '@maka/storage/execution-stores'; -import { openInteractiveExecutionStoresForWrite } from '@maka/storage/execution-stores'; +import { + openInteractiveExecutionStoresForWrite, + type ExecutionStoresWriter, +} from '@maka/storage/execution-stores'; +import { createSqliteRuntimeStore } from '@maka/storage/sqlite-runtime-store'; import { resolveStorageRoot, tryAcquireInteractiveRootOwner } from '@maka/storage/root-authority'; import { boundedFailureDiagnostic } from '../server/failure-diagnostic.js'; import { @@ -1058,6 +1063,211 @@ test('continues serving durable transcript rows with only soft projection diagno }); }); +for (const hasBefore of [true, false]) + test(`omits an oversized invocation ${hasBefore ? 'between turns' : 'at session start'}`, async () => { + const base = await mkdtemp(join(tmpdir(), 'maka-oversized-transcript-')); + const store = createSqliteRuntimeStore(join(base, 'runtime.sqlite')); + const sessionId = 'oversized-session'; + const append = (runId: string, id: string, overrides: Partial) => + store.appendRuntimeEvent( + sessionId, + runId, + runtimeEvent(sessionId, { + id, + invocationId: runId, + runId, + turnId: `turn-${runId}`, + ...overrides, + }), + ); + try { + if (hasBefore) { + await seedInvocation(store, { + sessionId, + runId: 'before', + turnId: 'turn-before', + openedAt: 0, + }); + await append('before', 'before-text', { + role: 'model', + author: 'agent', + content: { kind: 'text', text: 'before' }, + refs: { storedMessageId: 'before-text' }, + }); + await append('before', 'before-end', { + status: 'completed', + actions: { endInvocation: true }, + }); + } + + await seedInvocation(store, { sessionId, runId: 'large', turnId: 'turn-large', openedAt: 1 }); + // Opening + 8,191 non-presentation events + terminal = 8,193. + // Insert a valid historical ledger in one transaction. Going through the + // production writer would make this read-path regression spend minutes on + // 8,191 separately validated writes. + const db = new DatabaseSync(join(base, 'runtime.sqlite')); + try { + const nextSequence = ( + db + .prepare( + 'SELECT MAX(event_seq) + 1 AS next FROM runtime_events WHERE invocation_id = ?', + ) + .get('large') as { next: number } + ).next; + const nextOrdinal = ( + db + .prepare( + 'SELECT MAX(ordinal) + 1 AS next FROM runtime_session_event_ordinals WHERE session_id = ?', + ) + .get(sessionId) as { next: number } + ).next; + const insertEvent = db.prepare(`INSERT INTO runtime_events + (event_id, session_id, invocation_id, run_id, turn_id, event_seq, event_kind, payload_json, committed_at) + VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)`); + const insertOrdinal = db.prepare(`INSERT INTO runtime_session_event_ordinals + (session_id, ordinal, event_id) VALUES (?, ?, ?)`); + db.exec('BEGIN IMMEDIATE'); + try { + for (let index = 0; index < 8_191; index++) { + const event = runtimeEvent(sessionId, { + id: `large-fact-${index}`, + invocationId: 'large', + runId: 'large', + turnId: 'turn-large', + ts: index + 2, + actions: { stateDelta: { unclaimed: true } }, + }); + insertEvent.run( + event.id, + sessionId, + 'large', + 'large', + 'turn-large', + nextSequence + index, + 'runtime_fact', + JSON.stringify(event), + event.ts, + ); + insertOrdinal.run(sessionId, nextOrdinal + index, event.id); + } + db.exec('COMMIT'); + } catch (error) { + db.exec('ROLLBACK'); + throw error; + } + } finally { + db.close(); + } + await append('large', 'large-end', { status: 'completed', actions: { endInvocation: true } }); + await seedInvocation(store, { sessionId, runId: 'after', turnId: 'turn-after', openedAt: 2 }); + await append('after', 'after-text', { + role: 'model', + author: 'agent', + content: { kind: 'text', text: 'after' }, + refs: { storedMessageId: 'after-text' }, + }); + await append('after', 'after-end', { status: 'completed', actions: { endInvocation: true } }); + + const read = createSessionTranscriptReader({ + stores: { runtimeEventStore: store } as unknown as ExecutionStoresWriter<'interactive'>, + canonicalPermissionOutcomes: { readPermissionOutcome: async () => undefined }, + }); + + const page = await read.readDurableRecords(sessionId, { + direction: 'newer', + maxMessages: 16, + maxStoredBytes: 64 * 1024, + }); + assert.deepEqual( + page.records + .filter(({ message }) => message.type !== 'turn_state') + .map(({ message }) => message.type), + hasBefore ? ['assistant', 'system_note', 'assistant'] : ['system_note', 'assistant'], + ); + const omitted = page.records.find(({ message }) => message.type === 'system_note')?.message; + assert.ok(omitted?.type === 'system_note'); + assert.equal(omitted.kind, 'transcript_omitted'); + assert.equal(omitted.turnId, 'turn-large'); + const inspect = new DatabaseSync(join(base, 'runtime.sqlite')); + const midOrdinal = ( + inspect + .prepare('SELECT ordinal FROM runtime_session_event_ordinals WHERE event_id = ?') + .get('large-fact-100') as { ordinal: number } + ).ordinal; + inspect.close(); + const catchUp = await read.readDurableRecords(sessionId, { + direction: 'newer', + position: midOrdinal * 8, + maxMessages: 16, + maxStoredBytes: 64 * 1024, + }); + assert.ok( + catchUp.records.some( + ({ message }) => message.type === 'system_note' && message.kind === 'transcript_omitted', + ), + ); + await assert.rejects( + createTurnResultReader({ + stores: { runtimeEventStore: store } as unknown as ExecutionStoresWriter<'interactive'>, + canonicalPermissionOutcomes: { readPermissionOutcome: async () => undefined }, + })(sessionId, 'turn-large'), + /omitted transcript result/, + ); + + const older = await read.readDurableRecords(sessionId, { + direction: 'older', + maxMessages: 16, + maxStoredBytes: 64 * 1024, + }); + assert.deepEqual( + older.records + .filter(({ message }) => message.type !== 'turn_state') + .map(({ message }) => message.type), + hasBefore ? ['assistant', 'system_note', 'assistant'] : ['assistant', 'system_note'], + ); + if (!hasBefore) { + let bounded = await read.readDurablePage(sessionId, { + direction: 'older', + maxMessages: 1, + maxBytes: 64 * 1024, + }); + let foundOmission = false; + for (let pageIndex = 0; pageIndex < 8; pageIndex++) { + const message = JSON.parse(bounded.fragments[0]!.data.toString('utf8')) as StoredMessage; + if (message.type === 'system_note' && message.kind === 'transcript_omitted') { + foundOmission = true; + break; + } + const cursor = bounded.next; + assert.ok(cursor, 'a bounded page must lead to the omission'); + bounded = await read.readDurablePage(sessionId, { + direction: 'older', + maxMessages: 1, + maxBytes: 64 * 1024, + throughSequence: bounded.throughSequence, + position: cursor.position, + ...(cursor.byteOffset === null ? {} : { byteOffset: cursor.byteOffset }), + }); + } + assert.ok(foundOmission, 'following the cursor reaches the omitted first invocation'); + } + const bootstrap = await read.readDurablePage(sessionId, { + direction: 'older', + maxMessages: 16, + maxBytes: 64 * 1024, + }); + assert.ok( + bootstrap.fragments.some( + (fragment) => + (JSON.parse(fragment.data.toString('utf8')) as StoredMessage).type === 'system_note', + ), + ); + } finally { + store.close(); + await rm(base, { recursive: true, force: true }); + } + }); + const seed = ( stores: Awaited>, sessionId: string, diff --git a/packages/runtime-host/src/protocol/index.ts b/packages/runtime-host/src/protocol/index.ts index 27e6bffdc5..92ede35b96 100644 --- a/packages/runtime-host/src/protocol/index.ts +++ b/packages/runtime-host/src/protocol/index.ts @@ -104,7 +104,9 @@ export const RUNTIME_HOST_REGISTRATION_SCHEMA_VERSION = 1 as const; export const RUNTIME_HOST_PROTOCOL_VERSION = 0 as const; // Increment when the same protocol version no longer guarantees safe Client-Host // interoperability. Mismatches are rejected before domain commands are admitted. -export const RUNTIME_HOST_COMPATIBILITY_EPOCH = 202 as const; +export const RUNTIME_HOST_COMPATIBILITY_EPOCH = 203 as const; +// 203: Transcript omission notes and subscription.open transcript_unavailable require +// newer clients to decode the note and handle isolated seed failures. // 202: `session.remove.preview` takes a bounded list of Sessions and reports the // child tasks, worktrees and optionally the bytes their removal would delete; // `session.remove` takes `requireArchivedForMs` and may answer `too_recent`. diff --git a/packages/runtime-host/src/protocol/operation-spec.ts b/packages/runtime-host/src/protocol/operation-spec.ts index 6e69d1986c..5a61ff8c8a 100644 --- a/packages/runtime-host/src/protocol/operation-spec.ts +++ b/packages/runtime-host/src/protocol/operation-spec.ts @@ -30,6 +30,7 @@ export type HostOperationErrorCode = | 'session_busy' | 'session_binding_conflict' | 'transcript_preparing' + | 'transcript_unavailable' | 'candidate_set_stale' | 'operation_conflict' | 'capability_unavailable' diff --git a/packages/runtime-host/src/protocol/session-continuity.ts b/packages/runtime-host/src/protocol/session-continuity.ts index f376c04519..ca5c58fc45 100644 --- a/packages/runtime-host/src/protocol/session-continuity.ts +++ b/packages/runtime-host/src/protocol/session-continuity.ts @@ -317,6 +317,7 @@ const SUBSCRIPTION_OPEN_ERRORS = [ 'not_found', 'operation_conflict', 'persistence_failed', + 'transcript_unavailable', 'internal_failure', ] as const; diff --git a/packages/runtime-host/src/server/root-turn-coordinator.ts b/packages/runtime-host/src/server/root-turn-coordinator.ts index 1b3f59f458..87a16fe503 100644 --- a/packages/runtime-host/src/server/root-turn-coordinator.ts +++ b/packages/runtime-host/src/server/root-turn-coordinator.ts @@ -117,6 +117,7 @@ import { type TurnOperationHandlerMap, } from './operation-dispatcher.js'; import { RootAdmissionOwner } from './root-admission-owner.js'; +import { OmittedTurnResultError } from './session-transcript-reader.js'; import { type SessionAdmissionLease, SessionAdmissionGate } from './session-admission-gate.js'; import { type RuntimeSessionForwardedEvent, @@ -3468,6 +3469,7 @@ export class RootTurnCoordinator implements HostedExecutionAuthority { if ( !(error instanceof RuntimeHostedRootConflictError) && !(error instanceof RuntimeHostedRootUnavailableError) && + !(error instanceof OmittedTurnResultError) && !(error instanceof HostedRootAdmissionGateError) && !isShutdownCancelledInteractionAdmission(error) ) { diff --git a/packages/runtime-host/src/server/session-continuity-coordinator.ts b/packages/runtime-host/src/server/session-continuity-coordinator.ts index a25332df3f..23144178ee 100644 --- a/packages/runtime-host/src/server/session-continuity-coordinator.ts +++ b/packages/runtime-host/src/server/session-continuity-coordinator.ts @@ -949,6 +949,7 @@ export class SessionContinuityCoordinator implements SessionContinuityService { | 'operation_conflict' | 'operation_unavailable' | 'persistence_failed' + | 'transcript_unavailable' | 'transcript_preparing'; message: string; } @@ -1031,14 +1032,14 @@ export class SessionContinuityCoordinator implements SessionContinuityService { transcript = created.state; transcriptBootstrap = created.bootstrap; } catch (error) { - // Record the cause before the publication-failure hook can drain the Host. + // A transcript bootstrap read is scoped to this subscription. Keep + // serving other Sessions when this one cannot seed its transcript. console.error( `[runtime-host] subscription.open transcript bootstrap failed: ${boundedFailureDiagnostic(error)}`, ); - this.onPublicationFailure(error); return { ok: false as const, - code: 'persistence_failed' as const, + code: 'transcript_unavailable' as const, message: 'Session transcript is unavailable', }; } diff --git a/packages/runtime-host/src/server/session-transcript-reader.ts b/packages/runtime-host/src/server/session-transcript-reader.ts index 2749bf1372..a7335971e5 100644 --- a/packages/runtime-host/src/server/session-transcript-reader.ts +++ b/packages/runtime-host/src/server/session-transcript-reader.ts @@ -44,6 +44,7 @@ import type { RuntimeTranscriptRun, } from '@maka/storage/execution-stores'; import { foldTurnContribution } from '@maka/storage/session-message-projection'; +import { RuntimeTranscriptOversizedTurnError } from '@maka/storage/runtime-transcript-query'; import type { SessionTurnLandmark } from '../protocol/index.js'; const PERMISSION_OUTCOME_READ_CONCURRENCY = 8; @@ -62,6 +63,15 @@ const TRANSCRIPT_SOURCE_MAX_BYTES = /** How much a page may scan past to fill itself when a projection hides rows. */ const PAGE_HIDDEN_SCAN_MAX_BYTES = 16 * 1024 * 1024; +class TranscriptProjectionLimitError extends Error { + readonly name = 'TranscriptProjectionLimitError'; +} + +/** A delegated result cannot be delivered until its transcript is readable. */ +export class OmittedTurnResultError extends Error { + readonly name = 'OmittedTurnResultError'; +} + export function createSessionTranscriptReader(input: { stores: ExecutionStoresWriter<'interactive'>; canonicalPermissionOutcomes: CanonicalPermissionOutcomeReader; @@ -172,7 +182,14 @@ function createDurableLedgerTranscriptReader(input: { const projectTurn = async ( turn: PendingTranscriptRun, ): Promise<{ sequence: number; message: StoredMessage }[]> => { - const projected = await turn.projection.finish(input.canonicalPermissionOutcomes); + if (turn.oversized) return [omittedTranscriptRecord(turn)]; + let projected: Awaited>; + try { + projected = await turn.projection.finish(input.canonicalPermissionOutcomes); + } catch (error) { + if (error instanceof TranscriptProjectionLimitError) return [omittedTranscriptRecord(turn)]; + throw error; + } const hardDiagnostics = projected.diagnostics.filter(isHardRuntimeEventReadModelDiagnostic); if (hardDiagnostics.length > 0) { // Host diagnostics format the Error stack, so keep the failure's location @@ -242,9 +259,19 @@ function createDurableLedgerTranscriptReader(input: { (turn, events) => { const projection = createTranscriptProjection([turn.invocation]); const ordinals = new Map(); - for (const { event, ordinal } of events) { - ordinals.set(event.id, ordinal); - projection.push(event); + try { + for (const { event, ordinal } of events) { + ordinals.set(event.id, ordinal); + projection.push(event); + } + } catch (error) { + if ( + error instanceof RuntimeTranscriptOversizedTurnError || + error instanceof TranscriptProjectionLimitError + ) { + return { ...turn, projection, ordinals, oversized: true }; + } + throw error; } return { ...turn, projection, ordinals }; }, @@ -321,6 +348,11 @@ function createDurableLedgerTranscriptReader(input: { if (!run) break; if (run.invocation.turnId === turnId) { for (const { message } of await projectTurn(run)) { + if (message.type === 'system_note' && message.kind === 'transcript_omitted') { + throw new OmittedTurnResultError( + `Delegated turn ${turnId} has an omitted transcript result`, + ); + } if (message.type === 'assistant' && message.text.trim()) result = message.text; } } @@ -561,13 +593,13 @@ function ordinalOf(sequence: number): number { function assertTurnPresentationBounded(messages: readonly StoredMessage[]): void { if (messages.length > TRANSCRIPT_TURN_MAX_MESSAGES) { - throw new Error('Session transcript Turn exceeds its message limit'); + throw new TranscriptProjectionLimitError('Session transcript Turn exceeds its message limit'); } let encodedBytes = 0; for (const message of messages) { encodedBytes += Buffer.byteLength(JSON.stringify(message), 'utf8'); if (encodedBytes > TRANSCRIPT_TURN_MAX_BYTES) { - throw new Error('Session transcript Turn exceeds its byte limit'); + throw new TranscriptProjectionLimitError('Session transcript Turn exceeds its byte limit'); } } } @@ -590,7 +622,9 @@ async function readCanonicalPermissionOutcomes( if (!item.outcome) continue; encodedBytes += Buffer.byteLength(JSON.stringify(item.outcome), 'utf8'); if (encodedBytes > TRANSCRIPT_TURN_MAX_BYTES) { - throw new Error('Session permission outcomes exceed the transcript byte limit'); + throw new TranscriptProjectionLimitError( + 'Session permission outcomes exceed the transcript byte limit', + ); } outcomes.set(item.requestId, item.outcome); } @@ -601,6 +635,24 @@ async function readCanonicalPermissionOutcomes( interface PendingTranscriptRun extends RuntimeTranscriptRun { projection: ReturnType; ordinals: Map; + oversized?: boolean; +} + +/** A synthetic read-only row, anchored to this run so both paging directions advance. */ +function omittedTranscriptRecord(turn: RuntimeTranscriptRun): { + sequence: number; + message: StoredMessage; +} { + return { + sequence: turn.firstEventOrdinal * EVENT_SEQUENCE_STRIDE, + message: { + type: 'system_note', + id: `transcript-omitted:${turn.invocation.invocationId}:${turn.firstEventOrdinal}`, + turnId: turn.invocation.turnId, + ts: turn.invocation.openedAt, + kind: 'transcript_omitted', + }, + }; } /** Keep only presentation state while the storage snapshot visits complete facts. */ @@ -618,7 +670,9 @@ function createTranscriptProjection(invocations: readonly RuntimeInvocationRecor messageCount += 1; messageBytes += Buffer.byteLength(JSON.stringify(message)); if (messageCount > TRANSCRIPT_TURN_MAX_MESSAGES || messageBytes > TRANSCRIPT_TURN_MAX_BYTES) - throw new Error('Session transcript projection exceeds its presentation limit'); + throw new TranscriptProjectionLimitError( + 'Session transcript projection exceeds its presentation limit', + ); }, }); return { @@ -633,9 +687,9 @@ function createTranscriptProjection(invocations: readonly RuntimeInvocationRecor : event; sourceBytes += Buffer.byteLength(JSON.stringify(measured)); if (eventCount > TRANSCRIPT_SOURCE_MAX_EVENTS) - throw new Error('RuntimeEvent transcript exceeds its event limit'); + throw new TranscriptProjectionLimitError('RuntimeEvent transcript exceeds its event limit'); if (sourceBytes > TRANSCRIPT_TURN_MAX_BYTES) - throw new Error('RuntimeEvent transcript exceeds its byte limit'); + throw new TranscriptProjectionLimitError('RuntimeEvent transcript exceeds its byte limit'); projector.push(event); }, async finish(reader: CanonicalPermissionOutcomeReader) { diff --git a/packages/storage/package.json b/packages/storage/package.json index 43b40dc190..28fa893c08 100644 --- a/packages/storage/package.json +++ b/packages/storage/package.json @@ -39,6 +39,7 @@ "./root-authority": "./dist/root-authority.js", "./read-image-snapshot-store": "./dist/read-image-snapshot-store.js", "./runtime-event-persistence": "./dist/runtime-event-persistence.js", + "./runtime-transcript-query": "./dist/runtime-transcript-query.js", "./runtime-policy-stores": "./dist/runtime-policy-stores.js", "./scheduled-task-store": "./dist/scheduled-task-store.js", "./quiescent-session-snapshot": "./dist/quiescent-session-snapshot.js", diff --git a/packages/storage/src/__tests__/sqlite-runtime-store.test.ts b/packages/storage/src/__tests__/sqlite-runtime-store.test.ts index 05519450be..b36b1d9c01 100644 --- a/packages/storage/src/__tests__/sqlite-runtime-store.test.ts +++ b/packages/storage/src/__tests__/sqlite-runtime-store.test.ts @@ -370,6 +370,36 @@ describe('SqliteRuntimeStore', () => { }); }); + it('applies the event budget per invocation rather than per Session', async () => { + await withStore(async (store) => { + await appendSettledTurn(store, 1); + await appendSettledTurn(store, 2); + const request = { + direction: 'newer' as const, + throughOrdinal: Number.MAX_SAFE_INTEGER, + maxEvents: 3, + maxBytes: 64_000, + maxRecordBytes: 64_000, + }; + const first = await store.readTranscriptRun( + 'session-1', + { ...request, position: 0 }, + (run, events) => ({ run, events: [...events] }), + ); + assert.equal(first?.events.length, 3); + const second = await store.readTranscriptRun( + 'session-1', + { + ...request, + position: (first?.run.lastOrdinal ?? 0) + 1, + }, + (run, events) => ({ run, events: [...events] }), + ); + assert.equal(second?.events.length, 3); + assert.notEqual(first?.run.invocation.invocationId, second?.run.invocation.invocationId); + }); + }); + it('projects transcript RuntimeEvents from a bounded row iterator', async () => { await withStore(async (store) => { await appendSettledTurn(store, 1); diff --git a/packages/storage/src/runtime-transcript-query.ts b/packages/storage/src/runtime-transcript-query.ts index d47dc239a6..276fc364f4 100644 --- a/packages/storage/src/runtime-transcript-query.ts +++ b/packages/storage/src/runtime-transcript-query.ts @@ -53,6 +53,8 @@ export interface RuntimeTranscriptRun { readonly invocation: RuntimeInvocationRecord; readonly firstOrdinal: number; readonly lastOrdinal: number; + /** First actual event in this traversed stretch; boundary ordinals can be empty. */ + readonly firstEventOrdinal: number; } /** One Turn of a Session transcript: every visible invocation carrying its turnId. */ @@ -307,11 +309,23 @@ export class RuntimeTranscriptQuery { invocationId: seeked.invocation_id, }) as { ordinal: number } | undefined; const older = request.direction === 'older'; + const firstOrdinal = older ? (stop ? stop.ordinal + 1 : 0) : seeked.ordinal; + const lastOrdinal = older ? seeked.ordinal : stop ? stop.ordinal - 1 : request.throughOrdinal; + const firstEvent = this.db + .prepare(` + SELECT MIN(o.ordinal) AS ordinal + FROM runtime_session_event_ordinals o + JOIN runtime_events e ON e.event_id = o.event_id + WHERE o.session_id = ? AND e.invocation_id = ? + AND o.ordinal BETWEEN ? AND ? + `) + .get(sessionId, seeked.invocation_id, firstOrdinal, lastOrdinal) as { ordinal: number }; return project( { invocation: this.invocation(sessionId, seeked.invocation_id), - firstOrdinal: older ? (stop ? stop.ordinal + 1 : 0) : seeked.ordinal, - lastOrdinal: older ? seeked.ordinal : stop ? stop.ordinal - 1 : request.throughOrdinal, + firstOrdinal, + lastOrdinal, + firstEventOrdinal: firstEvent.ordinal, }, this.events(seeked.invocation_id, request), ); diff --git a/packages/storage/src/test-only/memory-execution-runtime.ts b/packages/storage/src/test-only/memory-execution-runtime.ts index dd7adcfbe9..6eddf963a2 100644 --- a/packages/storage/src/test-only/memory-execution-runtime.ts +++ b/packages/storage/src/test-only/memory-execution-runtime.ts @@ -489,13 +489,16 @@ function transcript( return { invocation, firstOrdinal: events.find((e) => e.event.content?.kind === 'invocation_opened')?.ordinal, + firstEventOrdinal: events[0]?.ordinal, lastOrdinal: ending?.ordinal ?? events.at(-1)?.ordinal, events, }; }) .filter( - (i): i is typeof i & { firstOrdinal: number; lastOrdinal: number } => - i.firstOrdinal !== undefined, + ( + i, + ): i is typeof i & { firstOrdinal: number; firstEventOrdinal: number; lastOrdinal: number } => + i.firstOrdinal !== undefined && i.firstEventOrdinal !== undefined, ) .sort((x, y) => x.firstOrdinal - y.firstOrdinal); } @@ -1189,10 +1192,16 @@ export function createMemoryRuntimeStore(a: MemoryExecutionAuthority): Execution (older ? e.ordinal < seeked.ordinal : e.ordinal > seeked.ordinal) && e.event.invocationId !== seeked.event.invocationId, ); + const firstOrdinal = older ? (stop ? stop.ordinal + 1 : 0) : seeked.ordinal; + const lastOrdinal = older ? seeked.ordinal : stop ? stop.ordinal - 1 : throughOrdinal; return { ...visible.get(seeked.event.invocationId)!, - firstOrdinal: older ? (stop ? stop.ordinal + 1 : 0) : seeked.ordinal, - lastOrdinal: older ? seeked.ordinal : stop ? stop.ordinal - 1 : throughOrdinal, + firstOrdinal, + lastOrdinal, + firstEventOrdinal: visible + .get(seeked.event.invocationId)! + .events.find((event) => event.ordinal >= firstOrdinal && event.ordinal <= lastOrdinal)! + .ordinal, }; }); if (!selected) return undefined; diff --git a/packages/ui/src/__tests__/materialize.test.ts b/packages/ui/src/__tests__/materialize.test.ts index b9c8c62d78..03d9dad22a 100644 --- a/packages/ui/src/__tests__/materialize.test.ts +++ b/packages/ui/src/__tests__/materialize.test.ts @@ -278,6 +278,20 @@ describe("materializeTurns message metadata", () => { ); }); + test("labels an omitted transcript invocation without hiding the Turn", () => { + const messages: StoredMessage[] = [{ + type: "system_note", + id: "transcript-omitted:invocation-1:1", + turnId: "turn-1", + ts: 1, + kind: "transcript_omitted", + }]; + assert.equal(materializeTurns(messages, "en")[0]?.notes[0]?.text, + "This part of the task transcript exceeds the display limit and was omitted. Earlier and later records remain available."); + assert.equal(materializeTurns(messages, "zh-CN")[0]?.notes[0]?.text, + "这段任务记录超出显示上限,已省略。前后的记录仍可查看。"); + }); + test("preserves an explicit empty reference projection as the new-format marker", () => { const messages: StoredMessage[] = [ { diff --git a/packages/ui/src/conversation-copy.ts b/packages/ui/src/conversation-copy.ts index 28bdde55b5..40bff6b370 100644 --- a/packages/ui/src/conversation-copy.ts +++ b/packages/ui/src/conversation-copy.ts @@ -329,6 +329,7 @@ export interface ConversationCopy { contextUsageUnavailable: string; contextUsageOpen: string; stepLimit: string; + transcriptOmitted: string; }; }; chat: { @@ -537,6 +538,7 @@ const CONVERSATION_COPY = { contextUsageUnavailable: '暂无用量数据', contextUsageOpen: '打开用量追踪', stepLimit: '已达到本轮工具步骤上限,任务可能尚未完成。发送“继续”即可接着处理。', + transcriptOmitted: '这段任务记录超出显示上限,已省略。前后的记录仍可查看。', }, }, chat: { @@ -661,6 +663,7 @@ const CONVERSATION_COPY = { contextUsageUnavailable: '暫無用量資料', contextUsageOpen: '開啟用量追蹤', stepLimit: '已達到本輪工具步驟上限,任務可能尚未完成。傳送“繼續”即可接著處理。', + transcriptOmitted: '這段任務記錄超出顯示上限,已省略。前後的記錄仍可查看。', }, }, chat: { @@ -782,6 +785,7 @@ const CONVERSATION_COPY = { contextUsageUnavailable: 'No usage data is available for this request.', contextUsageOpen: 'Open usage trace', stepLimit: 'Reached the configured step limit. The task may be incomplete. Send “continue” to resume.', + transcriptOmitted: 'This part of the task transcript exceeds the display limit and was omitted. Earlier and later records remain available.', }, }, chat: { diff --git a/packages/ui/src/materialize.ts b/packages/ui/src/materialize.ts index 57ce751495..686998640a 100644 --- a/packages/ui/src/materialize.ts +++ b/packages/ui/src/materialize.ts @@ -180,6 +180,7 @@ function systemNoteLabel(kind: string, data: unknown, locale: UiLocale): string return copy.contextWindowSuggestion(tokens, declared); } if (kind === "step_limit") return copy.stepLimit; + if (kind === "transcript_omitted") return copy.transcriptOmitted; return kind; }