Luigit
repositories / pi-ext

pi-ext

bugabingas pi extensions

owned by admin

extensions/decay/implementation.ts

Raw
import {
	type Api,
	contentText,
	type Model,
	normalizeContext,
	type ToolCall,
	type Usage,
	uuidv7,
} from "@earendil-works/pi-ai";
import type {
	ExtensionAPI,
	ExtensionContext,
	SessionBeforeCompactEvent,
	SessionCompactEvent,
} from "@earendil-works/pi-coding-agent";
import {
	allocateAtoms,
	CHUNK_TARGET_TOKENS,
	combineUsage,
	createSourceChunks,
	DECAY_DETAILS_VERSION,
	type DecayDetails,
	type DegradedStage,
	estimateKeptTokens,
	findPreviousDecayDetails,
	formatFileAppendix,
	indexRecords,
	isProtectedKind,
	isTransientFailure,
	type MemoryAtom,
	mergeFileLists,
	mergeMemory,
	missingChunks,
	parseDecayDetails,
	parseModelReference,
	type RecordChunk,
	reconcileRecords,
	renderFallbackSummary,
	type SourceChunk,
	validateRecordChunk,
} from "./core.js";
import {
	buildClassifierContext,
	buildSummarizerContext,
	CLASSIFIER_TOOL_NAME,
} from "./prompts.js";
import { CanonicalRecordLog } from "./record-log.js";
import { dbg, span } from "./src/debug.ts";
import {
	resolveSetting,
	type SettingDeclaration,
} from "./src/pi-ext-settings.ts";
import {
	DecayComponent,
	type DecayComponent as DecayComponentType,
	type DecayDisplayState,
} from "./ui.js";

type AnyModel = Model<Api>;

interface DecayRunResult {
	summary: string;
	usage: Usage;
	details: DecayDetails;
	estimatedTokensAfter?: number;
	degraded?: { stage: DegradedStage; error: Error };
}

type RunOutcome =
	| { type: "success"; result: DecayRunResult }
	| { type: "error"; error: Error; cancelled: boolean };

interface RunCallbacks {
	onRecord(record: RecordChunk): void;
	onMerge(
		atoms: readonly MemoryAtom[],
		selected: readonly MemoryAtom[],
		model: AnyModel,
	): void;
	onSummary?(tokens: number): void;
}

interface ClassificationResult {
	usage: Usage;
	model: AnyModel;
}

type DecayFailureStage =
	| "model"
	| "prepare"
	| "classifier"
	| "reconcile"
	| "merge"
	| "summarizer"
	| "cleanup";

class DecayRunError extends Error {
	constructor(
		readonly stage: DecayFailureStage,
		cause: Error,
	) {
		super(cause.message);
		this.name = "DecayRunError";
	}
}

const DECAY_DIAGNOSTIC_ENTRY = "decay-diagnostic";
const CLASSIFIER_RETRIES = 2;
const DECAY_WIDGET = "decay-progress";

// Precedence: trusted project, user, omitted. null explicitly selects the fallback model.
function modelSetting(
	key: string,
): SettingDeclaration<string | null | undefined> {
	return {
		key,
		parse: (raw) =>
			raw === null
				? null
				: typeof raw === "string" && raw.trim()
					? raw.trim()
					: undefined,
		default: undefined,
	};
}
const MODEL_SETTING = modelSetting("decay.model");
const SUMMARIZER_MODEL_SETTING = modelSetting("decay.summarizerModel");

// Invalid values become "", which resolveModel rejects with the fallback message below.
function readModelSetting(
	pi: ExtensionAPI,
	ctx: ExtensionContext,
	declaration: SettingDeclaration<string | null | undefined>,
): string | null | undefined {
	const result = resolveSetting(pi, ctx, declaration);
	return result.ok ? result.value : "";
}

