diff --git a/apps/desktop/src/main/runtime-host-client.ts b/apps/desktop/src/main/runtime-host-client.ts index 06f37663bb..bec3cb65ca 100644 --- a/apps/desktop/src/main/runtime-host-client.ts +++ b/apps/desktop/src/main/runtime-host-client.ts @@ -21,11 +21,13 @@ import { createHash, randomUUID } from "node:crypto"; import type { AttachmentRef, ShellRunUpdate } from "@maka/core/events"; import type { PlanSessionState, PlanUserControlInput } from "@maka/core/plan"; import { - decodeStoredMessage, + decodeStoredMessage as decodePersistedStoredMessage, type StoredMessage, type TurnRecord, } from "@maka/core/session"; +import { markPersisted } from "@maka/core/persisted-value"; import type { Task } from "@maka/core/task-ledger"; + import type { ConnectionCatalogSnapshot, ConnectionVersionBasis, @@ -131,6 +133,8 @@ import { type WorkspaceProjection, } from "@maka/runtime-host/protocol"; +const decodeStoredMessage = (value: unknown): StoredMessage => + decodePersistedStoredMessage(markPersisted(value)); const MAX_OPTIMISTIC_ATTEMPTS = 3; const MAX_SESSION_REVISION_ATTEMPTS = 8; const MAX_PRICING_SNAPSHOT_ATTEMPTS = 3; diff --git a/apps/desktop/src/renderer/desktop-transcript-range-store.ts b/apps/desktop/src/renderer/desktop-transcript-range-store.ts index 37e61f339d..de5500bd3a 100644 --- a/apps/desktop/src/renderer/desktop-transcript-range-store.ts +++ b/apps/desktop/src/renderer/desktop-transcript-range-store.ts @@ -18,6 +18,7 @@ */ import { decodeStoredMessage, type StoredMessage } from '@maka/core/session'; +import { markPersisted } from '@maka/core/persisted-value'; import type { DesktopTranscriptBatchPayload, DesktopTranscriptFragment, @@ -308,7 +309,7 @@ export class DesktopTranscriptRangeStore { const encoded = new TextDecoder('utf-8', { fatal: true }).decode(pending.bytes); const message = projectDesktopStoredMessage( { hostId: this.#hostId }, - decodeStoredMessage(JSON.parse(encoded) as unknown), + decodeStoredMessage(markPersisted(JSON.parse(encoded))), ); const projected = JSON.stringify(message); this.#pending.delete(key); diff --git a/packages/cli/src/runtime-host-session-channel.ts b/packages/cli/src/runtime-host-session-channel.ts index 801fbee163..aecf548d21 100644 --- a/packages/cli/src/runtime-host-session-channel.ts +++ b/packages/cli/src/runtime-host-session-channel.ts @@ -17,7 +17,11 @@ * under the License. */ -import { decodeStoredMessage, type StoredMessage } from '@maka/core/session'; +import { + decodeStoredMessage as decodePersistedStoredMessage, + type StoredMessage, +} from '@maka/core/session'; +import { markPersisted } from '@maka/core/persisted-value'; import { type SessionEvent } from '@maka/core/events'; import { createRuntimeHostSessionProjectionSeed, @@ -32,6 +36,7 @@ import { type RuntimeHostConnection, type RuntimeHostSessionSubscription, } from '@maka/runtime-host/client'; + import { InteractionAnsweredSnapshot, InteractionPendingSnapshot, @@ -42,6 +47,8 @@ import { } from '@maka/runtime-host/protocol'; import type { MakaPreparedSessionTurn } from './session-driver.js'; +const decodeStoredMessage = (value: unknown): StoredMessage => + decodePersistedStoredMessage(markPersisted(value)); const MAX_PENDING_FRAMES = 512; const MAX_PENDING_EVENTS_PER_TURN = 1_024; const LAG_REARM_PENDING_EVENTS = MAX_PENDING_EVENTS_PER_TURN / 2; diff --git a/packages/cli/src/runtime-host-session-driver.ts b/packages/cli/src/runtime-host-session-driver.ts index bb324b62c4..d4e21c8f27 100644 --- a/packages/cli/src/runtime-host-session-driver.ts +++ b/packages/cli/src/runtime-host-session-driver.ts @@ -21,17 +21,19 @@ import { randomUUID } from 'node:crypto'; import type { CreateSessionInput } from '@maka/core/runtime-inputs'; import { DEFAULT_SESSION_NAME } from '@maka/core/session-name'; import { - decodeStoredMessage, + decodeStoredMessage as decodePersistedStoredMessage, userFacingText, type SessionSummary, type StoredMessage, } from '@maka/core/session'; +import { markPersisted } from '@maka/core/persisted-value'; import { type ActiveInteractionRequestEvent, type QueueEnqueueOutcome, type SessionEvent, type ShellRunUpdate, } from '@maka/core/events'; + import type { OrchestrationMode } from '@maka/core/orchestration'; import type { PermissionMode } from '@maka/core/permission'; @@ -89,6 +91,8 @@ import { inspectGitCwdChanges, resolveMoveCwd, } from './session-driver-policy.js'; +const decodeStoredMessage = (value: unknown): StoredMessage => + decodePersistedStoredMessage(markPersisted(value)); const MAX_CATALOG_ATTEMPTS = 3; /** Optimistic-control retries for goal pause/resume/clear (mirrors the desktop client). */ diff --git a/packages/core/package.json b/packages/core/package.json index aa51d1cdde..70e987b69b 100644 --- a/packages/core/package.json +++ b/packages/core/package.json @@ -85,6 +85,7 @@ "./model-metadata": "./dist/model-metadata.js", "./model-web-search": "./dist/model-web-search.js", "./model-thinking": "./dist/model-thinking.js", + "./persisted-value": "./dist/persisted-value.js", "./chat-model-choice": "./dist/chat-model-choice.js", "./connection-readiness": "./dist/connection-readiness.js", "./provider-auth": "./dist/provider-auth.js", diff --git a/packages/core/src/__tests__/agent-run-authority.test.ts b/packages/core/src/__tests__/agent-run-authority.test.ts index 650222b165..f7bf6d7afb 100644 --- a/packages/core/src/__tests__/agent-run-authority.test.ts +++ b/packages/core/src/__tests__/agent-run-authority.test.ts @@ -19,7 +19,12 @@ import assert from 'node:assert/strict'; import { test } from 'node:test'; -import { decodeAgentRunHeader, type AgentRunHeader } from '../agent-run.js'; +import { + decodeAgentRunHeader, + decodePersistedAgentRunHeader, + type AgentRunHeader, +} from '../agent-run.js'; +import { markPersisted } from '../persisted-value.js'; test('rejects a Run header with multiple hosted root authorities', () => { assert.throws( @@ -34,14 +39,38 @@ test('rejects a Run header with multiple hosted root authorities', () => { }); test('decodes a released Automation Run as read-only legacy provenance', () => { - const decoded = decodeAgentRunHeader({ - ...runHeader(), - automationId: 'automation-1', - }); + const decoded = decodePersistedAgentRunHeader( + markPersisted({ + ...runHeader(), + automationId: 'automation-1', + }), + ); assert.equal(decoded.legacyAutomationId, 'automation-1'); assert.equal(Object.hasOwn(decoded, 'automationId'), false); }); +test('folds all retired AgentRun values only at the persistence boundary', () => { + const persisted = { + ...runHeader(), + status: 'waiting_permission', + permissionMode: 'execute', + }; + + const decoded = decodePersistedAgentRunHeader(markPersisted(persisted)); + assert.equal(decoded.status, 'waiting_for_user'); + assert.equal(decoded.permissionMode, 'ask'); + + assert.throws(() => decodeAgentRunHeader(persisted), /Invalid AgentRun header schema/); + assert.throws( + () => decodeAgentRunHeader({ ...runHeader(), automationId: 'automation-1' }), + /Invalid AgentRun header schema/, + ); + assert.throws( + () => decodeAgentRunHeader({ ...runHeader(), permissionMode: 'execute' }), + /Invalid AgentRun header schema/, + ); +}); + function runHeader(): AgentRunHeader { return { runId: 'run-1', diff --git a/packages/core/src/__tests__/persisted-value-contract.test.ts b/packages/core/src/__tests__/persisted-value-contract.test.ts new file mode 100644 index 0000000000..0579af238b --- /dev/null +++ b/packages/core/src/__tests__/persisted-value-contract.test.ts @@ -0,0 +1,70 @@ +/* + * 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 assert from 'node:assert/strict'; +import { test } from 'node:test'; +import { markPersisted, type PersistedValue } from '../persisted-value.js'; + +interface ExampleRecord { + readonly id: string; +} + +interface WiderRecord { + readonly id: string; +} + +interface NarrowerRecord extends WiderRecord { + readonly detail: string; +} + +test('persisted values cross into domain types only through an explicit decoder', () => { + const persisted = markPersisted({ id: 'record-1' }); + const decode = (value: PersistedValue): ExampleRecord => + value as unknown as ExampleRecord; + + const assignWithoutDecode = () => { + // @ts-expect-error A persisted value is not a decoded domain record. + const current: ExampleRecord = persisted; + return current; + }; + const passUnknownWithoutMarking = (value: unknown) => { + // @ts-expect-error Unknown input must be marked at a persistence read seam first. + return decode(value); + }; + const passCurrentWithoutMarking = (value: ExampleRecord) => { + // @ts-expect-error Current domain values are not persisted decoder inputs. + return decode(value); + }; + + assert.deepEqual(decode(persisted), { id: 'record-1' }); + assert.equal(typeof assignWithoutDecode, 'function'); + assert.equal(typeof passUnknownWithoutMarking, 'function'); + assert.equal(typeof passCurrentWithoutMarking, 'function'); +}); + +test('persisted values are invariant in their domain type', () => { + const narrower = markPersisted({ id: 'record-1', detail: 'detail' }); + const rejectWidening = () => { + // @ts-expect-error PersistedValue must not widen across related record types. + const wider: PersistedValue = narrower; + return wider; + }; + + assert.equal(typeof rejectWidening, 'function'); +}); diff --git a/packages/core/src/__tests__/runtime-event.test.ts b/packages/core/src/__tests__/runtime-event.test.ts index 4443fad269..120d2eac3a 100644 --- a/packages/core/src/__tests__/runtime-event.test.ts +++ b/packages/core/src/__tests__/runtime-event.test.ts @@ -34,7 +34,7 @@ import { type RuntimeEvent, type RuntimeEventActions, } from '../runtime-event.js'; -import { decodeStoredMessage } from '../session.js'; +import { decodeCanonicalMessage } from '../session.js'; import { decodeTurnOrigin } from '../turn-origin.js'; /** Minimal valid RuntimeEvent; callers spread overrides on top. */ @@ -68,7 +68,7 @@ test('Stored assistant reasoning parts survive recovery decoding', () => { }, }, ]; - const stored = decodeStoredMessage({ + const stored = decodeCanonicalMessage({ type: 'assistant', id: 'message-1', turnId: 'turn-1', @@ -84,7 +84,7 @@ test('Stored assistant reasoning parts survive recovery decoding', () => { }); test('decodes released Automation origins as read-only legacy provenance', () => { - const message = decodeStoredMessage({ + const message = decodeCanonicalMessage({ type: 'user', id: 'message-1', turnId: 'turn-1', @@ -419,7 +419,7 @@ describe('RuntimeEvent content variants', () => { event.content && 'quotes' in event.content ? event.content.quotes?.[0] : undefined, quotes[0], ); - const stored = decodeStoredMessage({ + const stored = decodeCanonicalMessage({ type: 'user', id: 'message-1', turnId: 'turn-1', @@ -443,7 +443,7 @@ describe('RuntimeEvent content variants', () => { assert.notEqual(stored.quotes?.[0], quotes[0]); assert.throws( () => - decodeStoredMessage({ + decodeCanonicalMessage({ type: 'user', id: 'message-1', turnId: 'turn-1', diff --git a/packages/core/src/__tests__/scheduled-task.test.ts b/packages/core/src/__tests__/scheduled-task.test.ts index 88308fd39b..e2082e3223 100644 --- a/packages/core/src/__tests__/scheduled-task.test.ts +++ b/packages/core/src/__tests__/scheduled-task.test.ts @@ -19,6 +19,7 @@ import assert from 'node:assert/strict'; import { describe, it } from 'node:test'; +import { markPersisted } from '../persisted-value.js'; import { computeNextFireAt, decodePersistedScheduledTask, @@ -263,7 +264,7 @@ describe('decodePersistedScheduledTask', () => { const stored = JSON.parse( JSON.stringify(base).replace('"permissionMode":"ask"', '"permissionMode":"execute"'), ) as ScheduledTask; - const decoded = decodePersistedScheduledTask(stored); + const decoded = decodePersistedScheduledTask(markPersisted(stored)); assert.equal( decoded.effect.kind === 'agent_run' ? decoded.effect.execution.permissionMode : undefined, 'ask', @@ -271,11 +272,40 @@ describe('decodePersistedScheduledTask', () => { }); it('returns the same task when nothing needs folding', () => { - assert.equal(decodePersistedScheduledTask(base), base); + assert.equal(decodePersistedScheduledTask(markPersisted(base)), base); }); it('leaves effects without an execution template alone', () => { const notify: ScheduledTask = { ...base, effect: { kind: 'notify', channel: 'local' } }; - assert.equal(decodePersistedScheduledTask(notify), notify); + assert.equal(decodePersistedScheduledTask(markPersisted(notify)), notify); + }); + + it('rejects unknown permission modes in persisted execution templates', () => { + assert.throws( + () => + decodePersistedScheduledTask( + markPersisted({ + ...base, + effect: { + ...base.effect, + execution: { + ...(base.effect.kind === 'agent_run' ? base.effect.execution : {}), + permissionMode: 'future-mode', + }, + }, + }), + ), + /Invalid persisted ScheduledTask permission mode/, + ); + }); + + it('rejects unknown effect kinds in persisted records', () => { + assert.throws( + () => + decodePersistedScheduledTask( + markPersisted({ ...base, effect: { kind: 'future-effect' } }), + ), + /Invalid persisted ScheduledTask effect/, + ); }); }); diff --git a/packages/core/src/__tests__/tool-result-record-schema.test.ts b/packages/core/src/__tests__/tool-result-record-schema.test.ts index 0e38d57fe0..dfa8d0f8f5 100644 --- a/packages/core/src/__tests__/tool-result-record-schema.test.ts +++ b/packages/core/src/__tests__/tool-result-record-schema.test.ts @@ -20,8 +20,13 @@ import assert from 'node:assert/strict'; import { describe, test } from 'node:test'; import type { StoredMessage } from '../session.js'; -import { decodeStoredMessage } from '../session.js'; -import { decodeCanonicalToolResultContent } from '../tool-result-record-schema.js'; +import { decodeCanonicalMessage, decodeStoredMessage } from '../session.js'; +import { markPersisted } from '../persisted-value.js'; +import { + decodeCanonicalToolResultContent, + decodePersistedToolResultContent, +} from '../tool-result-record-schema.js'; +import type { ToolResultContent } from '../events.js'; describe('sandbox denial tool result metadata', () => { test('accepts a canonical text result carrying a sandbox denial signal', () => { @@ -35,7 +40,7 @@ describe('sandbox denial tool result metadata', () => { } as const; assert.deepEqual(decodeCanonicalToolResultContent(result), result); - assert.deepEqual(toolResultContent(decodeStoredMessage(storedToolResult(result))), result); + assert.deepEqual(toolResultContent(decodePersistedMessage(storedToolResult(result))), result); }); test('rejects malformed or widened sandbox denial signals', () => { @@ -70,7 +75,7 @@ describe('sandbox boundary failure tool result metadata', () => { } as const; assert.deepEqual(decodeCanonicalToolResultContent(result), result); - assert.deepEqual(toolResultContent(decodeStoredMessage(storedToolResult(result))), result); + assert.deepEqual(toolResultContent(decodePersistedMessage(storedToolResult(result))), result); }); test('preserves the client capability source for actionable bypass recovery', () => { @@ -84,7 +89,7 @@ describe('sandbox boundary failure tool result metadata', () => { } as const; assert.deepEqual(decodeCanonicalToolResultContent(result), result); - assert.deepEqual(toolResultContent(decodeStoredMessage(storedToolResult(result))), result); + assert.deepEqual(toolResultContent(decodePersistedMessage(storedToolResult(result))), result); }); test('rejects malformed or widened boundary failure signals', () => { @@ -120,7 +125,7 @@ describe('uncertain tool outcome metadata', () => { } as const; assert.deepEqual(decodeCanonicalToolResultContent(result), result); - assert.deepEqual(toolResultContent(decodeStoredMessage(storedToolResult(result))), result); + assert.deepEqual(toolResultContent(decodePersistedMessage(storedToolResult(result))), result); }); test('rejects widened or retryable uncertain outcome signals', () => { @@ -155,18 +160,26 @@ describe('retired permission modes in stored subagent results', () => { } as const; test('folds a legacy mode to its live equivalent instead of returning it verbatim', () => { - const decoded = decodeCanonicalToolResultContent(stored); + const decoded = decodePersistedToolResultContent(markPersisted(stored)); assert.equal(decoded.kind === 'subagent' ? decoded.permissionMode : undefined, 'ask'); assert.deepEqual(decoded, { ...stored, permissionMode: 'ask' }); }); test('folds through the stored-message decoder as well', () => { - assert.deepEqual(toolResultContent(decodeStoredMessage(storedToolResult(stored))), { + assert.deepEqual(toolResultContent(decodePersistedMessage(storedToolResult(stored))), { ...stored, permissionMode: 'ask', }); }); + test('rejects retired values at canonical tool-result and message boundaries', () => { + assert.throws(() => decodeCanonicalToolResultContent(stored), /Invalid tool result content/); + assert.throws( + () => decodeCanonicalMessage(storedToolResult(stored)), + /Invalid tool result content/, + ); + }); + test('leaves a live mode untouched', () => { const live = { ...stored, permissionMode: 'bypass' } as const; assert.deepEqual(decodeCanonicalToolResultContent(live), live); @@ -192,6 +205,10 @@ function storedToolResult(content: unknown) { }; } +function decodePersistedMessage(value: unknown): StoredMessage { + return decodeStoredMessage(markPersisted(value)); +} + function toolResultContent(message: StoredMessage) { if (message.type !== 'tool_result') throw new Error('Expected tool result'); return message.content; diff --git a/packages/core/src/agent-run.ts b/packages/core/src/agent-run.ts index e47bdaa762..730d1ca8d5 100644 --- a/packages/core/src/agent-run.ts +++ b/packages/core/src/agent-run.ts @@ -17,7 +17,12 @@ * under the License. */ -import { decodePersistedPermissionMode, type PermissionMode } from './permission.js'; +import { + decodePersistedPermissionMode, + isPermissionMode, + type PermissionMode, +} from './permission.js'; +import type { PersistedValue } from './persisted-value.js'; import { isCollaborationMode, type CollaborationMode } from './collaboration.js'; import { isAgentSwarmAuthorizationSource, @@ -592,7 +597,14 @@ const AGENT_RUN_EVENT_SHAPE = defineObjectShape()( ['message', 'data'], ); -export function decodeAgentRunHeader(value: unknown): AgentRunHeader { +const RETIRED_AGENT_RUN_STATUSES: Readonly> = { + waiting_permission: 'waiting_for_user', +}; + +export function decodePersistedAgentRunHeader( + persisted: PersistedValue, +): AgentRunHeader { + let value = persisted as unknown; if ( isRecord(value) && value.automationId !== undefined && @@ -601,25 +613,33 @@ export function decodeAgentRunHeader(value: unknown): AgentRunHeader { const { automationId, ...current } = value; value = { ...current, legacyAutomationId: automationId }; } + if (isRecord(value)) { + const status = + typeof value.status === 'string' + ? (RETIRED_AGENT_RUN_STATUSES[value.status] ?? value.status) + : value.status; + const permissionMode = decodePersistedPermissionMode(value.permissionMode); + if (status !== value.status || permissionMode !== value.permissionMode) { + value = { ...value, status, permissionMode }; + } + } + return decodeAgentRunHeader(value); +} + +export function decodeAgentRunHeader(value: unknown): AgentRunHeader { if (!isRecord(value) || !hasExactShape(value, AGENT_RUN_HEADER_SHAPE)) { throw new Error('Invalid AgentRun header schema'); } - const status = - value.status === 'waiting_permission' ? ('waiting_for_user' as const) : value.status; - // Same shape as the status fold above: a run written before a mode was - // retired is old, not malformed, so it decodes to the live equivalent - // rather than making the run unreadable. - const permissionMode = decodePersistedPermissionMode(value.permissionMode); const valid = typeof value.runId === 'string' && typeof value.sessionId === 'string' && typeof value.turnId === 'string' && - (AGENT_RUN_STATUSES as readonly unknown[]).includes(status) && + (AGENT_RUN_STATUSES as readonly unknown[]).includes(value.status) && isPersistedBackendKind(value.backendKind) && typeof value.llmConnectionSlug === 'string' && typeof value.modelId === 'string' && typeof value.cwd === 'string' && - permissionMode !== undefined && + isPermissionMode(value.permissionMode) && (value.collaborationMode === undefined || isCollaborationMode(value.collaborationMode)) && (value.orchestrationMode === undefined || isOrchestrationMode(value.orchestrationMode)) && (value.orchestrationSource === undefined || @@ -663,9 +683,6 @@ export function decodeAgentRunHeader(value: unknown): AgentRunHeader { (value.continuationSource === undefined || isAgentRunContinuationSource(value.continuationSource)); if (!valid) throw new Error('Invalid AgentRun header schema'); - if (status !== value.status || permissionMode !== value.permissionMode) { - return { ...value, status, permissionMode } as unknown as AgentRunHeader; - } return value as unknown as AgentRunHeader; } diff --git a/packages/core/src/persisted-value.ts b/packages/core/src/persisted-value.ts new file mode 100644 index 0000000000..5eee5dc946 --- /dev/null +++ b/packages/core/src/persisted-value.ts @@ -0,0 +1,37 @@ +/* + * 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. + */ + +declare const persistedValueBrand: unique symbol; + +/** + * An untrusted value read from durable state for the domain type `T`. + * + * The function-valued brand keeps `T` invariant, so related domain types + * cannot be substituted across persistence seams. The value has no runtime + * wrapper: a persistence adapter marks the raw value, and a decoder is the + * only intended place that casts it back to `unknown` for inspection. + */ +export type PersistedValue = { + readonly [persistedValueBrand]: (value: T) => T; +}; + +/** Mark an untrusted value at the point where durable state enters the process. */ +export function markPersisted(value: unknown): PersistedValue { + return value as PersistedValue; +} diff --git a/packages/core/src/scheduled-task.ts b/packages/core/src/scheduled-task.ts index 27661f14c1..e741812c34 100644 --- a/packages/core/src/scheduled-task.ts +++ b/packages/core/src/scheduled-task.ts @@ -34,6 +34,7 @@ import { type PermissionMode, } from './permission.js'; import { isBotDeliveryProvider, type BotProvider } from './bot-chat-settings.js'; +import type { PersistedValue } from './persisted-value.js'; export const SCHEDULED_TASK_TITLE_MAX_CHARS = 120; export const SCHEDULED_TASK_INTENT_MAX_CHARS = 8_000; @@ -665,14 +666,33 @@ function fail(message: string): { ok: false; message: string } { * `normalizeCreateScheduledTaskInput`, which validates *new* input and is * deliberately strict. Without this fold a task written before a value was * retired would carry that value straight into execution, where nothing - * recognizes it any more. This is not a schema validator: a record that is - * malformed in any other way stays as stored. + * recognizes it any more. This decoder also rejects unknown effect kinds and + * invalid execution permission modes instead of admitting them as domain data; + * validation of new task input remains the normalizers' responsibility. */ -export function decodePersistedScheduledTask(task: ScheduledTask): ScheduledTask { +export function decodePersistedScheduledTask( + persisted: PersistedValue, +): ScheduledTask { + const value = persisted as unknown; + if (!isObject(value) || !isObject(value.effect) || typeof value.effect.kind !== 'string') { + throw new Error('Invalid persisted ScheduledTask effect'); + } + const task = value as ScheduledTask; const { effect } = task; - if (effect.kind !== 'agent_run') return task; + if (effect.kind === 'notify' || effect.kind === 'session_resume') { + return task; + } + if (effect.kind !== 'agent_run') { + throw new Error('Invalid persisted ScheduledTask effect'); + } + if (!isObject(effect.execution)) { + throw new Error('Invalid persisted ScheduledTask execution template'); + } const permissionMode = decodePersistedPermissionMode(effect.execution.permissionMode); - if (permissionMode === undefined || permissionMode === effect.execution.permissionMode) { + if (permissionMode === undefined) { + throw new Error('Invalid persisted ScheduledTask permission mode'); + } + if (permissionMode === effect.execution.permissionMode) { return task; } return { diff --git a/packages/core/src/session.ts b/packages/core/src/session.ts index d9fc065994..c0cd50434f 100644 --- a/packages/core/src/session.ts +++ b/packages/core/src/session.ts @@ -42,7 +42,11 @@ import { } from './record-schema.js'; import { isPermissionDecisionFields } from './interaction-record-schema.js'; import { isTokenUsageFields, type TokenUsageFields } from './usage-record-schema.js'; -import { decodeCanonicalToolResultContent } from './tool-result-record-schema.js'; +import { + decodeCanonicalToolResultContent, + decodePersistedToolResultContent, +} from './tool-result-record-schema.js'; +import { markPersisted, type PersistedValue } from './persisted-value.js'; import type { SubagentWorkspaceBinding } from './subagent-workspace.js'; import { decodeTurnOrigin, type TurnOrigin } from './turn-origin.js'; @@ -1021,8 +1025,21 @@ const SYSTEM_NOTE_KINDS = new Set([ 'abort', ]); -export function decodeStoredMessage(value: unknown): StoredMessage { - const message = decodeStoredMessageContent(value, decodeCanonicalToolResultContent); +export function decodeCanonicalMessage(value: unknown): StoredMessage { + return decodeMessage(value, decodeCanonicalToolResultContent); +} + +export function decodeStoredMessage(persisted: PersistedValue): StoredMessage { + return decodeMessage(persisted as unknown, (content) => + decodePersistedToolResultContent(markPersisted(content)), + ); +} + +function decodeMessage( + value: unknown, + decodeToolResultContent: (content: unknown) => ToolResultContent, +): StoredMessage { + const message = decodeStoredMessageContent(value, decodeToolResultContent); if (!isRecord(message)) throw new Error('Invalid stored message schema'); switch (message.type) { case 'user': diff --git a/packages/core/src/tool-result-record-schema.ts b/packages/core/src/tool-result-record-schema.ts index a6a4466ec7..cf1499dce7 100644 --- a/packages/core/src/tool-result-record-schema.ts +++ b/packages/core/src/tool-result-record-schema.ts @@ -21,7 +21,8 @@ import { decodeCanonicalShellToolResultContent, isSandboxDenialSignal, } from './shell-run-result.js'; -import { decodePersistedPermissionMode } from './permission.js'; +import { decodePersistedPermissionMode, isPermissionMode } from './permission.js'; +import type { PersistedValue } from './persisted-value.js'; import { isStorageRef, type ToolResultContent } from './events.js'; import { validateSandboxBoundaryExpansion } from './sandbox-boundary.js'; import { @@ -215,21 +216,18 @@ export function decodeCanonicalToolResultContent(value: unknown): ToolResultCont if (!isNonShellToolResultContent(value)) { throw new Error('Invalid tool result content'); } - return foldRetiredPermissionMode(value); + return value; } -/** - * Transcript records written before a permission mode was retired still carry - * the old spelling. The shape validators accept it on purpose — rejecting would - * make the Turn unreadable — so canonicalize it here, at the single exit every - * stored tool result passes through, rather than leaving a value the return - * type forbids. - */ -function foldRetiredPermissionMode(content: ToolResultContent): ToolResultContent { - if (content.kind !== 'subagent') return content; - const permissionMode = decodePersistedPermissionMode(content.permissionMode); - if (permissionMode === undefined || permissionMode === content.permissionMode) return content; - return { ...content, permissionMode }; +export function decodePersistedToolResultContent( + persisted: PersistedValue, +): ToolResultContent { + const value = persisted as unknown; + if (!isRecord(value) || value.kind !== 'subagent') { + return decodeCanonicalToolResultContent(value); + } + const permissionMode = decodePersistedPermissionMode(value.permissionMode); + return decodeCanonicalToolResultContent({ ...value, permissionMode }); } function isNonShellToolResultContent(value: unknown): value is ToolResultContent { @@ -341,7 +339,7 @@ function hasValidSubagentResultFields(value: Record): boolean { typeof value.agentName === 'string' && typeof value.turnId === 'string' && isOptionalString(value.runId) && - decodePersistedPermissionMode(value.permissionMode) !== undefined && + isPermissionMode(value.permissionMode) && typeof value.summary === 'string' && isStringArray(value.artifactIds) && isOptionalFiniteNumber(value.startedAt) && diff --git a/packages/runtime-host/protocol-compatible-changes/session-turns-strict-decoder.json b/packages/runtime-host/protocol-compatible-changes/session-turns-strict-decoder.json new file mode 100644 index 0000000000..5ed2fb236d --- /dev/null +++ b/packages/runtime-host/protocol-compatible-changes/session-turns-strict-decoder.json @@ -0,0 +1,5 @@ +{ + "epoch": 42, + "files": ["packages/runtime-host/src/protocol/session-turns.ts"], + "reason": "Uses the canonical message decoder at the wire boundary without changing the accepted Session turn contribution shape" +} diff --git a/packages/runtime-host/src/__tests__/execution-host.test.ts b/packages/runtime-host/src/__tests__/execution-host.test.ts index 4a2b3cf1c5..874c02a86c 100644 --- a/packages/runtime-host/src/__tests__/execution-host.test.ts +++ b/packages/runtime-host/src/__tests__/execution-host.test.ts @@ -40,7 +40,11 @@ import { canonicalToolArgsHash } from '@maka/core/tool-args-identity'; import type { AgentRunHeader } from '@maka/core/agent-run'; import type { MessageContent } from '@maka/core/events'; import type { ConnectionCatalogEntry } from '@maka/core/runtime-policy'; -import { decodeStoredMessage, type StoredMessage } from '@maka/core/session'; +import { + decodeStoredMessage as decodePersistedStoredMessage, + type StoredMessage, +} from '@maka/core/session'; +import { markPersisted } from '@maka/core/persisted-value'; import type { Task } from '@maka/core/task-ledger'; import type { ScheduledTask } from '@maka/core/scheduled-task'; import { isTerminalRuntimeEvent } from '@maka/core/runtime-event'; @@ -116,6 +120,9 @@ import { withTimeout, } from './fixtures/execution-host-suite.js'; +const decodeStoredMessage = (value: unknown): StoredMessage => + decodePersistedStoredMessage(markPersisted(value)); + test('production Host resumes a Session through the ScheduledTask authority', { timeout: 30_000, }, async () => { diff --git a/packages/runtime-host/src/__tests__/session-subscription-client.test.ts b/packages/runtime-host/src/__tests__/session-subscription-client.test.ts index 1734826ef9..fb1d0a2b40 100644 --- a/packages/runtime-host/src/__tests__/session-subscription-client.test.ts +++ b/packages/runtime-host/src/__tests__/session-subscription-client.test.ts @@ -29,7 +29,11 @@ import { prepareStorageRootControlDirectory, resolveStorageRoot, } from '@maka/storage/root-authority'; -import { decodeStoredMessage } from '@maka/core/session'; +import { + decodeStoredMessage as decodePersistedStoredMessage, + type StoredMessage, +} from '@maka/core/session'; +import { markPersisted } from '@maka/core/persisted-value'; import { connectRuntimeHost, RuntimeHostSubscriptionError, @@ -56,6 +60,8 @@ import { import { FramedTransport } from '../transport/framed-transport.js'; import { frameLocalIpcProtocolMessage } from '../transport/local-ipc-framing.js'; +const decodeStoredMessage = (value: unknown): StoredMessage => + decodePersistedStoredMessage(markPersisted(value)); const PROTOCOL = { min: RUNTIME_HOST_PROTOCOL_VERSION, max: RUNTIME_HOST_PROTOCOL_VERSION, diff --git a/packages/runtime-host/src/protocol/session-turns.ts b/packages/runtime-host/src/protocol/session-turns.ts index 2484581435..c615b7152a 100644 --- a/packages/runtime-host/src/protocol/session-turns.ts +++ b/packages/runtime-host/src/protocol/session-turns.ts @@ -17,7 +17,7 @@ * under the License. */ -import { decodeStoredMessage, type TurnRecord, type TurnStateMessage } from '@maka/core/session'; +import { decodeCanonicalMessage, type TurnRecord, type TurnStateMessage } from '@maka/core/session'; import { truncateUtf8 } from '@maka/core/diagnostic-log'; import { requireCount, @@ -406,7 +406,7 @@ function decodeSessionTurnContribution(value: unknown): SessionTurnContribution 'sequence', 'message', ]); - const message = decodeStoredMessage(state.message); + const message = decodeCanonicalMessage(state.message); if (message.type !== 'turn_state') { throw invalidProtocolFrame('Invalid Session turn state contribution'); } diff --git a/packages/runtime/src/__tests__/conversation-copy.test.ts b/packages/runtime/src/__tests__/conversation-copy.test.ts index 6819c7f362..8b9b4116ed 100644 --- a/packages/runtime/src/__tests__/conversation-copy.test.ts +++ b/packages/runtime/src/__tests__/conversation-copy.test.ts @@ -37,6 +37,7 @@ import { import { archivedToolResultContainsConversationOwnedReferences, cloneConversationRuntimeLedger, + collectConversationCopyLinkedChildReferences, createConversationCopySlice, prepareConversationRuntimeLedgerCopy, rewriteConversationCopyMessage, @@ -75,6 +76,22 @@ test('archived tool-result copy preflight detects conversation-owned references' ), true, ); + assert.equal( + archivedToolResultContainsConversationOwnedReferences( + serialized({ + kind: 'subagent', + agentName: 'Researcher', + turnId: 'retired-turn', + runId: 'retired-run', + status: 'completed', + permissionMode: 'execute', + summary: 'done', + artifactIds: [], + }), + 'session-source', + ), + true, + ); const linkedChildReferences = new Map([ [ 'child-session', @@ -159,6 +176,61 @@ test('archived tool-result copy preflight detects conversation-owned references' ); }); +test('conversation copy discovers linked children in persisted retired tool results', () => { + const result = { + kind: 'subagent', + childSessionId: 'child-session', + agentName: 'Researcher', + turnId: 'child-turn', + runId: 'child-run', + status: 'completed', + permissionMode: 'execute', + summary: 'done', + artifactIds: ['child-artifact'], + }; + const runtimeEvent = { + id: 'event-retired-result', + invocationId: 'invocation-source', + runId: 'run-source', + sessionId: 'session-source', + turnId: 'turn-source', + ts: 1, + partial: false, + role: 'tool', + author: 'tool', + content: { + kind: 'function_response', + id: 'tool-1', + name: 'subagent', + result, + }, + } as RuntimeEvent; + + assert.deepEqual( + collectConversationCopyLinkedChildReferences({ + messages: [], + runtimeEvents: [runtimeEvent], + archivedResults: [JSON.stringify(result)], + }), + [ + { + childSessionId: 'child-session', + runId: 'child-run', + turnId: 'child-turn', + artifactIds: ['child-artifact'], + status: 'completed', + }, + { + childSessionId: 'child-session', + runId: 'child-run', + turnId: 'child-turn', + artifactIds: ['child-artifact'], + status: 'completed', + }, + ], + ); +}); + test('conversation copy slices exact turns on inclusive and exclusive boundaries', () => { const messages = [ { type: 'user', id: 'user-1', turnId: 'turn-1', ts: 1, text: 'first' }, @@ -973,7 +1045,7 @@ test('conversation copy clones one terminal Runtime ledger with new owned identi turnId: 'turn-1', runId: 'run-source', status: 'completed', - permissionMode: 'ask', + permissionMode: 'execute' as never, summary: 'done', artifactIds: ['artifact-deleted'], }, @@ -1242,6 +1314,7 @@ test('conversation copy clones one terminal Runtime ledger with new owned identi ? targetEvents[4].content.result : undefined; const typedResult = decodeCanonicalToolResultContent(typedResultValue); + assert.equal(typedResult.kind === 'subagent' ? typedResult.permissionMode : undefined, 'ask'); assert.deepEqual(typedResult.kind === 'subagent' ? typedResult.artifactIds : undefined, [ 'artifact-target-deleted', ]); diff --git a/packages/runtime/src/__tests__/runtime-event-read-model.test.ts b/packages/runtime/src/__tests__/runtime-event-read-model.test.ts index ca7f35a795..8c698187a4 100644 --- a/packages/runtime/src/__tests__/runtime-event-read-model.test.ts +++ b/packages/runtime/src/__tests__/runtime-event-read-model.test.ts @@ -619,6 +619,43 @@ describe('projectRuntimeEventsToStoredMessages', () => { expect(replay.diagnostics).toEqual([]); }); + test('folds retired permission modes while projecting persisted tool results', () => { + const out = projectRuntimeEventsToStoredMessages( + [ + ev({ + id: 'evt-persisted-subagent-result', + role: 'tool', + author: 'tool', + content: { + kind: 'function_response', + id: 'tool-subagent', + name: 'subagent', + result: { + kind: 'subagent', + agentName: 'Researcher', + turnId: 'child-turn', + runId: 'child-run', + status: 'completed', + permissionMode: 'execute', + summary: 'done', + artifactIds: [], + } as never, + }, + refs: { toolCallId: 'tool-subagent' }, + }), + ], + { runHeaders: [header] }, + ); + + const projected = out.messages.find((message) => message.type === 'tool_result'); + expect( + projected?.type === 'tool_result' && projected.content.kind === 'subagent' + ? projected.content.permissionMode + : undefined, + ).toEqual('ask'); + expect(out.diagnostics).toEqual([]); + }); + test('restores a settled Agent Swarm function response', () => { const result = { kind: 'agent_swarm' as const, diff --git a/packages/runtime/src/conversation-copy.ts b/packages/runtime/src/conversation-copy.ts index 9b725722ae..bff9cc5309 100644 --- a/packages/runtime/src/conversation-copy.ts +++ b/packages/runtime/src/conversation-copy.ts @@ -26,8 +26,9 @@ import type { import type { RuntimeEvent } from '@maka/core/runtime-event'; import type { RuntimeEventStore } from '@maka/core/runtime-event-store'; import type { StorageRef, ToolResultContent } from '@maka/core/events'; +import { markPersisted } from '@maka/core/persisted-value'; import type { StoredMessage } from '@maka/core/session'; -import { decodeCanonicalToolResultContent } from '@maka/core/tool-result-record-schema'; +import { decodePersistedToolResultContent } from '@maka/core/tool-result-record-schema'; import { isEmittedAgentRunEventType, isSessionInlineRun } from '@maka/core/agent-run'; import { TOOL_RECOVERY_DECISION_FACT_KIND } from '@maka/core/tool-recovery-fact'; import { @@ -496,7 +497,7 @@ export function archivedToolResultContainsConversationOwnedReferences( let content: ToolResultContent; try { - content = decodeCanonicalToolResultContent(value); + content = decodePersistedToolResultContent(markPersisted(value)); } catch { return false; } @@ -575,7 +576,9 @@ export function collectConversationCopyLinkedChildReferences(input: { if (isArchivedToolResultPlaceholder(value)) return; try { references.push( - ...conversationCopyLinkedChildReferences(decodeCanonicalToolResultContent(value)), + ...conversationCopyLinkedChildReferences( + decodePersistedToolResultContent(markPersisted(value)), + ), ); } catch { // Opaque tool results have no typed linked-child references. @@ -1184,7 +1187,7 @@ function rewriteRuntimeToolResult( } let content: ToolResultContent; try { - content = decodeCanonicalToolResultContent(value); + content = decodePersistedToolResultContent(markPersisted(value)); } catch { return value; } diff --git a/packages/runtime/src/runtime-event-read-model.ts b/packages/runtime/src/runtime-event-read-model.ts index 6143a940fe..29a16105e7 100644 --- a/packages/runtime/src/runtime-event-read-model.ts +++ b/packages/runtime/src/runtime-event-read-model.ts @@ -21,6 +21,7 @@ import type { AgentRunHeader } from '@maka/core/agent-run'; import type { AssistantStepContentKind, StoredMessage, TurnStatus } 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, @@ -34,7 +35,7 @@ import { isTerminalRuntimeEventStatus, } from '@maka/core/runtime-event'; -import { decodeCanonicalToolResultContent } from '@maka/core/tool-result-record-schema'; +import { decodePersistedToolResultContent } from '@maka/core/tool-result-record-schema'; /** The statuses a settled boundary decision can carry — every status but `pending`. */ type SettledSandboxBoundaryStatus = Exclude< @@ -849,7 +850,9 @@ function projectFunctionResponse( let normalizedResult: ToolResultContent | undefined; if (!archivedPlaceholder) { try { - normalizedResult = decodeCanonicalToolResultContent(compatibleResult); + normalizedResult = decodePersistedToolResultContent( + markPersisted(compatibleResult), + ); } catch (error) { diagnostic( state, diff --git a/packages/storage/src/__tests__/claude-code-session-adapter.test.ts b/packages/storage/src/__tests__/claude-code-session-adapter.test.ts index 9328032b4d..717ec80559 100644 --- a/packages/storage/src/__tests__/claude-code-session-adapter.test.ts +++ b/packages/storage/src/__tests__/claude-code-session-adapter.test.ts @@ -22,7 +22,7 @@ import { mkdir, mkdtemp, rm, writeFile } from 'node:fs/promises'; import { tmpdir } from 'node:os'; import { join } from 'node:path'; import { describe, test } from 'node:test'; -import { decodeStoredMessage, type StoredMessage } from '@maka/core/session'; +import { decodeCanonicalMessage, type StoredMessage } from '@maka/core/session'; import { ClaudeCodeSessionAdapter } from '../claude-code-session-adapter.js'; import { createExternalSessionAdapterRegistry } from '../external-session-adapters.js'; @@ -392,7 +392,7 @@ describe('ClaudeCodeSessionAdapter', () => { const messages = await read(home, 'aaaaaaaa-0000-4000-8000-000000000009'); assert.ok(messages.length > 0); for (const message of messages) { - assert.doesNotThrow(() => decodeStoredMessage(JSON.parse(JSON.stringify(message)))); + assert.doesNotThrow(() => decodeCanonicalMessage(JSON.parse(JSON.stringify(message)))); } }); }); diff --git a/packages/storage/src/__tests__/codex-session-adapter.test.ts b/packages/storage/src/__tests__/codex-session-adapter.test.ts index d22b960daf..16c6f623b5 100644 --- a/packages/storage/src/__tests__/codex-session-adapter.test.ts +++ b/packages/storage/src/__tests__/codex-session-adapter.test.ts @@ -23,7 +23,7 @@ import { tmpdir } from 'node:os'; import { join } from 'node:path'; import { describe, test } from 'node:test'; import { fileURLToPath } from 'node:url'; -import { decodeStoredMessage } from '@maka/core/session'; +import { decodeCanonicalMessage } from '@maka/core/session'; import { CodexSessionAdapter } from '../codex-session-adapter.js'; import { createExternalSessionAdapterRegistry } from '../external-session-adapters.js'; @@ -190,7 +190,7 @@ describe('CodexSessionAdapter', () => { }); assert.equal(session.messages.length, 9); for (const message of session.messages) { - assert.deepEqual(decodeStoredMessage(message), message); + assert.deepEqual(decodeCanonicalMessage(message), message); } assert.deepEqual(session.messages[0], { diff --git a/packages/storage/src/__tests__/session-store.test.ts b/packages/storage/src/__tests__/session-store.test.ts index 71e81abe19..6247e5f7e1 100644 --- a/packages/storage/src/__tests__/session-store.test.ts +++ b/packages/storage/src/__tests__/session-store.test.ts @@ -25,6 +25,7 @@ import { join } from 'node:path'; import { DatabaseSync } from 'node:sqlite'; import { describe, test } from 'node:test'; import type { CreateSessionInput } from '@maka/core/runtime-inputs'; +import type { StoredMessage } from '@maka/core/session'; import { EXTERNAL_SESSION_IMPORT_LOOKUP_MAX_RECENT_SESSION_IDS, EXTERNAL_SESSION_IMPORT_LOOKUP_MAX_SOURCE_IDS, @@ -52,6 +53,75 @@ describe('SQLite SessionStore', () => { } }); + test('folds retired Session and transcript values only on persisted reads', async () => { + const root = await mkdtemp(join(tmpdir(), 'maka-session-persisted-decode-')); + const store = createSessionStore(root); + const currentMessage = { + type: 'tool_result', + id: 'result-1', + turnId: 'turn-1', + ts: 1, + toolUseId: 'call-1', + isError: false, + content: { + kind: 'subagent', + childSessionId: 'child-1', + agentName: 'Explore', + turnId: 'child-turn-1', + status: 'completed', + permissionMode: 'ask', + summary: 'done', + artifactIds: [], + }, + } as const satisfies StoredMessage; + let sessionId: string; + try { + const session = await store.create(makeInput({ permissionMode: 'ask' })); + sessionId = session.id; + await store.appendMessage(session.id, currentMessage); + await assert.rejects( + () => + store.appendMessage(session.id, { + ...currentMessage, + id: 'result-retired', + content: { ...currentMessage.content, permissionMode: 'execute' }, + } as unknown as StoredMessage), + /Invalid tool result content/, + ); + } finally { + await store.close?.(); + } + + const database = new DatabaseSync(join(root, OPERATIONAL_STATE_DATABASE_NAME)); + try { + database.exec(` + UPDATE session_metadata + SET payload_json = json_set(payload_json, '$.permissionMode', 'execute') + WHERE session_id = '${sessionId!}'; + UPDATE session_messages + SET record_json = json_set(record_json, '$.content.permissionMode', 'execute') + WHERE session_id = '${sessionId!}'; + `); + } finally { + database.close(); + } + + const reopened = createSessionStore(root); + try { + assert.equal((await reopened.readHeaderSnapshot(sessionId!)).permissionMode, 'ask'); + const [message] = await reopened.readMessages(sessionId!); + assert.equal( + message?.type === 'tool_result' && message.content.kind === 'subagent' + ? message.content.permissionMode + : undefined, + 'ask', + ); + } finally { + await reopened.close?.(); + await rm(root, { recursive: true, force: true }); + } + }); + test('looks up complete published import counts with bounded newest Session ids', async () => { const root = await mkdtemp(join(tmpdir(), 'maka-session-external-origin-lookup-')); const store = createSessionStore(root); diff --git a/packages/storage/src/__tests__/sqlite-core-execution-store.test.ts b/packages/storage/src/__tests__/sqlite-core-execution-store.test.ts index 9d6ce1acc4..5f0e590fdc 100644 --- a/packages/storage/src/__tests__/sqlite-core-execution-store.test.ts +++ b/packages/storage/src/__tests__/sqlite-core-execution-store.test.ts @@ -67,6 +67,49 @@ describe('SQLite core execution stores', () => { }); }); + test('folds retired AgentRun values only when reading persisted rows', async () => { + await withRoot(async (root) => { + const store = createSqliteAgentRunStore(root); + await store.createRun(runHeader()); + await assert.rejects( + () => + store.createRun({ + ...runHeader({ runId: 'run-retired', turnId: 'turn-retired' }), + permissionMode: 'execute', + } as unknown as AgentRunHeader), + /Invalid AgentRun header schema/, + ); + store.close?.(); + + const database = new DatabaseSync(join(root, 'runtime.sqlite')); + try { + const row = database + .prepare("SELECT record_json AS recordJson FROM core_agent_runs WHERE run_id = 'run-1'") + .get() as { recordJson: string }; + const retired = JSON.parse(row.recordJson) as Record; + retired.status = 'waiting_permission'; + retired.permissionMode = 'execute'; + retired.automationId = 'automation-1'; + database + .prepare("UPDATE core_agent_runs SET record_json = ? WHERE run_id = 'run-1'") + .run(JSON.stringify(retired)); + } finally { + database.close(); + } + + const reopened = createSqliteAgentRunStore(root); + try { + const decoded = await reopened.readRun('session-1', 'run-1'); + assert.equal(decoded.status, 'waiting_for_user'); + assert.equal(decoded.permissionMode, 'ask'); + assert.equal(decoded.legacyAutomationId, 'automation-1'); + assert.equal(Object.hasOwn(decoded, 'automationId'), false); + } finally { + reopened.close?.(); + } + }); + }); + test('advances the model-call high-water index with the authority append', async () => { await withRoot(async (root) => { const store = createSqliteAgentRunStore(root); diff --git a/packages/storage/src/__tests__/sqlite-runtime-store.test.ts b/packages/storage/src/__tests__/sqlite-runtime-store.test.ts index e6544c8056..be1d955a10 100644 --- a/packages/storage/src/__tests__/sqlite-runtime-store.test.ts +++ b/packages/storage/src/__tests__/sqlite-runtime-store.test.ts @@ -777,6 +777,45 @@ describe('SqliteRuntimeStore', () => { }); }); + it('decodes a persisted continuation target without widening new claims', async () => { + await withStore(async (store, dbPath) => { + const claim = continuationClaim(); + await persistImmutablePrefix(store, continuationSourcePrefix()); + await assert.rejects( + () => + store.claimContinuation({ + claim: { + ...claim, + targetRunHeader: { + ...claim.targetRunHeader, + permissionMode: 'execute', + } as unknown as ContinuationClaimV1['targetRunHeader'], + }, + }), + /Invalid AgentRun header schema/, + ); + assert.equal((await store.claimContinuation({ claim })).kind, 'acquired'); + + const database = new DatabaseSync(dbPath); + try { + database.exec(` + UPDATE runtime_continuation_claims + SET target_run_header_json = json_set( + target_run_header_json, + '$.permissionMode', + 'execute' + ) + WHERE claim_id = 'claim-1'; + `); + } finally { + database.close(); + } + + const persisted = await store.readContinuationClaimByBoundary(claim.boundaryDigest); + assert.equal(persisted?.targetRunHeader.permissionMode, 'ask'); + }); + }); + it('rejects a continuation claim whose immediate source boundary is not durable', async () => { await withStore(async (store) => { const claim = continuationClaim(); diff --git a/packages/storage/src/__tests__/sqlite-workflow-store.test.ts b/packages/storage/src/__tests__/sqlite-workflow-store.test.ts index 918d393948..dacb942619 100644 --- a/packages/storage/src/__tests__/sqlite-workflow-store.test.ts +++ b/packages/storage/src/__tests__/sqlite-workflow-store.test.ts @@ -22,6 +22,7 @@ import { mkdtemp, rm } from 'node:fs/promises'; import { tmpdir } from 'node:os'; import { join } from 'node:path'; import { after, describe, test } from 'node:test'; +import { DatabaseSync } from 'node:sqlite'; import { createSqliteDeepResearchStore } from '../deep-research-store.js'; import { openInteractiveScheduledTaskStoreForWrite } from '../scheduled-task-store.js'; import { createSqlitePlanStore } from '../plan-store.js'; @@ -404,6 +405,93 @@ describe('SQLite workflow stores', () => { }); }); + test('folds retired permission modes in tasks and pending fire claims', async () => { + await withRoot(async (root) => { + const now = Date.now(); + const { owner, open } = await scheduledTaskStoreRoot(root); + const store = await open(); + await assert.rejects( + () => + store.create( + { + title: 'Reject retired input', + intentBody: 'run', + schedule: { kind: 'once', runAt: now + 1_000 }, + effect: { + kind: 'agent_run', + execution: { + cwd: '/workspace', + llmConnectionSlug: 'default', + model: 'test-model', + permissionMode: 'execute', + collaborationMode: 'agent', + orchestrationMode: 'default', + }, + }, + createdBy: { kind: 'user' }, + }, + now, + ), + /execution.permissionMode is required/, + ); + const task = await store.create( + { + title: 'Decode retired rows', + intentBody: 'run', + schedule: { kind: 'once', runAt: now + 1_000 }, + effect: { + kind: 'agent_run', + execution: { + cwd: '/workspace', + llmConnectionSlug: 'default', + model: 'test-model', + permissionMode: 'ask', + collaborationMode: 'agent', + orchestrationMode: 'default', + }, + }, + createdBy: { kind: 'user' }, + }, + now, + ); + await store.claimNow(task.id, now); + store.close(); + + const database = new DatabaseSync(join(root, 'runtime.sqlite')); + try { + database.exec(` + UPDATE workflow_scheduled_tasks + SET record_json = json_set(record_json, '$.effect.execution.permissionMode', 'execute'); + UPDATE workflow_scheduled_task_fires + SET record_json = json_set(record_json, '$.task.effect.execution.permissionMode', 'execute'); + `); + } finally { + database.close(); + } + + const reopened = await open(); + try { + const decodedTask = (await reopened.list())[0]; + const decodedClaim = (await reopened.listPendingFires())[0]; + assert.equal( + decodedTask?.effect.kind === 'agent_run' + ? decodedTask.effect.execution.permissionMode + : undefined, + 'ask', + ); + assert.equal( + decodedClaim?.task.effect.kind === 'agent_run' + ? decodedClaim.task.effect.execution.permissionMode + : undefined, + 'ask', + ); + } finally { + reopened.close(); + await owner.close(); + } + }); + }); + test('does not lower maxFires below the task fire count', async () => { await withRoot(async (root) => { const now = Date.now(); diff --git a/packages/storage/src/agent-run-store.ts b/packages/storage/src/agent-run-store.ts index 8da7b31c82..069635e85d 100644 --- a/packages/storage/src/agent-run-store.ts +++ b/packages/storage/src/agent-run-store.ts @@ -24,6 +24,7 @@ import type { DatabaseSync } from 'node:sqlite'; import { decodeAgentRunEvent, decodeAgentRunHeader, + decodeCurrentAgentRunHeader, decodeRuntimeEvent, } from './execution-record-codec.js'; import { immutableSteeringMessageId } from './runtime-event-invariants.js'; @@ -317,7 +318,7 @@ class SqliteAgentRunStore implements DurableAgentRunStore { header: AgentRunHeader, _options: { durable?: boolean } = {}, ): Promise { - const normalized = normalizeAgentRunHeader(header, header.sessionId, header.runId); + const normalized = normalizeCurrentAgentRunHeader(header, header.sessionId, header.runId); this.#lease.transaction('write', () => { const inserted = this.#lease.database .prepare(` @@ -378,7 +379,7 @@ class SqliteAgentRunStore implements DurableAgentRunStore { throw new Error('AgentRun Run Composition is immutable'); } } - const next = normalizeAgentRunHeader( + const next = normalizeCurrentAgentRunHeader( { ...current, ...patch, sessionId, runId }, sessionId, runId, @@ -422,7 +423,7 @@ class SqliteAgentRunStore implements DurableAgentRunStore { if (typeof row.session_id !== 'string' || typeof row.record_json !== 'string') { throw new Error('Invalid SQLite AgentRun identity row'); } - return normalizeAgentRunHeader(JSON.parse(row.record_json), row.session_id, runId); + return decodePersistedAgentRunHeader(JSON.parse(row.record_json), row.session_id, runId); }); return { runs, truncated }; } @@ -447,7 +448,7 @@ class SqliteAgentRunStore implements DurableAgentRunStore { if (typeof row.run_id !== 'string' || typeof row.record_json !== 'string') { throw new Error('Invalid SQLite AgentRun row'); } - return normalizeAgentRunHeader(JSON.parse(row.record_json), sessionId, row.run_id); + return decodePersistedAgentRunHeader(JSON.parse(row.record_json), sessionId, row.run_id); }); return { runs, truncated }; } @@ -503,7 +504,7 @@ class SqliteAgentRunStore implements DurableAgentRunStore { ) { throw new Error('Invalid SQLite AgentRun page row'); } - return normalizeAgentRunHeader(JSON.parse(row.record_json), sessionId, row.run_id); + return decodePersistedAgentRunHeader(JSON.parse(row.record_json), sessionId, row.run_id); }); const last = pageRows.at(-1); return { @@ -532,7 +533,7 @@ class SqliteAgentRunStore implements DurableAgentRunStore { if (typeof row.run_id !== 'string' || typeof row.record_json !== 'string') { throw new Error('Invalid SQLite AgentRun row'); } - return normalizeAgentRunHeader(JSON.parse(row.record_json), sessionId, row.run_id); + return decodePersistedAgentRunHeader(JSON.parse(row.record_json), sessionId, row.run_id); }); } @@ -860,15 +861,29 @@ class SqliteAgentRunStore implements DurableAgentRunStore { } } -function normalizeAgentRunHeader(value: unknown, sessionId: string, runId: string): AgentRunHeader { +function normalizeCurrentAgentRunHeader( + value: unknown, + sessionId: string, + runId: string, +): AgentRunHeader { assertSafeId(sessionId, 'Invalid session id'); assertSafeId(runId, 'Invalid run id'); - return decodeAgentRunHeader(JSON.parse(JSON.stringify(value, sanitizeJson)), { + return decodeCurrentAgentRunHeader(JSON.parse(JSON.stringify(value, sanitizeJson)), { sessionId, runId, }); } +function decodePersistedAgentRunHeader( + value: unknown, + sessionId: string, + runId: string, +): AgentRunHeader { + assertSafeId(sessionId, 'Invalid session id'); + assertSafeId(runId, 'Invalid run id'); + return decodeAgentRunHeader(value, { sessionId, runId }); +} + function readSqliteAgentRun(db: DatabaseSync, sessionId: string, runId: string): AgentRunHeader { const row = db .prepare(` @@ -883,7 +898,7 @@ function readSqliteAgentRun(db: DatabaseSync, sessionId: string, runId: string): throw error; } if (typeof row.record_json !== 'string') throw new Error('Invalid SQLite AgentRun row'); - return normalizeAgentRunHeader(JSON.parse(row.record_json), sessionId, runId); + return decodePersistedAgentRunHeader(JSON.parse(row.record_json), sessionId, runId); } function readSqliteAgentRunEvents( diff --git a/packages/storage/src/execution-record-codec.ts b/packages/storage/src/execution-record-codec.ts index acb472a1f8..d12f76046f 100644 --- a/packages/storage/src/execution-record-codec.ts +++ b/packages/storage/src/execution-record-codec.ts @@ -20,6 +20,7 @@ import { decodeAgentRunEvent as decodeCanonicalAgentRunEvent, decodeAgentRunHeader as decodeCanonicalAgentRunHeader, + decodePersistedAgentRunHeader, type AgentRunEvent, type AgentRunHeader, } from '@maka/core/agent-run'; @@ -29,16 +30,22 @@ import { type RuntimeEvent, } from '@maka/core/runtime-event'; -import { decodeStoredMessage } from '@maka/core/session'; +import { + decodeStoredMessage as decodePersistedStoredMessage, + type StoredMessage, +} from '@maka/core/session'; +import { markPersisted } from '@maka/core/persisted-value'; -export { decodeStoredMessage }; +export function decodeStoredMessage(value: unknown): StoredMessage { + return decodePersistedStoredMessage(markPersisted(value)); +} export function decodeAgentRunHeader( value: unknown, expected: { sessionId: string; runId: string }, ): AgentRunHeader { try { - const header = decodeCanonicalAgentRunHeader(value); + const header = decodePersistedAgentRunHeader(markPersisted(value)); if (header.sessionId !== expected.sessionId || header.runId !== expected.runId) { throw new Error('AgentRun header identity does not match its path'); } @@ -50,6 +57,17 @@ export function decodeAgentRunHeader( } } +export function decodeCurrentAgentRunHeader( + value: unknown, + expected: { sessionId: string; runId: string }, +): AgentRunHeader { + const header = decodeCanonicalAgentRunHeader(value); + if (header.sessionId !== expected.sessionId || header.runId !== expected.runId) { + throw new Error('AgentRun header identity does not match its path'); + } + return header; +} + export function decodeAgentRunEvent( value: unknown, expected: { sessionId: string; runId: string; turnId: string }, diff --git a/packages/storage/src/index.ts b/packages/storage/src/index.ts index 1aaecd2a19..5f3a099322 100644 --- a/packages/storage/src/index.ts +++ b/packages/storage/src/index.ts @@ -25,6 +25,7 @@ export { assertSafeSessionId, createSessionStore, createUserMessage, + decodePersistedSessionHeader, isSafeSessionId, isSessionNotFoundError, normalizeSessionHeader, diff --git a/packages/storage/src/project-catalog.ts b/packages/storage/src/project-catalog.ts index 1ea3177aba..1e2b2f8bae 100644 --- a/packages/storage/src/project-catalog.ts +++ b/packages/storage/src/project-catalog.ts @@ -24,12 +24,13 @@ import { basename, dirname, isAbsolute, join, normalize, relative, resolve, sep import { promisify } from 'node:util'; import type { ProjectLocation, ProjectRecord } from '@maka/core/project'; import type { SessionHeader } from '@maka/core/session'; +import { markPersisted } from '@maka/core/persisted-value'; import { hasEnclosingGitEntry } from './git-entry.js'; import { acquireOperationalStateDatabase, type OperationalStateDatabaseLease, } from './operational-state-store.js'; -import { normalizeSessionHeader } from './session-store.js'; +import { decodePersistedSessionHeader, normalizeSessionHeader } from './session-store.js'; export type { ProjectLocation, ProjectRecord } from '@maka/core/project'; @@ -742,8 +743,8 @@ function reassignProjectSessions( ); const updatedSessionIds: string[] = []; for (const row of rows) { - const header = normalizeSessionHeader( - JSON.parse(row.payload_json) as SessionHeader, + const header = decodePersistedSessionHeader( + markPersisted(JSON.parse(row.payload_json)), row.session_id, ); let patch: Pick | undefined; diff --git a/packages/storage/src/scheduled-task-store.ts b/packages/storage/src/scheduled-task-store.ts index 58f745779c..9f28ec3dcc 100644 --- a/packages/storage/src/scheduled-task-store.ts +++ b/packages/storage/src/scheduled-task-store.ts @@ -34,6 +34,7 @@ import { type ScheduledTaskRun, type ScheduledTaskSchedule, } from '@maka/core/scheduled-task'; +import { markPersisted } from '@maka/core/persisted-value'; import { acquireOperationalStateDatabase, type OperationalStateDatabaseLease, @@ -597,7 +598,9 @@ class SqliteScheduledTaskStore implements ScheduledTaskStore { if (typeof row.record_json !== 'string') { throw new Error(`Invalid scheduled task at row ${index + 1}`); } - return decodePersistedScheduledTask(JSON.parse(row.record_json) as ScheduledTask); + return decodePersistedScheduledTask( + markPersisted(JSON.parse(row.record_json)), + ); }); const claimRows = this.#lease.database .prepare(` @@ -611,7 +614,10 @@ class SqliteScheduledTaskStore implements ScheduledTaskStore { throw new Error(`Invalid scheduled task fire claim at row ${index + 1}`); } const claim = JSON.parse(row.record_json) as ScheduledTaskFireClaim; - return { ...claim, task: decodePersistedScheduledTask(claim.task) }; + return { + ...claim, + task: decodePersistedScheduledTask(markPersisted(claim.task)), + }; }); return { tasks, claims }; } diff --git a/packages/storage/src/session-store.ts b/packages/storage/src/session-store.ts index a85b307f08..a6f71f0e28 100644 --- a/packages/storage/src/session-store.ts +++ b/packages/storage/src/session-store.ts @@ -37,7 +37,7 @@ import { } from './operational-state-store.js'; import { DEFAULT_SESSION_NAME, normalizeUserSessionName } from '@maka/core/session-name'; import { - decodeStoredMessage, + decodeCanonicalMessage, deriveTurnRecords, isSessionBlockedReason, isSessionConversationCopy, @@ -49,7 +49,8 @@ import { } from '@maka/core/session'; import { isCollaborationMode } from '@maka/core/collaboration'; import { isOrchestrationMode } from '@maka/core/orchestration'; -import { decodePersistedPermissionMode } from '@maka/core/permission'; +import { decodePersistedPermissionMode, isPermissionMode } from '@maka/core/permission'; +import type { PersistedValue } from '@maka/core/persisted-value'; import { isSubagentWorkspaceBinding } from '@maka/core/subagent-workspace'; import { WORKSPACE_AUTHORITY_SESSION_ID } from '@maka/core/workspace-version-authority'; import type { @@ -447,7 +448,7 @@ class SqliteSessionStore implements SessionAuthorityStore { throw new Error('Subagent spawn metadata requires createSubagent()'); } const canonicalMessages = messages.map((message) => - decodeStoredMessage(JSON.parse(JSON.stringify(message)) as unknown), + decodeCanonicalMessage(JSON.parse(JSON.stringify(message)) as unknown), ); const header: SessionHeader = { ...buildSessionHeader(this.workspaceRoot, input), @@ -1081,10 +1082,6 @@ export function normalizeSessionHeader( header: SessionHeader, sessionId: string = header.id, ): SessionHeader { - // A retired mode decodes to its live equivalent rather than failing the - // header: such a record is old, not malformed, and rejecting it would make - // the Session unopenable. - const permissionMode = decodePersistedPermissionMode(header.permissionMode); const valid = header.id === sessionId && typeof header.workspaceRoot === 'string' && @@ -1118,7 +1115,7 @@ export function normalizeSessionHeader( typeof header.connectionLocked === 'boolean' && typeof header.model === 'string' && (header.toolProfile === undefined || isSessionToolProfile(header.toolProfile)) && - permissionMode !== undefined && + isPermissionMode(header.permissionMode) && isCollaborationMode(header.collaborationMode) && isOrchestrationMode(header.orchestrationMode) && (header.transcriptLedgerVersion === undefined || @@ -1131,9 +1128,24 @@ export function normalizeSessionHeader( const normalizedName = normalizeSessionName(header.name); if (header.blockedReason === undefined) { const { blockedReason: _blockedReason, ...withoutBlockedReason } = header; - return { ...withoutBlockedReason, name: normalizedName, permissionMode }; + return { ...withoutBlockedReason, name: normalizedName }; } - return { ...header, name: normalizedName, permissionMode }; + return { ...header, name: normalizedName }; +} + +export function decodePersistedSessionHeader( + persisted: PersistedValue, + sessionId?: string, +): SessionHeader { + const header = persisted as unknown as SessionHeader; + const permissionMode = decodePersistedPermissionMode(header.permissionMode); + if (permissionMode === undefined) { + return normalizeSessionHeader(header, sessionId ?? header.id); + } + return normalizeSessionHeader( + permissionMode === header.permissionMode ? header : { ...header, permissionMode }, + sessionId ?? header.id, + ); } function isValidSessionExternalOrigin(origin: SessionHeader['externalOrigin']): boolean { diff --git a/packages/storage/src/sqlite-runtime-store.ts b/packages/storage/src/sqlite-runtime-store.ts index 6a406734ee..4ebc1b4832 100644 --- a/packages/storage/src/sqlite-runtime-store.ts +++ b/packages/storage/src/sqlite-runtime-store.ts @@ -59,6 +59,8 @@ import { import { type ToolRecoveryDecisionFact } from '@maka/core/tool-recovery-fact'; import { canonicalToolArgsHash, stableJsonStringify } from '@maka/core/tool-args-identity'; import { encodeCanonicalRuntimeEvent } from '@maka/core/canonical-runtime-event'; +import { decodePersistedAgentRunHeader, type AgentRunHeader } from '@maka/core/agent-run'; +import { markPersisted } from '@maka/core/persisted-value'; import { scanToolLedger, ToolLedgerCorruptionError, @@ -3718,6 +3720,9 @@ function decodeContinuationClaimRow(row: ContinuationClaimStorageRow): Continuat throw new Error(`Unsupported continuation claim protocol ${row.protocol_version}`); } const boundary = JSON.parse(row.boundary_json) as unknown; + const targetRunHeader = decodePersistedAgentRunHeader( + markPersisted(JSON.parse(row.target_run_header_json)), + ); const claim = decodeContinuationClaim({ protocol: 'continuation_claim_v1', claimId: row.claim_id, @@ -3731,7 +3736,7 @@ function decodeContinuationClaimRow(row: ContinuationClaimStorageRow): Continuat runId: row.target_run_id, turnId: row.target_turn_id, }, - targetRunHeader: JSON.parse(row.target_run_header_json) as unknown, + targetRunHeader, claimedAt: row.claimed_at, }); const source = claim.boundary.segments.at(-1)!; diff --git a/packages/storage/src/sqlite-session-metadata-store.ts b/packages/storage/src/sqlite-session-metadata-store.ts index 27e818c293..a00893b039 100644 --- a/packages/storage/src/sqlite-session-metadata-store.ts +++ b/packages/storage/src/sqlite-session-metadata-store.ts @@ -92,8 +92,10 @@ import { type SessionHeaderPatch, type StoredMessage, type SubagentSessionParent, - decodeStoredMessage, + decodeCanonicalMessage, + decodeStoredMessage as decodePersistedStoredMessage, } from '@maka/core/session'; +import { markPersisted } from '@maka/core/persisted-value'; import { type AgentGraphIntentAdmissionSnapshot, type AgentGraphTimelineMetadataSnapshot, @@ -110,6 +112,7 @@ import { import { type SessionListFilter } from '@maka/core/runtime-inputs'; import { assertSafeSessionId, + decodePersistedSessionHeader, normalizeSessionHeader, SessionNotFoundError, type ExternalSessionImportLookupResult, @@ -145,6 +148,10 @@ const SQLITE_TURN_CONTRIBUTION_MAX_SOURCE_MESSAGES = 1_024; const SQLITE_TURN_CONTRIBUTION_MAX_SOURCE_BYTES = 4 * 1024 * 1024; const SQLITE_TURN_LANDMARK_LEGACY_NEIGHBOR_MESSAGES = 32; +function decodeStoredMessage(value: unknown): StoredMessage { + return decodePersistedStoredMessage(markPersisted(value)); +} + const require = createRequire(import.meta.url); const AGENT_GRAPH_CONTROL_DELETE_TABLES = SQLITE_AGENT_GRAPH_CONTROL_TABLES.filter( (table) => table !== 'agent_graph_epochs', @@ -1335,7 +1342,7 @@ export class SqliteSessionMetadataStore { // through JSON so the stored form matches what the recovery path reads. const encoded = messages.map((message) => { const json = JSON.stringify(message); - const canonical = decodeStoredMessage(JSON.parse(json) as unknown); + const canonical = decodeCanonicalMessage(JSON.parse(json) as unknown); return { message: canonical, json }; }); return this.transaction(() => { @@ -1441,7 +1448,7 @@ export class SqliteSessionMetadataStore { if (messages.length === 0) return; const encoded = messages.map((message) => { const json = JSON.stringify(message); - const canonical = decodeStoredMessage(JSON.parse(json) as unknown); + const canonical = decodeCanonicalMessage(JSON.parse(json) as unknown); return { message: canonical, json }; }); this.transaction(() => { @@ -5274,7 +5281,7 @@ function decodeRecord(row: SessionMetadataRow): SessionMetadataRecord { throw new Error(`Invalid SQLite session metadata record for ${row.session_id}`); } return { - header: normalizeSessionHeader(parsed, row.session_id), + header: decodePersistedSessionHeader(markPersisted(parsed), row.session_id), metadataVersion: row.metadata_version, committedAt: row.committed_at, }; @@ -5593,7 +5600,7 @@ function decodeStoredMessageRow( } try { const parsed = JSON.parse(row.record_json) as unknown; - return decodeStoredMessage(parsed); + return decodeStoredMessage(markPersisted(parsed)); } catch (error) { throw new StoredSessionMessageIncompatibleError(sessionId, sequence, { cause: error, @@ -5963,7 +5970,9 @@ function validateTranscriptRecord( ): void { try { decodeStoredMessage( - JSON.parse(typeof data === 'string' ? data : data.toString('utf8')) as unknown, + markPersisted( + JSON.parse(typeof data === 'string' ? data : data.toString('utf8')), + ), ); } catch (error) { throw new StoredSessionMessageIncompatibleError(sessionId, sequence, {