repositories / pi-ext
pi-ext
bugabingas pi extensions
owned by admin
extensions/decay/implementation.ts
Rawimport {
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,
};