export function createDecayRuntime(pi: ExtensionAPI) {
	let nextRunIndex = 0;
	const sessionBeforeCompact = async (
		event: SessionBeforeCompactEvent,
		ctx: ExtensionContext,
	) => {
		const runIndex = ++nextRunIndex;
		const settings = {
			model: readModelSetting(pi, ctx, MODEL_SETTING),
			summarizerModel: readModelSetting(pi, ctx, SUMMARIZER_MODEL_SETTING),
		};
		const selected = resolveModel(ctx, settings.model);
		if (!selected) {
			dbg?.("compaction.outcome", {
				runIndex,
				outcome: "model_unavailable",
				stage: "model",
			});
			const message =
				settings.model === ""
					? "Decay model setting must be provider/model, null, or omitted; using default compaction"
					: settings.model
						? `Decay model not found: ${settings.model}; using default compaction`
						: "Decay has no current model; using default compaction";
			persistDiagnostic(
				pi,
				"fallback",
				"model",
				settings.model || "current/unknown",
				new Error(message),
			);
			warn(ctx, message);
			return;
		}
		const summarizerFollowsClassifier =
			settings.summarizerModel === null ||
			settings.summarizerModel === undefined;
		const summarizer = summarizerFollowsClassifier
			? selected
			: resolveModel(ctx, settings.summarizerModel);
		if (!summarizer) {
			dbg?.("compaction.outcome", {
				runIndex,
				outcome: "model_unavailable",
				stage: "model",
			});
			const message =
				settings.summarizerModel === ""
					? "Decay summarizerModel setting must be provider/model, null, or omitted; using default compaction"
					: `Decay summarizer model not found: ${settings.summarizerModel}; using default compaction`;
			persistDiagnostic(
				pi,
				"fallback",
				"model",
				settings.summarizerModel || "classifier",
				new Error(message),
			);
			warn(ctx, message);
			return;
		}

		const modelReference = `${selected.provider}/${selected.id}`;
		const summarizerModelReference = `${summarizer.provider}/${summarizer.id}`;
		const signal = event.signal;
		let outcome: RunOutcome;
		if (ctx.mode === "tui") {
			outcome = await runWithWidget(
				event,
				ctx,
				selected,
				summarizer,
				modelReference,
				signal,
				summarizerFollowsClassifier,
				runIndex,
			);
		} else {
			try {
				const result = await runDecay(
					event,
					ctx,
					selected,
					summarizer,
					signal,
					{ onRecord: () => {}, onMerge: () => {} },
					summarizerFollowsClassifier,
					runIndex,
				);
				outcome = { type: "success", result };
			} catch (error) {
				outcome = {
					type: "error",
					error: asError(error),
					cancelled: signal.aborted,
				};
			}
		}

		if (outcome.type === "error") {
			if (outcome.cancelled) {
				dbg?.("compaction.outcome", {
					runIndex,
					outcome: "cancelled",
					stage: failureStage(outcome.error),
				});
				return { cancel: true };
			}
			const stage = failureStage(outcome.error);
			dbg?.("compaction.outcome", { runIndex, outcome: "fallback", stage });
			persistDiagnostic(
				pi,
				"fallback",
				stage,
				stage === "summarizer" ? summarizerModelReference : modelReference,
				outcome.error,
			);
			warn(
				ctx,
				`Decay failed: ${boundedError(outcome.error)}; using default compaction`,
			);
			return;
		}

		if (outcome.result.degraded) {
			dbg?.("compaction.outcome", {
				runIndex,
				outcome: "salvaged",
				stage: "summarizer",
			});
			persistDiagnostic(
				pi,
				"salvage",
				outcome.result.degraded.stage,
				summarizerModelReference,
				outcome.result.degraded.error,
			);
			warn(
				ctx,
				`Decay summarizer failed: ${boundedError(outcome.result.degraded.error)}; rendered classified memory without prose synthesis`,
			);
		} else {
			dbg?.("compaction.outcome", {
				runIndex,
				outcome: "success",
				reportedCost: outcome.result.usage.cost.total,
			});
		}

		return {
			compaction: {
				summary: outcome.result.summary,
				firstKeptEntryId: event.preparation.firstKeptEntryId,
				tokensBefore: event.preparation.tokensBefore,
				estimatedTokensAfter: outcome.result.estimatedTokensAfter,
				usage: outcome.result.usage,
				details: outcome.result.details,
			},
		};
	};

	const sessionCompact = (
		event: SessionCompactEvent,
		ctx: ExtensionContext,
	) => {
		if (!event.fromExtension) {
			dbg?.("compaction.outcome", { outcome: "ignored" });
			return;
		}
		const details = parseDecayDetails(event.compactionEntry.details)?.decay;
		if (!details) {
			dbg?.("compaction.outcome", { outcome: "ignored" });
			return;
		}
		ctx.ui.notify(
			`Decay · ${modelLabel(details.model)} · ${details.chunkCount}/${details.chunkCount} classified · ${details.estimatedKeptTokens === undefined ? "unknown" : `~${details.estimatedKeptTokens.toLocaleString("en-US")}`} kept (target ${details.keepRecentTokens.toLocaleString("en-US")}) · ${(details.elapsedMs / 1_000).toFixed(1)}s${details.degraded ? " · salvaged" : ""}`,
			details.degraded ? "warning" : "info",
		);
	};

	return { sessionBeforeCompact, sessionCompact };
}

