import 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 { 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): Promise { 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 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, ): Promise { 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["tools"], onTrace: (trace: RunnerTrace) => Promise, 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(); 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 | undefined { return value && typeof value === "object" && !Array.isArray(value) ? value as Record : 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, ): Promise> { const images: Array = []; 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; }