import { readFile, writeFile } from "node:fs/promises"; import { createHash } from "node:crypto"; import { dirname, join } from "node:path"; import type { CommonsConfig } from "./config"; import { COMMONS_ARTIFACT_CUSTOM_TYPE, extractSourceText, type CommonsArtifactMetadata } from "./discovery"; import { planInsertion } from "./insertion"; import { loadAssembleSelection, loadJsonlReduce, type ReductionModelAdapter } from "./reduce"; import type { CommonsSource } from "./types"; // v2: reduced artifacts emit a fully valid AssistantMessage (usage/stopReason/ // api/provider/timestamp); v1 artifacts crash Pi's usage accounting on resume. export const REDUCER_VERSION = "3"; export const PROMPT_VERSION = "4"; export interface ReductionOutcome { path: string; source: CommonsSource; warnings: string[]; } /** * Full Pi v3 AssistantMessage for a reduction summary. Pi's usage accounting * reads usage/cost on every assistant message at resume time, so artifacts * must carry the complete shape (validity principle: valid v3 out, never * "v3-ish"). */ function summaryAssistantMessage(text: string, modelId: string | undefined, timestamp: string) { return { role: "assistant", content: [{ type: "text", text }], api: "jsonl-reduce", provider: "jsonl-reduce", model: modelId ?? "extractive-fallback", usage: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, totalTokens: 0, cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 }, }, stopReason: "stop", timestamp: Date.parse(timestamp), }; } export function artifactPath(sourcePath: string, outputSessionId: string): string { return join(dirname(sourcePath), `${outputSessionId}.jsonl`); } /** z=-1: whole-span reduction to the single root record via jsonl-reduce. */ function selectableEntryIds(jsonl: string): string[] { const entries = jsonl.split("\n").filter(Boolean).flatMap((line) => { try { return [JSON.parse(line) as RawEntry]; } catch { return []; } }).filter((entry) => entry.type !== "session"); const projectedIds = new Set(entries.flatMap((entry) => entry.type === "custom" && entry.customType === "chi.commons.projected_title" && entry.data && typeof entry.data === "object" && typeof (entry.data as { sessionInfoEntryId?: unknown }).sessionInfoEntryId === "string" ? [(entry.data as { sessionInfoEntryId: string }).sessionInfoEntryId] : [], )); return entries .filter((entry) => !(entry.type === "custom" && entry.customType === "chi.commons.projected_title")) .filter((entry) => !(entry.type === "session_info" && entry.id && projectedIds.has(entry.id))) .map((entry) => String(entry.id ?? "")) .filter(Boolean); } function deltaJsonl(rawJsonl: string, coveredIds: ReadonlySet): string { const lines = rawJsonl.split("\n").filter(Boolean); const projectedIds = new Set(); for (const line of lines) { try { const entry = JSON.parse(line) as RawEntry; if (entry.type === "custom" && entry.customType === "chi.commons.projected_title" && entry.data && typeof entry.data === "object") { const id = (entry.data as { sessionInfoEntryId?: unknown }).sessionInfoEntryId; if (typeof id === "string") projectedIds.add(id); } } catch { /* malformed tails are ignored by the reducer too */ } } return lines.filter((line) => { try { const entry = JSON.parse(line) as RawEntry; if (entry.type === "session") return true; if (!entry.id || coveredIds.has(entry.id)) return false; if (entry.type === "custom" && entry.customType === "chi.commons.projected_title") return false; if (entry.type === "session_info" && projectedIds.has(entry.id)) return false; return true; } catch { return false; } }).join("\n") + "\n"; } export async function reduceAndPersist(args: { source: CommonsSource; config: CommonsConfig; model?: ReductionModelAdapter; /** Latest compatible immutable root; only uncovered source entries are added. */ previous?: CommonsSource; }): Promise { if (args.source.z !== 0) throw new Error("Only raw z=0 sources can be reduced"); const inputJsonl = await readFile(args.source.path, "utf8"); const reduce = await loadJsonlReduce(); let reductionInput: Parameters[0] = inputJsonl; const incrementalWarnings: string[] = []; if (args.previous?.lineage?.contentHash && args.previous.lineage.sourceSessionIds.length === 1 && args.previous.lineage.sourceSessionIds[0] === args.source.id) { const previousJsonl = await readFile(args.previous.path, "utf8"); const covered = new Set(args.previous.lineage.sourceEntryIds); const delta = deltaJsonl(inputJsonl, covered); const deltaIds = selectableEntryIds(delta); if (deltaIds.length === 0) { return { path: args.previous.path, source: args.previous, warnings: [] }; } try { const assemble = await loadAssembleSelection(); const previousIds = selectableEntryIds(previousJsonl); reductionInput = assemble([ { sessionJsonl: previousJsonl, orderedEntryIds: previousIds, sourceLeafId: previousIds.at(-1)!, artifactHash: args.previous.lineage.contentHash }, { sessionJsonl: delta, orderedEntryIds: deltaIds, sourceLeafId: deltaIds.at(-1)! }, ]); } catch (error) { if (!(error instanceof Error) || !error.message.includes("re-derive from raw leaves")) throw error; incrementalWarnings.push(error.message); } } const result = await reduce(reductionInput, { dialect: "pi-v3", reducerVersion: REDUCER_VERSION, promptVersion: PROMPT_VERSION, scope: "session-span", budget: { maxOutputTokens: Math.max(1, Math.min(1200, args.config.maxInsertTokens - 128)), targetCompression: 0.1, }, model: args.model, }); const path = artifactPath(args.source.path, result.metadata.outputSessionId); await writeFile(path, result.outputJsonl, "utf8"); return { path, warnings: [...incrementalWarnings, ...result.warnings], source: { id: result.metadata.outputSessionId, path, z: -1, title: result.metadata.title, timestamp: result.metadata.timeRange?.endedAt ?? args.source.timestamp, entryCount: 3, approxTokens: result.metadata.outputTokensEstimate, lineage: { sourceSessionIds: [...new Set(result.metadata.sources.map((source) => source.sessionId))], sourceEntryIds: result.metadata.sources.flatMap((source) => source.entryIds), reducerVersion: result.metadata.reducerVersion, promptVersion: result.metadata.promptVersion, model: result.metadata.modelId ?? undefined, contentHash: result.metadata.artifactHash, }, }, }; } interface RawEntry { type?: string; id?: string; timestamp?: string; [key: string]: unknown; } /** How many pair slices a z=1 reduction of this source would produce. */ export function pairCount(entryCount: number): number { return Math.ceil(Math.max(0, entryCount) / 2); } /** * z=1: adjacent pairwise reduction (binary-tree strategy, one level above the * leaves). Each adjacent pair of entries is reduced through jsonl-reduce; the * results assemble into one artifact session of ~N/2 summary records with a * planner metadata entry carrying exact lineage. */ export async function reducePairwiseAndPersist(args: { source: CommonsSource; config: CommonsConfig; model?: ReductionModelAdapter; }): Promise { if (args.source.z !== 0) throw new Error("Only raw z=0 sources can be reduced"); const lines = (await readFile(args.source.path, "utf8")).split("\n").filter((line) => line.trim()); const parsed = lines.map((line) => { try { return JSON.parse(line) as RawEntry; } catch { return null; } }); const headerIndex = parsed.findIndex((entry) => entry?.type === "session"); if (headerIndex === -1) throw new Error("source has no session header"); const headerLine = lines[headerIndex]; const entryLines = lines.filter((_, index) => index !== headerIndex && parsed[index] !== null); const entries = parsed.filter((entry, index) => index !== headerIndex && entry !== null) as RawEntry[]; if (entries.length === 0) throw new Error("source has no entries to reduce"); const pairs: Array<{ jsonl: string; entryIds: string[] }> = []; for (let index = 0; index < entryLines.length; index += 2) { const slice = entryLines.slice(index, index + 2); pairs.push({ jsonl: `${headerLine}\n${slice.join("\n")}\n`, entryIds: entries.slice(index, index + 2).map((entry) => String(entry.id ?? "")).filter(Boolean), }); } const reduce = await loadJsonlReduce(); const perPairTokens = Math.max(48, Math.min(400, Math.floor(args.config.maxInsertTokens / pairs.length))); const warnings: string[] = []; const summaries: Array<{ text: string; reductionKey: string; entryIds: string[]; timestamp: string }> = []; for (const pair of pairs) { const result = await reduce(pair.jsonl, { dialect: "pi-v3", reducerVersion: REDUCER_VERSION, promptVersion: PROMPT_VERSION, scope: "selection", budget: { maxOutputTokens: perPairTokens }, model: args.model, }); warnings.push(...result.warnings); const message = result.outputJsonl.split("\n").filter(Boolean).map((line) => JSON.parse(line)) .find((entry) => entry.type === "message"); summaries.push({ text: String(message?.message?.content?.[0]?.text ?? ""), reductionKey: result.metadata.reductionKey, entryIds: pair.entryIds, timestamp: result.metadata.timeRange?.endedAt ?? args.source.timestamp, }); } const key = createHash("sha256").update(JSON.stringify(summaries.map((s) => s.reductionKey))).digest("hex"); const outputSessionId = `z1_${key.slice(0, 24)}`; const timestamp = summaries.at(-1)?.timestamp ?? args.source.timestamp; const metadata: CommonsArtifactMetadata = { schemaVersion: 1, z: 1, strategy: "pairwise", title: `z=1 ยท ${args.source.title}`.slice(0, 120), sourceSessionIds: [args.source.id], sourceEntryIds: pairs.flatMap((pair) => pair.entryIds), reductionKeys: summaries.map((summary) => summary.reductionKey), }; const metaId = `m${key.slice(0, 7)}`; const outEntries: RawEntry[] = [ { type: "session", version: 3, id: outputSessionId, timestamp }, { type: "session_info", id: `i${key.slice(7, 14)}`, parentId: null, timestamp, name: metadata.title }, { type: "custom", id: metaId, parentId: `i${key.slice(7, 14)}`, timestamp, customType: COMMONS_ARTIFACT_CUSTOM_TYPE, data: metadata }, ]; let parentId = metaId; summaries.forEach((summary, index) => { const id = `p${index}${key.slice(14, 20)}`; outEntries.push({ type: "message", id, parentId, timestamp: summary.timestamp, message: summaryAssistantMessage(summary.text, args.model?.id, summary.timestamp), }); parentId = id; }); const path = artifactPath(args.source.path, outputSessionId); await writeFile(path, `${outEntries.map((entry) => JSON.stringify(entry)).join("\n")}\n`, "utf8"); return { path, warnings, source: { id: outputSessionId, path, z: 1, title: metadata.title, timestamp, entryCount: outEntries.length - 1, approxTokens: Math.ceil(summaries.reduce((total, summary) => total + summary.text.length, 0) / 4), lineage: { sourceSessionIds: metadata.sourceSessionIds, sourceEntryIds: metadata.sourceEntryIds, contentHash: metadata.reductionKeys.join(","), }, }, }; } export async function prepareInsertion(source: CommonsSource, config: CommonsConfig) { const text = extractSourceText(await readFile(source.path, "utf8")); return planInsertion({ text, source, config }); }