export type DecayRuntime = ReturnType<typeof createDecayRuntime>;

function resolveModel(
	ctx: ExtensionContext,
	configured: string | null | undefined,
): AnyModel | undefined {
	if (configured === "" || configured === null || configured === undefined) {
		return configured === "" ? undefined : ctx.model;
	}
	const parsed = parseModelReference(configured);
	if (!parsed) return undefined;
	return (
		ctx.modelRegistry
			.getAvailable()
			.find(
				(model) =>
					model.provider === parsed.provider && model.id === parsed.modelId,
			) ?? ctx.model
	);
}

async function runWithWidget(
	event: SessionBeforeCompactEvent,
	ctx: ExtensionContext,
	model: AnyModel,
	summarizerModel: AnyModel,
	modelReference: string,
	signal: AbortSignal,
	summarizerFollowsClassifier: boolean,
	runIndex: number,
): Promise<RunOutcome> {
	let component: DecayComponentType | undefined;
	let widgetShown = false;
	const state: DecayDisplayState = {
		model: modelReference,
		inputTokens: event.preparation.tokensBefore,
		keptTokens: estimateKeptTokens(event),
		totalChunks: previewChunkCount(event),
		phase: "classify",
		records: [],
		startedAt: Date.now(),
	};
	try {
		ctx.ui.setWidget(
			DECAY_WIDGET,
			(tui, theme) => {
				component = new DecayComponent(tui, theme, state);
				return component;
			},
			{ placement: "aboveEditor" },
		);
		widgetShown = true;
		const result = await runDecay(
			event,
			ctx,
			model,
			summarizerModel,
			signal,
			{
				onRecord: (record) => component?.record(record),
				onMerge: (atoms, selectedAtoms, activeModel) => {
					state.model = `${activeModel.provider}/${activeModel.id}`;
					component?.merge(atoms, selectedAtoms);
				},
				onSummary: (tokens) => component?.setSummary(tokens),
			},
			summarizerFollowsClassifier,
			runIndex,
		);
		return { type: "success", result };
	} catch (error) {
		return {
			type: "error",
			error: asError(error),
			cancelled: signal.aborted,
		};
	} finally {
		component?.dispose();
		if (widgetShown) ctx.ui.setWidget(DECAY_WIDGET, undefined);
	}
}

function previewChunkCount(event: SessionBeforeCompactEvent): number {
	const previous = findPreviousDecayDetails(event.branchEntries);
	return createSourceChunks({
		messages: event.preparation.messagesToSummarize,
		turnPrefixMessages: event.preparation.turnPrefixMessages,
		legacySummary: previous ? undefined : event.preparation.previousSummary,
	}).length;
}

