repositories / pi-ext
pi-ext
bugabingas pi extensions
owned by admin
extensions/ultra/journal.ts
Rawimport { 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<string, unknown>;
/** 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<string, unknown>;
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<JournalEntryData>;
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<string, RunStepResult> = {};
let sessions: Record<string, SubagentSessionRef> = {};
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);
},
};
}