import { createHash } from "node:crypto"; import type { WorkflowJournal } from "./engine.ts"; import type { RunStepResult } from "./runner.ts"; import type { WorkflowSpec } from "./spec.ts"; const JOURNAL_TYPE = "ultra-journal"; export interface SubagentSessionRef { id: string; path: string; } interface JournalEntryData { version: 1; id: string; kind: "step" | "session" | "complete"; key?: string; result?: RunStepResult; session?: SubagentSessionRef; } export interface SessionWorkflowJournal extends WorkflowJournal { /** Stable across retries of the same workflow name, cwd, and arguments. */ id: string; getSession(key: string): SubagentSessionRef | undefined; setSession(key: string, session: SubagentSessionRef): void; } export interface SessionEntryLike { type: string; customType?: string; data?: unknown; } export interface SessionWorkflowJournalOptions { cwd: string; spec: WorkflowSpec; args?: Record; /** Effective routes isolate checkpoints when tier settings or the session model change. */ modelRouting?: { model?: string; thinkingLevel?: string }[]; entries: readonly SessionEntryLike[]; appendEntry: (customType: string, data: unknown) => void; } function stableJson(value: unknown): string { if (Array.isArray(value)) return `[${value.map(stableJson).join(",")}]`; if (value && typeof value === "object") { const record = value as Record; return `{${Object.keys(record) .sort() .map((key) => `${JSON.stringify(key)}:${stableJson(record[key])}`) .join(",")}}`; } return JSON.stringify(value); } function journalId(opts: SessionWorkflowJournalOptions): string { return createHash("sha256") .update( stableJson({ cwd: opts.cwd, workflow: opts.spec.name, args: opts.args ?? null, modelRouting: opts.modelRouting ?? null, }), ) .digest("hex") .slice(0, 32); } function parseJournalEntry(value: unknown): JournalEntryData | undefined { if (!value || typeof value !== "object") return undefined; const data = value as Partial; if ( data.version !== 1 || typeof data.id !== "string" || (data.kind !== "step" && data.kind !== "session" && data.kind !== "complete") ) return undefined; if (data.kind === "complete") return data as JournalEntryData; if (typeof data.key !== "string") return undefined; if (data.kind === "session") { if ( !data.session || typeof data.session.id !== "string" || typeof data.session.path !== "string" ) return undefined; return data as JournalEntryData; } if ( !data.result || typeof data.result !== "object" || data.result.ok !== true ) return undefined; return data as JournalEntryData; } export function createSessionWorkflowJournal( opts: SessionWorkflowJournalOptions, ): SessionWorkflowJournal { const id = journalId(opts); let results: Record = {}; let sessions: Record = {}; for (const entry of opts.entries) { if (entry.type !== "custom" || entry.customType !== JOURNAL_TYPE) continue; const data = parseJournalEntry(entry.data); if (!data || data.id !== id) continue; if (data.kind === "complete") { results = {}; sessions = {}; } else if (data.kind === "step") { results[data.key as string] = data.result as RunStepResult; } else { sessions[data.key as string] = data.session as SubagentSessionRef; } } return { id, get(key) { return results[key]; }, set(key, result) { if (!result.ok) return; results[key] = result; opts.appendEntry(JOURNAL_TYPE, { version: 1, id, kind: "step", key, result, } satisfies JournalEntryData); }, getSession(key) { return sessions[key]; }, setSession(key, session) { sessions[key] = session; opts.appendEntry(JOURNAL_TYPE, { version: 1, id, kind: "session", key, session, } satisfies JournalEntryData); }, complete() { results = {}; sessions = {}; opts.appendEntry(JOURNAL_TYPE, { version: 1, id, kind: "complete", } satisfies JournalEntryData); }, }; }