async function runDecay(
	event: SessionBeforeCompactEvent,
	ctx: ExtensionContext,
	model: AnyModel,
	summarizerModel: AnyModel,
	signal: AbortSignal,
	callbacks: RunCallbacks,
	summarizerFollowsClassifier = false,
	runIndex = 0,
): Promise<DecayRunResult> {
	const startedAt = Date.now();
	let stage: DecayFailureStage = "prepare";
	let log: CanonicalRecordLog | undefined;
	const currentModel = ctx.model;
	let result: DecayRunResult | undefined;
	let failure: DecayRunError | undefined;
	try {
		const previous = findPreviousDecayDetails(event.branchEntries);
		const chunks = createSourceChunks({
			messages: event.preparation.messagesToSummarize,
			turnPrefixMessages: event.preparation.turnPrefixMessages,
			legacySummary: previous ? undefined : event.preparation.previousSummary,
		});
		if (chunks.length === 0) throw new Error("Decay has no source chunks");
		if (dbg) {
			dbg("compaction.plan", {
				runIndex,
				chunkCount: chunks.length,
				sourceTokens: chunks.reduce(
					(sum, chunk) => sum + chunk.estimatedTokens,
					0,
				),
				largestChunkTokens: chunks.reduce(
					(largest, chunk) => Math.max(largest, chunk.estimatedTokens),
					0,
				),
				overTargetChunks: chunks.filter(
					(chunk) => chunk.estimatedTokens > CHUNK_TARGET_TOKENS,
				).length,
				previousAtoms: previous?.atoms.length ?? 0,
				legacySummary: chunks[0]?.source === "legacy-summary",
				tokensBefore: event.preparation.tokensBefore,
				keepRecentTokens: event.preparation.settings.keepRecentTokens,
			});
		}
		stage = "classifier";
		log = await CanonicalRecordLog.create(callbacks.onRecord);
		const classification = await classifyAllChunks(
			ctx,
			model,
			currentModel,
			chunks,
			previous?.atoms ?? [],
			event.preparation.settings.reserveTokens,
			event.customInstructions,
			signal,
			log,
			runIndex,
		);
		stage = "reconcile";
		await log.drain();
		const records = reconcileRecords(log.records, chunks.length);
		stage = "merge";
		const atoms = mergeMemory(previous?.atoms ?? [], records, chunks);
		const classifierModel = classification.model;
		const activeSummarizerModel = summarizerFollowsClassifier
			? classifierModel
			: summarizerModel;
		const summarizerRole =
			currentModel &&
			!sameModel(activeSummarizerModel, summarizerModel) &&
			sameModel(activeSummarizerModel, currentModel)
				? "current_fallback"
				: "selected";
		let outputBudget = maximumOutputTokens(
			activeSummarizerModel,
			event.preparation.settings.reserveTokens,
		);
		let allocated = allocateAtoms(atoms, outputBudget);
		logMemoryPlan(runIndex, atoms, allocated, outputBudget, summarizerRole);
		callbacks.onMerge(atoms, allocated, activeSummarizerModel);
		stage = "summarizer";
		let prose: string;
		let summarizerUsage: Usage = combineUsage();
		let degraded: DecayRunResult["degraded"];
		let summarized: { prose: string; usage: Usage } | undefined;
		try {
			summarized = await summarize(
				ctx,
				activeSummarizerModel,
				allocated,
				outputBudget,
				event.customInstructions,
				signal,
				callbacks,
				runIndex,
				summarizerRole,
			);
		} catch (error) {
			if (signal.aborted) throw error;
			if (currentModel && !sameModel(activeSummarizerModel, currentModel)) {
				outputBudget = maximumOutputTokens(
					currentModel,
					event.preparation.settings.reserveTokens,
				);
				allocated = allocateAtoms(atoms, outputBudget);
				logMemoryPlan(
					runIndex,
					atoms,
					allocated,
					outputBudget,
					"current_fallback",
				);
				callbacks.onMerge(atoms, allocated, currentModel);
				callbacks.onSummary?.(0);
				try {
					summarized = await summarize(
						ctx,
						currentModel,
						allocated,
						outputBudget,
						event.customInstructions,
						signal,
						callbacks,
						runIndex,
						"current_fallback",
					);
				} catch (fallbackError) {
					if (signal.aborted) throw fallbackError;
					degraded = {
						stage: "summarizer",
						error: asError(fallbackError),
					};
				}
			} else {
				degraded = { stage: "summarizer", error: asError(error) };
			}
		}
		if (summarized) {
			prose = summarized.prose;
			summarizerUsage = summarized.usage;
		} else {
			prose = renderFallbackSummary(allocated);
			callbacks.onSummary?.(Math.ceil(prose.length / 4));
		}
		const files = mergeFileLists(previous, event.preparation.fileOps);
		const summary =
			prose + formatFileAppendix(files.readFiles, files.modifiedFiles);
		const elapsedMs = Date.now() - startedAt;
		const estimatedKeptTokens = estimateKeptTokens(event);
		result = {
			summary,
			usage: combineUsage(classification.usage, summarizerUsage),
			...(degraded ? { degraded } : {}),
			estimatedTokensAfter:
				estimatedKeptTokens === undefined
					? undefined
					: Math.ceil(summary.length / 4) + estimatedKeptTokens,
			details: {
				decay: {
					version: DECAY_DETAILS_VERSION,
					model: `${classifierModel.provider}/${classifierModel.id}`,
					elapsedMs,
					chunkCount: chunks.length,
					keepRecentTokens: event.preparation.settings.keepRecentTokens,
					...(estimatedKeptTokens === undefined ? {} : { estimatedKeptTokens }),
					...(degraded ? { degraded: degraded.stage } : {}),
					atoms,
					readFiles: files.readFiles,
					modifiedFiles: files.modifiedFiles,
				},
			},
		};
	} catch (error) {
		failure = decayRunError(stage, error);
	}
	if (log) {
		try {
			await log.close();
		} catch (error) {
			failure ??= decayRunError("cleanup", error);
		}
	}
	if (failure) throw failure;
	if (!result)
		throw new DecayRunError("prepare", new Error("Decay produced no result"));
	return result;
}

