Luigit
repositories / pi-ext

pi-ext

bugabingas pi extensions

owned by admin

extensions/ultra/journal.ts

Raw
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<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);
		},
	};
}