Something went wrong. Try again.
A local-first event pipeline for independent agents, built on Jazz.
Something went wrong. Try again.
35 kB · 775 lines
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776import type { AssistantMessage, ImageContent } from "@earendil-works/pi-ai";import { createHash, randomUUID } from "node:crypto";import path from "node:path";import { canonicalJson, type JsonObject } from "../core/json.js";import { privateCheckpointForDeclaration } from "./declarations.js";import type { InferenceUsage } from "../store/types.js";import { CONCEPTUALIZATION_OUTPUT_JSON_SCHEMA, createOutputContractRegistry, outputContractForDeclaration, outputContractIdentityJson, OutputContractValidationError, type ConversationCompactionOutput, type ObservationOutput, type OutputContractIdentity, type OutputContractRegistry, type SemanticOutput,} from "./output-contracts.js";import { createBuiltinProviderProfileResolver, type ProviderProfileResolver,} from "./provider-profiles.js";import { launchSandboxedModel, SandboxExecutionError } from "./sandbox/launcher.js";import { startProviderBroker } from "./sandbox/provider-broker.js";import { createRunTools, type AgentToolName, type RunToolOptions } from "./tools.js";import { AgentRunFailure, type AgentOutput, type AgentRunInput, type AgentRunner, type RunnerTrace, type ThoughtAgentDeclaration } from "./types.js";import { resolveImageArtifact } from "../connectors/telegram-images.js";import type { AgentContextPacket, ImageArtifactReference } from "./context.js";import { capturedProposalsSchema, proposalCapabilitiesSchema, validateCapturedProposalAgainstCapabilities, type CapturedProposal, type ProposalCapabilities,} from "./proposals.js";import { SANDBOX_PROTOCOL_VERSION } from "./sandbox/protocol.js";import { isTelegramHelpCommand, isTelegramFocusHelpCommand, telegramFocusDescription, STREAM_TELEGRAM_FOCUS_HELP_TEXT, STREAM_TELEGRAM_FOCUS_HELP_VERSION, STREAM_TELEGRAM_HELP_TEXT, STREAM_TELEGRAM_HELP_VERSION,} from "./telegram-help.js";
const MAX_FINAL_JSON_CHARS = 64_000;const MAX_CONVERSATION_TEXT_CHARS = 4_096;const MAX_COMPACTION_TEXT_CHARS = 24_000;
export interface PiAgentRunnerOptions extends Pick<RunToolOptions, "artifactRoot" | "fetchImpl" | "resolveHostname"> { providerProfiles?: ProviderProfileResolver; providerFetchImpl?: typeof fetch; workerBundlePath?: string; outputContracts?: OutputContractRegistry;}
export class PiAgentRunner implements AgentRunner { readonly mode = "pi" as const; private readonly providerProfiles: ProviderProfileResolver; private readonly outputContracts: OutputContractRegistry;
constructor(private readonly options: PiAgentRunnerOptions = {}) { this.providerProfiles = options.providerProfiles ?? createBuiltinProviderProfileResolver(); this.outputContracts = options.outputContracts ?? createOutputContractRegistry(); }
async run(input: AgentRunInput, onTrace: (trace: RunnerTrace) => Promise<void>): Promise<AgentOutput> { const declaration = input.declaration; const outputContractIdentity = outputContractForDeclaration(declaration); const outputContract = this.outputContracts.resolve(outputContractIdentity); if ( declaration.outputMode === "conversation-text" && declaration.contextStrategy === "telegram-conversation" && isTelegramHelpCommand(input.event) ) { await onTrace({ kind: "telegram.command.completed", data: { command: "help", version: STREAM_TELEGRAM_HELP_VERSION, providerRequests: 0 }, }); const output = this.outputContracts.validate(outputContractIdentity, { summary: STREAM_TELEGRAM_HELP_TEXT, tags: ["conversation", "help"], importance: "normal", confidence: 1, }); return { ...output, model: { provider: "trusted-parent", id: STREAM_TELEGRAM_HELP_VERSION }, }; } if ( declaration.outputMode === "conversation-text" && declaration.contextStrategy === "telegram-conversation" && isTelegramFocusHelpCommand(input.event) ) { await onTrace({ kind: "telegram.command.completed", data: { command: "focus", version: STREAM_TELEGRAM_FOCUS_HELP_VERSION, providerRequests: 0 }, }); const output = this.outputContracts.validate(outputContractIdentity, { summary: STREAM_TELEGRAM_FOCUS_HELP_TEXT, tags: ["conversation", "focus-help"], importance: "normal", confidence: 1, }); return { ...output, model: { provider: "trusted-parent", id: STREAM_TELEGRAM_FOCUS_HELP_VERSION }, }; } if (!declaration.model) throw new Error(`Pi agent ${declaration.id} has no model`); if (!declaration.providerProfile) throw new Error(`Pi agent ${declaration.id} has no trusted provider profile`); const providerModel = declaration.modelAdapter ? privateCheckpointForDeclaration(declaration) : declaration.model; const profile = this.providerProfiles.resolve(declaration.providerProfile, providerModel); const providerIdentity: JsonObject = { provider: profile.provider, providerProfile: profile.id, model: declaration.model, ...(declaration.modelAdapter ? { modelAdapter: declaration.modelAdapter as unknown as JsonObject } : {}), }; const acceptsImages = profile.imageInputModels.has(declaration.model); let observedRevision: string | undefined; let observedUsage: InferenceUsage | undefined; const toolSet = createRunTools({ event: input.event, names: declaration.tools as AgentToolName[], onTrace, artifactRoot: this.options.artifactRoot, fetchImpl: this.options.fetchImpl, resolveHostname: this.options.resolveHostname, }); const prefetched = toolSet.tools.length > 0 ? await prefetchReadOnlyEvidence(input.event, toolSet.tools, onTrace, declaration.timeoutMs, acceptsImages) : { text: "", images: [] as ImageContent[] }; // Resolve current-turn image artifact references from the context packet. // The trusted Pi parent validates and resolves content-addressed artifacts // beneath the artifact root before injecting base64 ImageContent into the sandbox. const contextImages = acceptsImages && this.options.artifactRoot && input.context.imageArtifacts ? await resolveContextImages(input.context.imageArtifacts, this.options.artifactRoot, onTrace) : []; const allImages = [...contextImages, ...prefetched.images]; const { systemPrompt, prompt, proposalCapabilities } = composePiModelInput(declaration, input.context, { readOnlyToolCount: toolSet.tools.length, readOnlyEvidenceText: prefetched.text, resolvedContextImageCount: contextImages.length, outputContractPrompt: outputContract.prompt, }); await onTrace({ kind: "system_prompt", data: stringMetadata(systemPrompt) }); await onTrace({ kind: "prompt", data: stringMetadata(prompt) });
const broker = await startProviderBroker({ profile, model: providerModel, traceModel: declaration.model, runId: input.runId, timeoutMs: declaration.timeoutMs, maxOutputTokens: declaration.maxOutputTokens, ...(this.options.providerFetchImpl ? { fetchImpl: this.options.providerFetchImpl } : {}), onTrace: async (kind, data) => onTrace({ kind, data }), }).catch((error) => { throw new AgentRunFailure("Disposable model provider broker could not start", { cause: error, advanceProgress: false, diagnostic: { code: "broker-start-failed", stage: "sandbox-setup", ...providerIdentity }, }); });
try { const result = await launchSandboxedModel({ version: SANDBOX_PROTOCOL_VERSION, runId: input.runId, systemPrompt, prompt, ...(input.context.messages ? { messages: input.context.messages } : {}), images: allImages, ...(proposalCapabilities ? { proposals: proposalCapabilities } : {}), model: { provider: profile.provider, id: providerModel, reasoning: profile.provider === "tinker", thinkingLevel: declaration.thinkingLevel ?? "off", acceptsImages, jsonObjectResponseFormat: (declaration.outputMode ?? "strict-json") === "strict-json" && profile.jsonObjectResponseFormat, ...((declaration.outputMode ?? "strict-json") === "strict-json" && profile.jsonSchemaResponseFormat ? { jsonSchemaResponseFormat: CONCEPTUALIZATION_OUTPUT_JSON_SCHEMA } : {}), contextWindow: 131_072, maxTokens: declaration.maxOutputTokens, }, broker: { socketPath: "/broker/provider.sock", capability: broker.capability, origin: "http://provider.invalid", route: "/v1/chat/completions", }, }, { workerBundlePath: this.options.workerBundlePath ?? path.resolve(process.cwd(), "dist/sandbox/worker.cjs"), brokerDirectory: broker.socketDirectory, timeoutMs: declaration.timeoutMs + 2_000, }); observedRevision = result.observedRevision; await emitSandboxTraces(result.traces, onTrace); const finalMessage = validateFinalAssistant(result.finalMessage, outputContractIdentity); const proposals = validateProposalCompletion(finalMessage, result.proposals, proposalCapabilities, outputContractIdentity); if (telegramFocusDescription(input.event) !== undefined && proposals.filter((proposal) => proposal.kind === "focus-declaration").length !== 1) { throw invalidConversationText("focus-proposal-required", finalOutputDiagnostic(finalMessage), outputContractIdentity); } observedUsage = inferenceUsageFromAssistant(finalMessage); const parsed = declaration.outputMode === "conversation-text" ? parseConversationText(finalMessage, proposals, outputContractIdentity, this.outputContracts) : declaration.outputMode === "compaction-text" ? parseCompactionText(finalMessage, proposals, outputContractIdentity, this.outputContracts) : parseFinalOutput(finalMessage, proposals, outputContractIdentity, this.outputContracts); return { ...parsed, model: { provider: profile.provider, id: declaration.model, ...(declaration.modelAdapter ? { revision: `sha256:${declaration.modelAdapter.binding.checkpointReferenceSha256}` } : observedRevision ? { revision: observedRevision } : {}), }, ...(toolSet.outcomes.length > 0 ? { enrichments: toolSet.outcomes } : {}), ...(observedUsage ? { usage: observedUsage } : {}), ...(proposals.length > 0 ? { proposals } : {}), }; } catch (error) { if (error instanceof AgentRunFailure) { const failureUsage = error.usage ?? observedUsage; throw new AgentRunFailure(error.message, { cause: error, advanceProgress: error.advanceProgress, ...(failureUsage ? { usage: failureUsage } : {}), diagnostic: { ...(error.diagnostic ?? {}), ...providerIdentity, ...(declaration.modelAdapter ? { checkpointRevision: `sha256:${declaration.modelAdapter.binding.checkpointReferenceSha256}` } : observedRevision ? { checkpointRevision: observedRevision } : {}), }, }); } if (error instanceof SandboxExecutionError) { await emitSandboxTraces(error.traces, onTrace); const providerTimeout = error.code === "sandbox-worker-provider-timeout"; const providerClientError = error.code === "sandbox-worker-provider-client-error"; throw new AgentRunFailure(providerTimeout ? "Pi provider run timed out" : "Disposable model sandbox failed", { cause: error, advanceProgress: providerClientError, diagnostic: { code: providerTimeout ? "timeout" : providerClientError ? "provider-client-error" : error.code, stage: providerTimeout || providerClientError ? "provider" : "sandbox-execution", ...(error.workerReason ? { workerReason: error.workerReason } : {}), ...(error.proposalFailureCode ? { proposalFailureCode: error.proposalFailureCode } : {}), ...providerIdentity, }, }); } throw new AgentRunFailure("Pi provider run failed", { cause: error, advanceProgress: false, diagnostic: { code: "provider-run-failed", stage: "sandbox-execution", ...providerIdentity }, }); } finally { await broker.close(); } }}
export function composePiModelInput( declaration: ThoughtAgentDeclaration, context: AgentContextPacket, options: { readOnlyToolCount: number; readOnlyEvidenceText: string; resolvedContextImageCount: number; outputContractPrompt: string; },): { systemPrompt: string; prompt: string; proposalCapabilities?: ProposalCapabilities } { const proposalCapabilities = resolveProposalCapabilities(declaration, context.manifest); const visualInstructions = declaration.outputMode === "conversation-text" && declaration.contextStrategy === "telegram-conversation" ? options.resolvedContextImageCount > 0 ? `The current user turn includes ${options.resolvedContextImageCount} trusted image input${options.resolvedContextImageCount === 1 ? "" : "s"}. You may describe only what those current image parts support.` : "The current user turn includes no image input. Do not claim to see a current photo or reconstruct one from earlier text, placeholders, or assistant descriptions." : ""; const toolInstructions = options.readOnlyToolCount > 0 ? "The trusted parent has already acquired the configured read-only evidence below. Tool failures are evidence: report uncertainty rather than inventing missing context." : "No read-only evidence tool was used for this turn."; const proposalInstructions = proposalCapabilities ? renderProposalInstructions(proposalCapabilities) : ""; const outputPrompt = declaration.outputMode === "conversation-text" ? proposalCapabilities ? "Reply with concise visible text, use one relevant proposal function, or do both. Do not emit JSON or narrate function machinery." : "Reply with concise visible text." : declaration.outputMode === "compaction-text" ? "Return only the plain-text continuity summary. Do not emit JSON, Markdown fences, analysis, or a reply to the conversation." : `${options.outputContractPrompt} Do not emit analysis, a plan, Markdown, or <think> tags.`; const subscribedContext = context.systemText ? `\n\n${context.systemText}` : ""; const systemPrompt = [ declaration.systemPrompt, subscribedContext, toolInstructions, proposalInstructions, visualInstructions, outputPrompt, ].filter(Boolean).join("\n\n"); const sourceContext = options.readOnlyEvidenceText ? `${context.text}\n\n## Pre-fetched read-only evidence\n${options.readOnlyEvidenceText}` : context.text; // A Telegram conversation context already exposes the current inbound text // as `sourceContext`; `agent.prompt()` appends it as the final user message. // Appending output instructions here would create a synthetic user turn // after Cameron's actual message and make harness prose look conversational. const prompt = declaration.outputMode === "conversation-text" ? sourceContext : declaration.outputMode === "compaction-text" ? `${sourceContext}\n\n${outputPrompt}` : `${sourceContext}\n\n## Required final answer\n${outputPrompt}`; return { systemPrompt, prompt, ...(proposalCapabilities ? { proposalCapabilities } : {}) };}
async function emitSandboxTraces( traces: Array<{ kind: string; data: unknown }>, onTrace: (trace: RunnerTrace) => Promise<void>,): Promise<void> { const messageUpdateCount = traces.filter((trace) => trace.kind === "pi.message_update").length; for (const trace of traces.filter((candidate) => candidate.kind !== "pi.message_update")) { await onTrace({ kind: trace.kind, data: traceMetadata(trace.data) }); } if (messageUpdateCount > 0) { await onTrace({ kind: "pi.message_updates_coalesced", data: { count: messageUpdateCount } }); }}
async function prefetchReadOnlyEvidence( event: AgentRunInput["event"], tools: ReturnType<typeof createRunTools>["tools"], onTrace: (trace: RunnerTrace) => Promise<void>, timeoutMs: number, acceptsImages: boolean,): Promise<{ text: string; images: ImageContent[] }> { const signal = AbortSignal.timeout(Math.max(100, timeoutMs)); const markdownTool = tools.find((tool) => tool.name === "fetch_atproto_markdown"); const imageTool = tools.find((tool) => tool.name === "download_image"); const evidence: string[] = []; const images: ImageContent[] = []; let imageDataChars = 0; const imageUrls = new Set<string>(); const collection = typeof event.payload.collection === "string" ? event.payload.collection : ""; const targets: Array<"event" | "subject"> = ["event"]; if (collection === "app.bsky.feed.like" || collection === "app.bsky.feed.repost") targets.push("subject");
for (const target of targets) { if (!markdownTool) break; try { await onTrace({ kind: "pi.tool_prefetch_start", data: { toolName: markdownTool.name, argumentsRedacted: true }, }); const result = await markdownTool.execute(randomUUID(), { target }, signal); evidence.push(...result.content.filter((part) => part.type === "text").map((part) => part.text)); const details = result.details as { imageUrls?: unknown }; if (Array.isArray(details.imageUrls)) { for (const url of details.imageUrls) if (typeof url === "string") imageUrls.add(url); } await onTrace({ kind: "pi.tool_prefetch_end", data: { toolName: markdownTool.name, target, status: "succeeded" } }); } catch { evidence.push(`[${markdownTool.name} ${target} failed]`); await onTrace({ kind: "pi.tool_prefetch_end", data: { toolName: markdownTool.name, target, status: "failed", errorCode: "tool-execution-failed" }, }); } }
if (!acceptsImages && imageUrls.size > 0) { evidence.push(`[${imageUrls.size} referenced image${imageUrls.size === 1 ? " was" : "s were"} omitted because the trusted model profile is text-only]`); await onTrace({ kind: "pi.image_prefetch_skipped", data: { reason: "model-text-only", imageCount: imageUrls.size }, }); }
for (const url of acceptsImages ? [...imageUrls].slice(0, 8) : []) { if (!imageTool || imageDataChars >= 8 * 1024 * 1024) break; try { await onTrace({ kind: "pi.tool_prefetch_start", data: { toolName: imageTool.name, argumentsRedacted: true }, }); const result = await imageTool.execute(randomUUID(), { url }, signal); evidence.push(...result.content.filter((part) => part.type === "text").map((part) => part.text)); for (const image of result.content.filter((part): part is ImageContent => part.type === "image")) { if (imageDataChars + image.data.length > 8 * 1024 * 1024) { imageDataChars = 8 * 1024 * 1024; break; } imageDataChars += image.data.length; images.push(image); } await onTrace({ kind: "pi.tool_prefetch_end", data: { toolName: imageTool.name, status: "succeeded" } }); } catch { evidence.push(`[${imageTool.name} failed]`); await onTrace({ kind: "pi.tool_prefetch_end", data: { toolName: imageTool.name, status: "failed", errorCode: "tool-execution-failed" }, }); } }
return { text: evidence.join("\n\n"), images };}
function resolveProposalCapabilities( declaration: ThoughtAgentDeclaration, manifest: JsonObject,): ProposalCapabilities | undefined { const declared = declaration.proposals ?? []; if (declared.length === 0) { if (manifest.proposalCapabilities !== undefined) throw new Error("Context snapshot contains undeclared proposal capability"); return undefined; } const capabilities = proposalCapabilitiesSchema.parse(manifest.proposalCapabilities); if (capabilities.enabled.some((capability) => !declared.includes(capability))) { throw new Error("Context snapshot proposal capability exceeds the declaration"); } const snapshot = asRecord(manifest.contextSnapshot); if (!snapshot || typeof snapshot.id !== "string") throw new Error("Proposal capability requires an immutable context snapshot"); return capabilities;}
function renderProposalInstructions(capabilities: ProposalCapabilities): string { const lines = [ "Proposal functions create private inert suggestions for review. They do not apply, activate, train, or publish the proposed effect.", ]; if (capabilities.enabled.includes("memory-change")) { lines.push("Use request_memory_change only for a concrete continuity fact worth preserving; set evidence_event_ids to []."); } if (capabilities.enabled.includes("self-correction")) { lines.push("Use submit_correction with target_output `latest` for your last delivered reply; set evidence_event_ids to []."); } if (capabilities.enabled.includes("focus-declaration")) { lines.push("This turn is a validated `/focus` command. Call propose_focus exactly once. Convert its description into one complete conservative declaration: lowercase hyphenated id, version 1 unless revision is explicit, replay now, bounded budgets, and no external actions unless explicitly requested."); } return lines.join("\n");}
function validateFinalAssistant(value: unknown, identity: OutputContractIdentity): AssistantMessage { const message = asRecord(value); const content = message?.content; const malformed = !message || message.role !== "assistant" || !Array.isArray(content) || (message.model !== undefined && typeof message.model !== "string") || (message.stopReason !== undefined && typeof message.stopReason !== "string") || content.some((value) => { const part = asRecord(value); if (!part || typeof part.type !== "string") return true; if (part.type === "text") return typeof part.text !== "string"; if (part.type === "thinking") return typeof part.thinking !== "string"; if (part.type === "toolCall") return typeof part.id !== "string" || typeof part.name !== "string" || !part.arguments || typeof part.arguments !== "object" || Array.isArray(part.arguments); return true; }); if (malformed) { throw new AgentRunFailure("Pi final output rejected: invalid assistant message", { diagnostic: { code: "invalid-final-output", stage: "final-output-validation", reason: "malformed-assistant-message", assistantMessages: 0, outputContract: outputContractIdentityJson(identity), }, }); } return value as AssistantMessage;}
function validateProposalCompletion( message: AssistantMessage, rawProposals: unknown, capabilities: ProposalCapabilities | undefined, identity: OutputContractIdentity,): CapturedProposal[] { const diagnostic = finalOutputDiagnostic(message); let proposals: CapturedProposal[]; try { proposals = capturedProposalsSchema.parse(rawProposals); } catch { throw invalidConversationText("invalid-proposal-result", diagnostic, identity); } const toolCalls = message.content.filter((part) => part.type === "toolCall"); if (!capabilities) { if (proposals.length > 0 || toolCalls.length > 0) throw invalidConversationText("undeclared-proposal-call", diagnostic, identity); return []; } if (toolCalls.length !== proposals.length) throw invalidConversationText("proposal-call-result-mismatch", diagnostic, identity); const byId = new Map(toolCalls.map((call) => [call.id, call])); for (const proposal of proposals) { const call = byId.get(proposal.toolCallId); const expectedName = proposal.kind === "memory-change" ? "request_memory_change" : proposal.kind === "self-correction" ? "submit_correction" : "propose_focus"; if (!call || call.name !== expectedName) throw invalidConversationText("proposal-call-result-mismatch", diagnostic, identity); if (canonicalJson(call.arguments as JsonObject) !== canonicalJson(proposal.arguments as unknown as JsonObject)) { throw invalidConversationText("proposal-arguments-mismatch", diagnostic, identity); } try { validateCapturedProposalAgainstCapabilities(proposal, capabilities); } catch { throw invalidConversationText("proposal-capability-mismatch", diagnostic, identity); } } return proposals;}
function inferenceUsageFromAssistant(message: AssistantMessage): InferenceUsage | undefined { const usage = asRecord(message.usage); if (!usage) return undefined; const input = nonnegativeSafeInteger(usage.input); const output = nonnegativeSafeInteger(usage.output); const cacheRead = nonnegativeSafeInteger(usage.cacheRead) ?? 0; const cacheWrite = nonnegativeSafeInteger(usage.cacheWrite) ?? 0; if (input === undefined || output === undefined) return undefined; const inputTokens = input + cacheRead + cacheWrite; if (!Number.isSafeInteger(inputTokens) || inputTokens + output === 0) return undefined; return { inputTokens, outputTokens: output };}
function nonnegativeSafeInteger(value: unknown): number | undefined { return Number.isSafeInteger(value) && Number(value) >= 0 ? Number(value) : undefined;}
function traceMetadata(value: unknown): JsonObject { const serialized = JSON.stringify(value); return { chars: serialized.length, sha256: sha256Text(serialized), redacted: true, };}
function parseFinalOutput( message: AssistantMessage, proposals: CapturedProposal[], identity: OutputContractIdentity, registry: OutputContractRegistry,): SemanticOutput { const diagnostic = finalOutputDiagnostic(message); const textParts = message.content.filter((part) => part.type === "text"); const authoritativeParts = message.content.filter((part) => part.type !== "thinking"); if (proposals.length > 0 || authoritativeParts.length !== 1 || textParts.length !== 1) { throw invalidFinalOutput("expected-one-text-part", diagnostic, identity); } const text = textParts[0]!.text; if (text.trim().length === 0) throw invalidFinalOutput("empty-final-text", diagnostic, identity); if (text.length > MAX_FINAL_JSON_CHARS) throw invalidFinalOutput("final-text-too-large", diagnostic, identity); let parsed: unknown; try { parsed = JSON.parse(text); } catch { throw invalidFinalOutput("invalid-json", diagnostic, identity); } try { return registry.validate(identity, parsed); } catch (error) { if (!(error instanceof OutputContractValidationError)) throw error; throw invalidFinalOutput("output-contract-invalid", { ...diagnostic, validationIssues: error.issues.map((issue) => ({ code: issue.code, path: issue.path })), }, identity); }}
function parseConversationText( message: AssistantMessage, proposals: CapturedProposal[], identity: OutputContractIdentity, registry: OutputContractRegistry,): ObservationOutput { const diagnostic = finalOutputDiagnostic(message); const textParts = message.content.filter((part) => part.type === "text"); const toolCallParts = message.content.filter((part) => part.type === "toolCall"); const otherAuthoritativeParts = message.content.filter((part) => part.type !== "thinking" && part.type !== "text" && part.type !== "toolCall"); if (textParts.length > 1 || toolCallParts.length !== proposals.length || otherAuthoritativeParts.length > 0) { throw invalidConversationText("invalid-conversation-parts", diagnostic, identity); } const visible = textParts[0]?.text.trim() ?? ""; const summary = visible || deterministicProposalAcknowledgment(proposals); if (summary.length === 0) throw invalidConversationText("empty-final-text", diagnostic, identity); if (summary.length > MAX_CONVERSATION_TEXT_CHARS) { throw invalidConversationText("final-text-too-large", diagnostic, identity); } const output = registry.validate(identity, { summary, tags: ["conversation"], importance: "normal", confidence: 0.5, }); if (!("tags" in output)) throw invalidConversationText("output-contract-invalid", diagnostic, identity); return output;}
function parseCompactionText( message: AssistantMessage, proposals: CapturedProposal[], identity: OutputContractIdentity, registry: OutputContractRegistry,): ConversationCompactionOutput { const diagnostic = finalOutputDiagnostic(message); const textParts = message.content.filter((part) => part.type === "text"); const toolCallParts = message.content.filter((part) => part.type === "toolCall"); const otherAuthoritativeParts = message.content.filter((part) => ( part.type !== "thinking" && part.type !== "text" && part.type !== "toolCall" )); if (textParts.length !== 1 || toolCallParts.length > 0 || proposals.length > 0 || otherAuthoritativeParts.length > 0) { throw invalidConversationText("invalid-compaction-parts", diagnostic, identity, "compaction-text"); } const boundary = textParts[0]!.text.trim(); if (boundary.length === 0) { throw invalidConversationText("empty-final-text", diagnostic, identity, "compaction-text"); } if (boundary.length > MAX_COMPACTION_TEXT_CHARS) { throw invalidConversationText("final-text-too-large", diagnostic, identity, "compaction-text"); } const summary = boundary.length <= 2_000 ? boundary : `${boundary.slice(0, 1_999)}…`; const output = registry.validate(identity, { summary, boundary, openLoops: [], decisions: [], exactReferences: [], unresolved: [], lookupHints: [], confidence: 0.5, }); if (!("boundary" in output)) { throw invalidConversationText("output-contract-invalid", diagnostic, identity, "compaction-text"); } return output;}
function deterministicProposalAcknowledgment(proposals: CapturedProposal[]): string { const kinds = new Set(proposals.map((proposal) => proposal.kind)); if (kinds.size === 1 && kinds.has("focus-declaration")) return "I saved that as a focus proposal."; if (kinds.has("focus-declaration")) return "I saved those as focus and agent suggestions."; if (kinds.has("memory-change") && kinds.has("self-correction")) return "I saved those as memory and correction suggestions."; if (kinds.has("memory-change")) return "I saved that as a memory suggestion."; if (kinds.has("self-correction")) return "I saved that as a proposed correction."; return "";}
function invalidConversationText( reason: string, diagnostic: JsonObject, identity: OutputContractIdentity, outputMode: "conversation-text" | "compaction-text" = "conversation-text",): AgentRunFailure { const expected = outputMode === "conversation-text" ? "conversation text" : "compaction text"; return new AgentRunFailure(`Pi final output rejected: expected one bounded ${expected} part`, { diagnostic: { ...diagnostic, code: "invalid-final-output", stage: "final-output-validation", reason, outputMode, outputContract: outputContractIdentityJson(identity), }, });}
function invalidFinalOutput( reason: string, diagnostic: JsonObject, identity: OutputContractIdentity,): AgentRunFailure { return new AgentRunFailure("Pi final output rejected: expected one strict JSON object", { diagnostic: { ...diagnostic, code: "invalid-final-output", stage: "final-output-validation", reason, outputContract: outputContractIdentityJson(identity), }, });}
function finalOutputDiagnostic(message: AssistantMessage): JsonObject { const textParts = message.content.filter((part) => part.type === "text"); const thinkingParts = message.content.filter((part) => part.type === "thinking"); const toolCallParts = message.content.filter((part) => part.type === "toolCall"); const text = textParts.map((part) => part.text).join(""); const thinking = thinkingParts.map((part) => part.thinking).join(""); const stopReason = safeStopReason(message.stopReason); return { assistantMessages: 1, textParts: textParts.length, textChars: text.length, ...(text.length > 0 ? { textSha256: sha256Text(text) } : {}), thinkingParts: thinkingParts.length, thinkingChars: thinking.length, thinkingRedacted: true, ...(thinking.length > 0 ? { thinkingSha256: sha256Text(thinking) } : {}), toolCallParts: toolCallParts.length, otherParts: message.content.length - textParts.length - thinkingParts.length - toolCallParts.length, ...(stopReason ? { stopReason } : {}), };}
function stringMetadata(value: string): JsonObject { return { chars: value.length, sha256: sha256Text(value), redacted: true };}
function safeStopReason(value: unknown): string | undefined { return typeof value === "string" && ["stop", "length", "toolUse", "error", "aborted"].includes(value) ? value : undefined;}
function sha256Text(value: string): string { return createHash("sha256").update(value).digest("hex");}
function asRecord(value: unknown): Record<string, unknown> | undefined { return value && typeof value === "object" && !Array.isArray(value) ? value as Record<string, unknown> : undefined;}
/** * Resolve content-addressed image artifact references to base64 ImageContent. * The trusted Pi parent validates each artifact beneath the artifact root: * rejects symlinks, path escapes, hash/MIME/magic mismatches, and size violations. * Failed resolutions are content-dark in traces and fail the run closed. A * message that claims to contain an image must not silently become a text-only * turn because its durable evidence disappeared or changed. */async function resolveContextImages( artifacts: ImageArtifactReference[], artifactRoot: string, onTrace: (trace: RunnerTrace) => Promise<void>,): Promise<Array<ImageContent & { mimeType: ImageArtifactReference["mimeType"] }>> { const images: Array<ImageContent & { mimeType: ImageArtifactReference["mimeType"] }> = []; for (const artifact of artifacts) { try { const image = await resolveImageArtifact(artifact, artifactRoot); images.push(image); await onTrace({ kind: "pi.image_resolved", data: { sha256: artifact.sha256, mimeType: artifact.mimeType, sizeBytes: artifact.sizeBytes }, }); } catch (error) { await onTrace({ kind: "pi.image_resolution_failed", data: { sha256: artifact.sha256, code: "image-artifact-invalid" }, }); throw new AgentRunFailure("Current-turn image artifact could not be validated", { cause: error, advanceProgress: false, diagnostic: { code: "image-artifact-invalid", stage: "trusted-evidence-resolution" }, }); } } return images;}