function reportedUsage(usage: Usage | undefined) {
	return usage
		? {
				reportedInputTokens: usage.input,
				reportedOutputTokens: usage.output,
				reportedCacheReadTokens: usage.cacheRead,
				reportedCost: usage.cost?.total,
			}
		: {};
}

function logMemoryPlan(
	runIndex: number,
	atoms: readonly MemoryAtom[],
	selected: ReturnType<typeof allocateAtoms>,
	outputBudget: number,
	modelRole: "selected" | "current_fallback",
): void {
	if (!dbg) return;
	const active = atoms.filter((atom) => atom.status === "active");
	dbg("compaction.memory", {
		runIndex,
		modelRole,
		memoryAtoms: atoms.length,
		activeAtoms: active.length,
		selectedAtoms: selected.length,
		protectedAtoms: selected.filter((atom) => isProtectedKind(atom.kind))
			.length,
		allocatedTokens: selected.reduce((sum, atom) => sum + atom.allocation, 0),
		outputBudget,
		atomGoals: active.filter((atom) => atom.kind === "goal").length,
		atomConstraints: active.filter((atom) => atom.kind === "constraint").length,
		atomDecisions: active.filter((atom) => atom.kind === "decision").length,
		atomStates: active.filter((atom) => atom.kind === "state").length,
		atomBlockers: active.filter((atom) => atom.kind === "blocker").length,
		atomExacts: active.filter((atom) => atom.kind === "exact").length,
		atomContexts: active.filter((atom) => atom.kind === "context").length,
	});
}

