import { Agent, type AgentEvent, type AgentMessage, type AgentTool, type StreamFn } from "@earendil-works/pi-agent-core"; import type { AssistantMessage, Model, ToolResultMessage, UserMessage } from "@earendil-works/pi-ai"; import { streamSimple as streamSimpleOpenAICompletions } from "@earendil-works/pi-ai/api/openai-completions"; import { Type } from "typebox"; import net from "node:net"; import { capturedProposalsSchema, correctionArgumentsSchema, focusDeclarationArgumentsSchema, memoryChangeArgumentsSchema, validateCapturedProposalAgainstCapabilities, type CapturedProposal, type ProposalCapabilities, } from "../proposals.js"; import { decodeSingleFrame, encodeFrame, MAX_RESULT_FRAME_BYTES, MAX_RUN_PACKET_BYTES, MAX_TRACE_BYTES, MAX_TRACE_ITEMS, SANDBOX_PROTOCOL_VERSION, sandboxRunPacketSchema, type BrokerRequest, type BrokerResponse, type SandboxResult, } from "./protocol.js"; import type { ConversationMessage } from "../conversation-history.js"; async function main(): Promise { let runId = "invalid"; let providerFailure: "timeout" | "rejected" | "client" | "transient" | undefined; let failureReason: WorkerFailureReason | undefined; let proposalFailureCode: ProposalFailureCode | undefined; let proposalToolError = false; const traces: Array<{ kind: string; data: unknown }> = []; try { const raw = await readStream(process.stdin, MAX_RUN_PACKET_BYTES + 4); const packet = sandboxRunPacketSchema.parse(decodeSingleFrame(raw, MAX_RUN_PACKET_BYTES)); runId = packet.runId; const proposals: CapturedProposal[] = []; let observedRevision: string | undefined; let brokerUsed = false; globalThis.fetch = async (input: string | URL | Request, init?: RequestInit): Promise => { if (brokerUsed) throw new Error("Provider capability permits exactly one request"); brokerUsed = true; const request = input instanceof Request ? input : new Request(input, init); const url = new URL(request.url); if (url.origin !== packet.broker.origin || url.pathname !== packet.broker.route || request.method !== "POST") { throw new Error("Provider destination is not allowlisted"); } const body = await request.text(); const parsed = JSON.parse(body) as Record; const constrainedBody = packet.proposals ? JSON.stringify({ ...parsed, tool_choice: "auto" }) : packet.model.jsonSchemaResponseFormat ? JSON.stringify({ ...parsed, response_format: { type: "json_schema", json_schema: { name: "thoughtstream_conceptualization", strict: true, schema: packet.model.jsonSchemaResponseFormat, }, }, }) : packet.model.jsonObjectResponseFormat ? JSON.stringify({ ...parsed, response_format: { type: "json_object" } }) : body; const headers: Record = {}; for (const name of ["accept", "content-type"]) { const value = request.headers.get(name); if (value) headers[name] = value; } const brokerRequest: BrokerRequest = { version: 1, capability: packet.broker.capability, runId: packet.runId, origin: url.origin, path: url.pathname, method: "POST", model: typeof parsed.model === "string" ? parsed.model : "", headers, body: constrainedBody, }; try { const brokerResponse = await requestBroker(packet.broker.socketPath, brokerRequest); if (brokerResponse.status !== "completed") { providerFailure = brokerResponse.code === "provider-timeout" ? "timeout" : "rejected"; throw new Error("Provider broker rejected request"); } const response = new Response(brokerResponse.body ?? "", { status: brokerResponse.httpStatus ?? 502, ...(brokerResponse.headers ? { headers: brokerResponse.headers } : {}), }); if (response.status >= 400) { providerFailure = [408, 409, 425, 429].includes(response.status) || response.status >= 500 ? "transient" : "client"; } appendTrace(traces, "provider.proxy.completed", { status: response.status, bodyBytes: Buffer.byteLength(brokerResponse.body ?? ""), }); return response; } catch (error) { appendTrace(traces, "provider.proxy.failed", { classification: providerFailure ?? (error instanceof Error && error.name === "AbortError" ? "timeout" : "rejected-or-transport"), }); throw error; } }; const model = buildModel(packet); const priorMessages = packet.messages ?? []; const agent = new Agent({ initialState: { systemPrompt: packet.systemPrompt, model, thinkingLevel: packet.model.thinkingLevel, tools: createProposalTools(packet.proposals, proposals), messages: buildPriorAgentMessages(priorMessages), }, streamFn: streamSimpleOpenAICompletions as StreamFn, getApiKey: () => "broker-placeholder", onPayload: async (payload) => { appendTrace(traces, "provider.request", summarizePayload(payload)); return undefined; }, onResponse: async (response) => { observedRevision = safeRevision( response.headers["x-tinker-checkpoint"] ?? response.headers["x-model-revision"] ?? response.headers["x-checkpoint-revision"], ); appendTrace(traces, "provider.response", { status: response.status, headers: safeHeaders(response.headers) }); }, maxRetryDelayMs: 0, }); agent.subscribe(async (event) => { if (event.type === "tool_execution_end" && event.isError) { proposalToolError = true; proposalFailureCode = classifyProposalFailure(event.result); } appendTrace(traces, `pi.${event.type}`, compactAgentEvent(event)); }); await agent.prompt(packet.prompt, packet.images); let finalMessage: AssistantMessage; try { finalMessage = extractFinalAssistant(agent.state.messages); } catch (error) { failureReason = "assistant-missing"; throw error; } if (finalMessage.stopReason === "error") { failureReason = proposalToolError ? "proposal-tool-error" : "assistant-error-stop"; throw new Error("Model execution failed"); } let validatedProposals: CapturedProposal[]; try { validatedProposals = capturedProposalsSchema.parse(proposals); } catch (error) { failureReason = "proposal-result-invalid"; throw error; } writeResult({ version: SANDBOX_PROTOCOL_VERSION, runId: packet.runId, status: "completed", finalMessage, proposals: validatedProposals, traces, ...(observedRevision ? { observedRevision } : {}), }); } catch (error) { const invalidPacket = runId === "invalid"; writeResult({ version: SANDBOX_PROTOCOL_VERSION, runId, status: "failed", code: invalidPacket ? "invalid-packet" : classifyWorkerFailure(error, providerFailure), message: invalidPacket ? "Sandbox worker rejected its run packet" : "Sandboxed model execution failed", ...(!invalidPacket ? { reason: failureReason ?? "worker-exception", traces, ...(proposalFailureCode ? { proposalFailureCode } : {}), } : {}), }); } } const memoryChangeParameters = Type.Object({ operation: Type.Union([Type.Literal("append"), Type.Literal("replace-document")]), proposed_text: Type.String({ minLength: 1, maxLength: 32_768 }), reason: Type.String({ minLength: 1, maxLength: 1_000 }), evidence_event_ids: Type.Array(Type.String({ minLength: 1, maxLength: 500 }), { maxItems: 16, uniqueItems: true }), }, { additionalProperties: false }); const correctionParameters = Type.Object({ target_output: Type.String({ minLength: 1, maxLength: 500 }), replacement: Type.String({ minLength: 1, maxLength: 4_096 }), reason: Type.String({ minLength: 1, maxLength: 1_000 }), evidence_event_ids: Type.Array(Type.String({ minLength: 1, maxLength: 500 }), { maxItems: 16, uniqueItems: true }), }, { additionalProperties: false }); const focusDeclarationParameters = Type.Object({ declaration: Type.Object({ id: Type.String({ minLength: 1, maxLength: 100, pattern: "^[a-z0-9][a-z0-9-]*$" }), version: Type.Integer({ minimum: 1 }), name: Type.String({ minLength: 1, maxLength: 160 }), description: Type.String({ minLength: 1, maxLength: 2_000 }), parentFocusIds: Type.Array(Type.String({ minLength: 1, maxLength: 100, pattern: "^[a-z0-9][a-z0-9-]*$" }), { maxItems: 16, uniqueItems: true }), scope: Type.Object({ summary: Type.String({ minLength: 1, maxLength: 4_000 }), includeTags: Type.Array(Type.String({ minLength: 1, maxLength: 100, pattern: "^[a-z0-9]+(?:-[a-z0-9]+)*$" }), { maxItems: 100, uniqueItems: true }), excludeTags: Type.Array(Type.String({ minLength: 1, maxLength: 100, pattern: "^[a-z0-9]+(?:-[a-z0-9]+)*$" }), { maxItems: 100, uniqueItems: true }), }, { additionalProperties: false }), objective: Type.Object({ statement: Type.String({ minLength: 1, maxLength: 4_000 }) }, { additionalProperties: false }), subscriptions: Type.Object({ eventTypes: Type.Array(Type.String({ minLength: 1, maxLength: 300, pattern: "^stream\\.thought\\.[a-z0-9.*-]+$" }), { minItems: 1, maxItems: 100, uniqueItems: true }), sourcePatterns: Type.Array(Type.String({ minLength: 1, maxLength: 200, pattern: "^(?:\\*|[a-z0-9][a-z0-9._:-]*(?:\\*)?)$", }), { minItems: 1, maxItems: 100, uniqueItems: true }), privacy: Type.Array(Type.Union([Type.Literal("public-source"), Type.Literal("private"), Type.Literal("sensitive")]), { minItems: 1, maxItems: 3, uniqueItems: true }), replay: Type.Union([Type.Literal("beginning"), Type.Literal("now")]), }, { additionalProperties: false }), budgets: Type.Object({ inference: Type.Object({ maxCallsPerDay: Type.Integer({ minimum: 1, maximum: 1_000_000 }), maxInputTokensPerDay: Type.Integer({ minimum: 1, maximum: 1_000_000_000 }), maxOutputTokensPerDay: Type.Integer({ minimum: 1, maximum: 1_000_000_000 }), maxCostMicrousdPerDay: Type.Optional(Type.Integer({ minimum: 1, maximum: 1_000_000_000_000 })), }, { additionalProperties: false }), attention: Type.Object({ maxDeliveriesPerDay: Type.Integer({ minimum: 0, maximum: 10_000 }) }, { additionalProperties: false }), }, { additionalProperties: false }), permissions: Type.Object({ requestedReadTools: Type.Array(Type.String({ minLength: 1, maxLength: 200, pattern: "^[a-z0-9][a-z0-9._:-]*$" }), { maxItems: 32, uniqueItems: true }), requestedExternalActions: Type.Array(Type.String({ minLength: 1, maxLength: 200, pattern: "^[a-z0-9][a-z0-9._:-]*$" }), { maxItems: 32, uniqueItems: true }), }, { additionalProperties: false }), retirementRule: Type.String({ minLength: 1, maxLength: 2_000 }), }, { additionalProperties: false }), }, { additionalProperties: false }); function createProposalTools( capabilities: ProposalCapabilities | undefined, captured: CapturedProposal[], ): AgentTool[] { if (!capabilities) return []; const record = (proposal: CapturedProposal) => { if (captured.length >= 3) throw new Error("Proposal call limit exceeded"); if (captured.some((prior) => prior.kind === proposal.kind)) throw new Error("Duplicate proposal kind"); captured.push(validateCapturedProposalAgainstCapabilities(proposal, capabilities)); return { content: [{ type: "text" as const, text: "Suggestion captured for trusted review." }], details: { captured: true, kind: proposal.kind }, terminate: true, }; }; const tools: AgentTool[] = []; if (capabilities.enabled.includes("memory-change")) { tools.push({ name: "request_memory_change", label: "Suggest memory change", description: "Submit one bounded, non-operative memory suggestion for trusted human review. This does not change memory.", parameters: memoryChangeParameters, executionMode: "sequential", execute: async (toolCallId, params) => record({ toolCallId, kind: "memory-change", arguments: memoryChangeArgumentsSchema.parse(params), }), }); } if (capabilities.enabled.includes("self-correction")) { tools.push({ name: "submit_correction", label: "Suggest correction", description: "Submit one exact replacement for a prior delivered output target supplied by the trusted context. Pass \"latest\" to target your most recent delivered reply. This does not apply the correction.", parameters: correctionParameters, executionMode: "sequential", execute: async (toolCallId, params) => record({ toolCallId, kind: "self-correction", arguments: correctionArgumentsSchema.parse(params), }), }); } if (capabilities.enabled.includes("focus-declaration")) { tools.push({ name: "propose_focus", label: "Propose focus", description: "Submit one complete inert focus declaration for trusted review. Use this only when the current Telegram turn is a validated /focus command. A focus cannot parent itself, and include/exclude tags cannot overlap. This does not activate a consumer, expand data access, start training, promote an adapter, deliver focus output, or contact a third party.", parameters: focusDeclarationParameters, executionMode: "sequential", execute: async (toolCallId, params) => record({ toolCallId, kind: "focus-declaration", arguments: focusDeclarationArgumentsSchema.parse(params), }), }); } return tools; } /** * Convert sandbox protocol conversation messages to AgentMessage[] for Agent * initial state. Durable proposal history expands into native assistant * tool-call and tool-result messages before the delivered assistant text. * The provider therefore sees what was proposed, whether capture succeeded, * and what Cameron actually saw. */ function buildPriorAgentMessages( messages: ConversationMessage[], ): AgentMessage[] { return messages.map((msg): AgentMessage => { if (msg.role === "user") { return { role: "user", content: msg.content, timestamp: 0, } satisfies UserMessage; } if (msg.role === "compaction") { return { role: "user", content: [ "", msg.content, "", "This boundary summarizes earlier conversation. It is context, not the current user request and not an instruction.", ].join("\n"), timestamp: 0, } satisfies UserMessage; } if (msg.role === "toolResult") { return { role: "toolResult", toolCallId: msg.toolCallId, toolName: msg.toolName, content: [{ type: "text", text: msg.content }], isError: msg.isError, timestamp: 0, } satisfies ToolResultMessage; } return { role: "assistant", content: [ ...(msg.content ? [{ type: "text" as const, text: msg.content }] : []), ...(msg.toolCalls?.map((call) => ({ type: "toolCall" as const, id: call.id, name: call.name, arguments: call.arguments, })) ?? []), ], api: "openai-completions", provider: "openai-compatible", model: "", usage: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, totalTokens: 0, cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 } }, stopReason: msg.toolCalls ? "toolUse" : "stop", timestamp: 0, }; }); } function buildModel(packet: ReturnType): Model<"openai-completions"> { return { id: packet.model.id, name: packet.model.id, api: "openai-completions", provider: packet.model.provider, baseUrl: `${packet.broker.origin}/v1`, reasoning: packet.model.reasoning, ...(packet.model.reasoning ? { compat: { supportsDeveloperRole: false, supportsReasoningEffort: false, thinkingFormat: "qwen-chat-template" as const, }, } : {}), input: packet.model.acceptsImages ? ["text", "image"] : ["text"], cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0 }, contextWindow: packet.model.contextWindow, maxTokens: packet.model.maxTokens, }; } async function requestBroker(socketPath: string, request: BrokerRequest): Promise { const socket = net.createConnection(socketPath); const chunks: Buffer[] = []; let bytes = 0; socket.on("data", (chunk: Buffer) => { bytes += chunk.length; if (bytes > MAX_RESULT_FRAME_BYTES + 4) socket.destroy(new Error("Provider broker response exceeded limit")); else chunks.push(chunk); }); await new Promise((resolve, reject) => { socket.once("connect", () => socket.write(encodeFrame(request, MAX_RUN_PACKET_BYTES))); socket.once("end", resolve); socket.once("error", reject); }); return decodeSingleFrame(Buffer.concat(chunks), MAX_RESULT_FRAME_BYTES) as BrokerResponse; } function extractFinalAssistant(messages: unknown[]): AssistantMessage { const assistants = messages.filter((message): message is AssistantMessage => Boolean( message && typeof message === "object" && "role" in message && message.role === "assistant", )); const last = assistants.at(-1); if (!last) throw new Error("Model returned no assistant message"); return last; } function appendTrace(traces: Array<{ kind: string; data: unknown }>, kind: string, data: unknown): void { if (traces.length >= MAX_TRACE_ITEMS) return; const serialized = JSON.stringify(data); traces.push({ kind, data: Buffer.byteLength(serialized) <= MAX_TRACE_BYTES ? data : { omitted: "trace-size-limit", bytes: Buffer.byteLength(serialized) }, }); } function summarizePayload(payload: unknown): unknown { if (!payload || typeof payload !== "object") return { kind: typeof payload }; const body = payload as Record; return { model: body.model, stream: body.stream, messageCount: Array.isArray(body.messages) ? body.messages.length : undefined, toolCount: Array.isArray(body.tools) ? body.tools.length : 0, }; } function compactAgentEvent(event: AgentEvent): unknown { if (event.type === "message_update") { const update = event.assistantMessageEvent; return { type: event.type, updateType: update.type }; } if (event.type === "message_end") { const message = event.message as AssistantMessage; return { type: event.type, role: message.role, stopReason: message.stopReason, content: message.content.map((part) => ({ type: part.type, ...(part.type === "text" ? { chars: part.text.length } : {}), ...(part.type === "thinking" ? { chars: part.thinking.length } : {}), ...(part.type === "toolCall" ? { name: part.name } : {}), })), }; } if (event.type === "tool_execution_start") { return { type: event.type, toolName: event.toolName }; } if (event.type === "tool_execution_end") { return { type: event.type, toolName: event.toolName, isError: event.isError }; } return { type: event.type }; } function safeRevision(value: string | undefined): string | undefined { return value && /^[A-Za-z0-9._:/@+-]{1,200}$/.test(value) ? value : undefined; } function safeHeaders(headers: Record): Record { return Object.fromEntries(Object.entries(headers).filter(([name]) => [ "content-type", "x-tinker-checkpoint", "x-model-revision", "x-checkpoint-revision", ].includes(name.toLowerCase()))); } type WorkerFailureReason = "assistant-missing" | "assistant-error-stop" | "proposal-tool-error" | "proposal-result-invalid" | "worker-exception"; type ProposalFailureCode = "arguments-invalid" | "evidence-outside-snapshot" | "target-outside-snapshot" | "proposal-call-limit" | "duplicate-proposal-kind" | "proposal-tool-unknown"; function classifyProposalFailure(result: unknown): ProposalFailureCode { const serialized = safeErrorText(result); if (/evidence id is outside the context snapshot/i.test(serialized)) return "evidence-outside-snapshot"; if (/correction target is outside the context snapshot/i.test(serialized)) return "target-outside-snapshot"; if (/proposal call limit exceeded/i.test(serialized)) return "proposal-call-limit"; if (/duplicate proposal kind/i.test(serialized)) return "duplicate-proposal-kind"; if (/invalid.*argument|argument.*invalid|validation|schema/i.test(serialized)) return "arguments-invalid"; return "proposal-tool-unknown"; } function safeErrorText(value: unknown): string { if (!value || typeof value !== "object") return ""; const record = value as Record; const parts: string[] = []; if (typeof record.error === "string") parts.push(record.error); if (typeof record.message === "string") parts.push(record.message); if (Array.isArray(record.content)) { for (const item of record.content) { if (item && typeof item === "object" && "text" in item && typeof item.text === "string") parts.push(item.text); } } return parts.join("\n").slice(0, 4_096); } function classifyWorkerFailure( error: unknown, providerFailure: "timeout" | "rejected" | "client" | "transient" | undefined, ): "provider-rejected" | "provider-client-error" | "provider-transient-error" | "provider-timeout" | "model-failed" | "protocol-failed" | "output-limit" { if (providerFailure === "timeout") return "provider-timeout"; if (providerFailure === "rejected") return "provider-rejected"; if (providerFailure === "client") return "provider-client-error"; if (providerFailure === "transient") return "provider-transient-error"; const message = error instanceof Error ? error.message : String(error); if (/timeout/i.test(message)) return "provider-timeout"; if (/broker|provider|http|fetch/i.test(message)) return "provider-rejected"; if (/frame|protocol|packet|json/i.test(message)) return "protocol-failed"; if (/limit|too large|heap|memory/i.test(message)) return "output-limit"; return "model-failed"; } async function readStream(stream: NodeJS.ReadableStream, maxBytes: number): Promise { const chunks: Buffer[] = []; let bytes = 0; for await (const chunk of stream) { const buffer = Buffer.isBuffer(chunk) ? chunk : Buffer.from(chunk as string | Uint8Array); bytes += buffer.length; if (bytes > maxBytes) throw new Error("Sandbox input exceeded limit"); chunks.push(buffer); } return Buffer.concat(chunks); } function writeResult(result: SandboxResult): void { try { process.stdout.write(encodeFrame(result, MAX_RESULT_FRAME_BYTES)); } catch { const fallback: SandboxResult = { version: SANDBOX_PROTOCOL_VERSION, runId: result.runId, status: "failed", code: "output-limit", message: "Sandbox result exceeded its output limit", }; process.stdout.write(encodeFrame(fallback, MAX_RESULT_FRAME_BYTES)); } } void main();