async function summarize(
	ctx: ExtensionContext,
	model: AnyModel,
	atoms: ReturnType<typeof allocateAtoms>,
	outputBudget: number,
	manualInstructions: string | undefined,
	signal: AbortSignal,
	callbacks: RunCallbacks,
	runIndex: number,
	modelRole: "selected" | "current_fallback",
): Promise<{ prose: string; usage: Usage }> {
	const finish = span?.("summarizer.request", {
		runIndex,
		modelRole,
		selectedAtoms: atoms.length,
		outputBudget,
	});
	let usage: Usage | undefined;
	let stopReason:
		| "stop"
		| "toolUse"
		| "length"
		| "error"
		| "aborted"
		| "deferred"
		| "pending"
		| undefined;
	let requestResult:
		| "accepted"
		| "unexpected_tool"
		| "empty_summary"
		| "output_limit"
		| "provider_error"
		| "cancelled"
		| "other" = "other";
	try {
		const stream = ctx.modelRegistry.streamSimple(
			model,
			normalizeContext(buildSummarizerContext(atoms, manualInstructions)),
			{
				maxTokens: outputBudget,
				signal,
				cacheRetention: "none",
				sessionId: uuidv7(),
				reasoning: model.reasoning ? "low" : undefined,
			},
		);
		for await (const item of stream) {
			if (item.type === "text_delta") {
				callbacks.onSummary?.(
					Math.ceil(contentText(item.partial.content).length / 4),
				);
			}
		}
		const response = await stream.result();
		usage = response.usage;
		stopReason = response.stopReason;
		if (stopReason !== "stop") {
			requestResult =
				stopReason === "length" ? "output_limit" : "provider_error";
			throw new Error(
				`Decay summarizer stopped with ${stopReason}: ${response.errorMessage ?? "no detail"}`,
			);
		}
		if (response.content.some((content) => content.type === "toolCall")) {
			requestResult = "unexpected_tool";
			throw new Error("Decay summarizer returned a tool call");
		}
		const prose = contentText(response.content).trim();
		if (!prose) {
			requestResult = "empty_summary";
			throw new Error("Decay summarizer returned an empty summary");
		}
		finish?.("finish", {
			runIndex,
			requestResult: "accepted",
			stopReason,
			summaryTokens: Math.ceil(prose.length / 4),
			...reportedUsage(usage),
		});
		return { prose, usage };
	} catch (error) {
		finish?.("error", {
			runIndex,
			requestResult: signal.aborted ? "cancelled" : requestResult,
			stopReason,
			...reportedUsage(usage),
		});
		throw error;
	}
}

async function classifyAllChunks(
	ctx: ExtensionContext,
	model: AnyModel,
	currentModel: AnyModel | undefined,
	chunks: readonly SourceChunk[],
	priorAtoms: readonly MemoryAtom[],
	reserveTokens: number,
	manualInstructions: string | undefined,
	signal: AbortSignal,
	log: CanonicalRecordLog,
	runIndex: number,
): Promise<ClassificationResult> {
	const usages: Usage[] = [];
	let activeModel = model;
	for (const chunk of chunks) {
		const catalog = mergeMemory(priorAtoms, log.records, chunks);
		let attempt = 0;
		for (let round = 0; ; round++) {
			attempt++;
			let failure: Error | undefined;
			try {
				usages.push(
					await classify(
						ctx,
						activeModel,
						chunk,
						catalog,
						reserveTokens,
						manualInstructions,
						signal,
						log,
						runIndex,
						attempt,
						sameModel(activeModel, model) ? "selected" : "current_fallback",
					),
				);
			} catch (error) {
				if (signal.aborted) throw error;
				failure = asError(error);
			}
			try {
				await log.drain();
			} catch (error) {
				dbg?.("classifier.decision", {
					runIndex,
					chunkIndex: chunk.chunk,
					attempt,
					decision: "give_up",
					requestResult: "record_log_error",
				});
				throw error;
			}
			if (!missingChunks(log.records, chunks.length).includes(chunk.chunk))
				break;
			if (failure && currentModel && !sameModel(activeModel, currentModel)) {
				dbg?.("classifier.decision", {
					runIndex,
					chunkIndex: chunk.chunk,
					attempt,
					decision: "switch_model",
					transient: isTransientFailure(failure),
				});
				activeModel = currentModel;
				round = -1;
				continue;
			}
			if (failure && !isTransientFailure(failure)) {
				dbg?.("classifier.decision", {
					runIndex,
					chunkIndex: chunk.chunk,
					attempt,
					decision: "give_up",
					transient: false,
				});
				throw failure;
			}
			if (round >= CLASSIFIER_RETRIES) {
				dbg?.("classifier.decision", {
					runIndex,
					chunkIndex: chunk.chunk,
					attempt,
					decision: "give_up",
					transient: failure ? true : undefined,
					requestResult: failure ? undefined : "omitted",
				});
				throw (
					failure ?? new Error(`Decay classifier omitted chunk: ${chunk.chunk}`)
				);
			}
			dbg?.("classifier.decision", {
				runIndex,
				chunkIndex: chunk.chunk,
				attempt,
				decision: "retry",
				transient: failure ? true : undefined,
				requestResult: failure ? undefined : "omitted",
			});
		}
	}
	return { usage: combineUsage(...usages), model: activeModel };
}

async function classify(
	ctx: ExtensionContext,
	model: AnyModel,
	chunk: SourceChunk,
	priorAtoms: readonly MemoryAtom[],
	reserveTokens: number,
	manualInstructions: string | undefined,
	signal: AbortSignal,
	log: CanonicalRecordLog,
	runIndex = 0,
	attempt = 1,
	modelRole: "selected" | "current_fallback" = "selected",
): Promise<Usage> {
	const outputBudget = maximumOutputTokens(model, reserveTokens);
	let catalogAtoms = 0;
	let catalogChars = 0;
	const classifierContext = buildClassifierContext(
		[chunk],
		priorAtoms,
		manualInstructions,
		span
			? (includedAtoms, characters) => {
					catalogAtoms = includedAtoms;
					catalogChars = characters;
				}
			: undefined,
	);
	const finish = span?.("classifier.request", {
		runIndex,
		chunkIndex: chunk.chunk,
		chunkCount: chunk.of,
		attempt,
		modelRole,
		sourceKind: chunk.source,
		sourceTokens: chunk.estimatedTokens,
		catalogAvailableAtoms: priorAtoms.length,
		catalogAtoms,
		catalogChars,
		outputBudget,
	});
	let usage: Usage | undefined;
	let stopReason:
		| "stop"
		| "toolUse"
		| "length"
		| "error"
		| "aborted"
		| "deferred"
		| "pending"
		| undefined;
	let requestResult:
		| "accepted"
		| "omitted"
		| "unknown_tool"
		| "invalid_record"
		| "wrong_shard"
		| "conflicting_records"
		| "output_limit"
		| "incomplete_stream"
		| "provider_error"
		| "record_log_error"
		| "cancelled"
		| "other" = "other";
	let toolCalls = 0;
	const records: RecordChunk[] = [];
	try {
		const stream = ctx.modelRegistry.streamSimple(
			model,
			normalizeContext(classifierContext),
			{
				signal,
				cacheRetention: "none",
				sessionId: uuidv7(),
				reasoning:
					model.provider === "openai-codex" && model.id === "gpt-6-luna"
						? undefined
						: model.reasoning
							? "low"
							: undefined,
				maxTokens: outputBudget,
			},
		);
		for await (const item of stream) {
			if (item.type === "toolcall_end") {
				toolCalls++;
				try {
					records.push(parseToolCall(item.toolCall, chunk));
				} catch (error) {
					requestResult =
						item.toolCall.name !== CLASSIFIER_TOOL_NAME
							? "unknown_tool"
							: !validateRecordChunk(item.toolCall.arguments)
								? "invalid_record"
								: "wrong_shard";
					throw error;
				}
			} else if (item.type === "done") {
				usage = item.message.usage;
				stopReason = item.message.stopReason;
				if (stopReason !== "stop" && stopReason !== "toolUse") {
					requestResult =
						stopReason === "length" ? "output_limit" : "provider_error";
					throw new Error(
						`Decay classifier stream ended without completion (${stopReason})`,
					);
				}
			} else if (item.type === "error") {
				usage = item.error.usage;
				stopReason = item.error.stopReason;
				requestResult = "provider_error";
				throw new Error(
					item.error.errorMessage ||
						`Decay classifier stopped with ${item.reason}`,
				);
			}
		}
		if (!usage) {
			requestResult = "incomplete_stream";
			throw new Error("Decay classifier stream ended without completion");
		}
		let record: RecordChunk | undefined;
		try {
			record = indexRecords(records, chunk.of).get(chunk.chunk);
		} catch (error) {
			requestResult = "conflicting_records";
			throw error;
		}
		if (record) {
			try {
				await log.append(record);
			} catch (error) {
				requestResult = "record_log_error";
				throw error;
			}
		}
		requestResult = record ? "accepted" : "omitted";
		finish?.("finish", {
			runIndex,
			chunkIndex: chunk.chunk,
			attempt,
			requestResult,
			stopReason,
			toolCalls,
			validRecords: records.length,
			committedRecords: record ? 1 : 0,
			receivedAtoms: records.reduce((sum, item) => sum + item.atoms.length, 0),
			recordAtoms: record?.atoms.length ?? 0,
			...reportedUsage(usage),
		});
		return usage;
	} catch (error) {
		finish?.("error", {
			runIndex,
			chunkIndex: chunk.chunk,
			attempt,
			requestResult: signal.aborted ? "cancelled" : requestResult,
			stopReason,
			toolCalls,
			validRecords: records.length,
			committedRecords: 0,
			receivedAtoms: records.reduce((sum, item) => sum + item.atoms.length, 0),
			...reportedUsage(usage),
		});
		throw error;
	}
}

function parseToolCall(toolCall: ToolCall, chunk: SourceChunk): RecordChunk {
	if (toolCall.name !== CLASSIFIER_TOOL_NAME)
		throw new Error(`Decay classifier called unknown tool: ${toolCall.name}`);
	const record = validateRecordChunk(toolCall.arguments);
	if (!record)
		throw new Error("Decay classifier returned an invalid chunk record");
	if (record.chunk !== chunk.chunk || record.of !== chunk.of)
		throw new Error("Decay classifier returned a record for another chunk");
	return record;
}

function maximumOutputTokens(model: AnyModel, reserveTokens: number): number {
	const reserveBudget = Math.max(1, Math.floor(reserveTokens * 0.8));
	return model.maxTokens > 0
		? Math.min(reserveBudget, model.maxTokens)
		: reserveBudget;
}

function sameModel(left: AnyModel, right: AnyModel): boolean {
	return left.provider === right.provider && left.id === right.id;
}

function modelLabel(reference: string): string {
	const separator = reference.indexOf("/");
	return separator >= 0 ? reference.slice(separator + 1) : reference;
}

function decayRunError(
	stage: DecayFailureStage,
	value: unknown,
): DecayRunError {
	return value instanceof DecayRunError
		? value
		: new DecayRunError(stage, asError(value));
}

function failureStage(error: Error): DecayFailureStage {
	return error instanceof DecayRunError ? error.stage : "prepare";
}

function persistDiagnostic(
	pi: ExtensionAPI,
	outcome: "fallback" | "salvage",
	stage: DecayFailureStage,
	model: string,
	error: Error,
): void {
	try {
		pi.appendEntry(DECAY_DIAGNOSTIC_ENTRY, {
			version: 1,
			outcome,
			stage,
			model: sanitizeDiagnostic(model).slice(0, 160),
			reason: boundedError(error),
			timestamp: Date.now(),
		});
	} catch {
		// Diagnostic persistence must never prevent Pi's default compaction.
	}
}

function boundedError(error: Error): string {
	return sanitizeDiagnostic(error.message).slice(0, 240);
}

function sanitizeDiagnostic(value: string): string {
	return value
		.replace(
			/([?&](?:api_?key|access_token|token|secret)=)[^&\s]+/gi,
			"$1[redacted]",
		)
		.replace(/\b(Bearer)\s+\S+/gi, "$1 [redacted]")
		.replace(
			/\b(api[-_ ]?key|authorization|token|secret|password)\b(\s*[:=]\s*)[^\s,;]+/gi,
			"$1$2[redacted]",
		)
		.replace(/\s+/g, " ")
		.trim();
}

function warn(ctx: ExtensionContext, message: string): void {
	if (ctx.hasUI) ctx.ui.notify(message, "warning");
}

function asError(value: unknown): Error {
	return value instanceof Error ? value : new Error(String(value));
}

export const __decayTest = {
	resolveModel,
	previewChunkCount,
	runDecay,
	classify,
	classifyAllChunks,
	summarize,
	parseToolCall,
	maximumOutputTokens,
	boundedError,
};