Something went wrong. Try again.
A local-first event pipeline for independent agents, built on Jazz.
Something went wrong. Try again.
60 kB · 1412 lines
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101110211031104110511061107110811091110111111121113111411151116111711181119112011211122112311241125112611271128112911301131113211331134113511361137113811391140114111421143114411451146114711481149115011511152115311541155115611571158115911601161116211631164116511661167116811691170117111721173117411751176117711781179118011811182118311841185118611871188118911901191119211931194119511961197119811991200120112021203120412051206120712081209121012111212121312141215121612171218121912201221122212231224122512261227122812291230123112321233123412351236123712381239124012411242124312441245124612471248124912501251125212531254125512561257125812591260126112621263126412651266126712681269127012711272127312741275127612771278127912801281128212831284128512861287128812891290129112921293129412951296129712981299130013011302130313041305130613071308130913101311131213131314131513161317131813191320132113221323132413251326132713281329133013311332133313341335133613371338133913401341134213431344134513461347134813491350135113521353135413551356135713581359136013611362136313641365136613671368136913701371137213731374137513761377137813791380138113821383138413851386138713881389139013911392139313941395139613971398139914001401140214031404140514061407140814091410141114121413import http from "node:http";import path from "node:path";import os from "node:os";import fs from "node:fs/promises";import { createHash } from "node:crypto";import YAML from "yaml";import { canonicalJson } from "../src/core/json.js";import { afterEach, describe, expect, test } from "vitest";import { buildContextPacket } from "../src/agents/context.js";import { loadAgentDeclarations } from "../src/agents/declarations.js";import { PiAgentRunner } from "../src/agents/pi.js";import { staticProviderProfileResolver, type ProviderProfile } from "../src/agents/provider-profiles.js";import type { ThoughtAgentDeclaration } from "../src/agents/types.js";import type { ThoughtEvent } from "../src/events/types.js";import { telegramFocusDescription } from "../src/agents/telegram-help.js";import { CONVERSATION_COMPACTION_OUTPUT_CONTRACT } from "../src/agents/output-contracts.js";
const servers: http.Server[] = [];const temporaryRoots: string[] = [];const workerBundlePath = path.resolve("dist/sandbox/worker.cjs");
afterEach(async () => { delete process.env.THOUGHTSTREAM_TEST_API_KEY; delete process.env.THOUGHTSTREAM_TINKER_ALLOWED_MODELS; await Promise.all(temporaryRoots.splice(0).map((root) => fs.rm(root, { recursive: true, force: true }))); await Promise.all(servers.splice(0).map((server) => new Promise<void>((resolve) => { server.closeAllConnections(); server.close(() => resolve()); })));});
describe("PiAgentRunner", () => { test("normalizes private and current-bot /focus descriptions only", () => { const event = fixtureRunInput(fixtureDeclaration()).event; event.type = "stream.thought.source.telegram.message"; event.sourceKind = "telegram"; event.payload = { text: "/focus AI news", accountUsername: "TheStreamBot" }; expect(telegramFocusDescription(event)).toBe("AI news"); event.payload = { text: "/focus@TheStreamBot AI news", accountUsername: "TheStreamBot" }; expect(telegramFocusDescription(event)).toBe("AI news"); event.payload = { text: "/focus@OtherBot AI news", accountUsername: "TheStreamBot" }; expect(telegramFocusDescription(event)).toBeUndefined(); event.payload = { text: "/focus", accountUsername: "TheStreamBot" }; expect(telegramFocusDescription(event)).toBeUndefined(); });
test("runs the model in the disposable sandbox through exactly one provider request", async () => { let requestCount = 0; let authorizationWasInjected = false; let receivedBody = ""; const server = await startServer((request, response) => { requestCount += 1; authorizationWasInjected = request.headers.authorization === "Bearer fixture-secret"; let body = ""; request.on("data", (chunk) => { body += chunk; }); request.on("end", () => { receivedBody = body; respondWithOutput(response, validOutput("The fixture changed"), { promptTokens: 321, completionTokens: 45 }); }); }); process.env.THOUGHTSTREAM_TEST_API_KEY = "fixture-secret"; const declaration = fixtureDeclaration(); const traces: Array<{ kind: string; data: unknown }> = []; const input = fixtureRunInput(declaration); input.context.systemText = "TRUSTED SUBSCRIBED DOCUMENT SENTINEL";
const output = await fixtureRunner(server).run( input, async (trace) => { traces.push(trace); }, );
expect(output.summary).toBe("The fixture changed"); expect(output.model).toEqual({ provider: "openai-compatible", id: "fixture-model" }); expect(output.usage).toEqual({ inputTokens: 321, outputTokens: 45 }); expect(requestCount).toBe(1); expect(authorizationWasInjected).toBe(true); expect(receivedBody).toContain("thoughtstream-source-event"); expect(receivedBody).not.toContain("fixture-secret"); const requestBody = JSON.parse(receivedBody) as { messages: Array<{ role: string; content: string | Array<{ type: string; text?: string }> }>; response_format: { type: string }; }; expect(requestBody).toMatchObject({ response_format: { type: "json_object" } }); const systemContent = requestBody.messages.find((message) => message.role === "system")?.content ?? ""; const renderedSystem = typeof systemContent === "string" ? systemContent : systemContent.map((part) => part.text ?? "").join(""); expect(renderedSystem).toContain("TRUSTED SUBSCRIBED DOCUMENT SENTINEL"); const finalContent = requestBody.messages.at(-1)?.content ?? ""; const finalPrompt = typeof finalContent === "string" ? finalContent : finalContent.map((part) => part.text ?? "").join(""); expect(finalPrompt).toContain("## Required final answer"); expect(finalPrompt).not.toContain("TRUSTED SUBSCRIBED DOCUMENT SENTINEL"); expect(finalPrompt).toContain('importance value must be exactly one of "low", "normal", or "high"'); expect(finalPrompt.indexOf("thoughtstream-source-event")).toBeLessThan(finalPrompt.indexOf("## Required final answer")); expect(traces.map((trace) => trace.kind)).toEqual(expect.arrayContaining([ "sandbox.provider.request", "sandbox.provider.response", "provider.request", "provider.response", "pi.agent_start", "pi.agent_end", "pi.message_updates_coalesced", ])); expect(traces.map((trace) => trace.kind)).not.toContain("pi.message_update"); const serializedTraces = JSON.stringify(traces); expect(serializedTraces).not.toContain("fixture-secret"); expect(serializedTraces).not.toContain("fixture.md"); expect(serializedTraces).not.toContain("The fixture changed"); expect(serializedTraces).not.toContain("TRUSTED SUBSCRIBED DOCUMENT SENTINEL"); expect(serializedTraces).toContain("redacted"); }, 15_000);
test("uses a process-local learned checkpoint while traces and output retain only public adapter identity", async () => { const privateCheckpoint = "tinker://private-adapter-checkpoint-fixture"; let receivedBody = ""; const server = await startServer((request, response) => { request.on("data", (part) => { receivedBody += part; }); request.on("end", () => { response.setHeader("x-tinker-checkpoint", privateCheckpoint); respondWithOutput(response, validOutput("Adapted output")); }); }); process.env.THOUGHTSTREAM_TEST_API_KEY = "fixture-secret"; const declaration = await compiledAdapterDeclaration(privateCheckpoint); const traces: Array<{ kind: string; data: unknown }> = [];
const output = await fixtureRunner(server, { allowedModels: [privateCheckpoint], profileId: "tinker-default", provider: "tinker", }).run( fixtureRunInput(declaration), async (trace) => { traces.push(trace); }, );
expect(JSON.parse(receivedBody)).toMatchObject({ model: privateCheckpoint }); expect(output.model).toEqual({ provider: "tinker", id: "Qwen/Qwen3.5-4B", revision: `sha256:${createHash("sha256").update(privateCheckpoint).digest("hex")}`, }); expect(JSON.stringify(traces)).not.toContain(privateCheckpoint); expect(JSON.stringify(traces)).toContain("Qwen/Qwen3.5-4B"); expect(JSON.stringify(output)).not.toContain(privateCheckpoint);
receivedBody = ""; await expect(fixtureRunner(server, { allowedModels: [privateCheckpoint], profileId: "tinker-default", provider: "tinker", }).run( fixtureRunInput(structuredClone(declaration)), async () => undefined, )).rejects.toThrow("not compiled by the trusted declaration loader"); expect(receivedBody).toBe(""); }, 15_000);
test("never persists a raw learned-checkpoint revision on adapter failure", async () => { const privateCheckpoint = "tinker://private-adapter-failure-revision"; const server = await startServer((_request, response) => { response.setHeader("x-tinker-checkpoint", privateCheckpoint); respondWithText(response, "not-json"); }); process.env.THOUGHTSTREAM_TEST_API_KEY = "fixture-secret"; const declaration = await compiledAdapterDeclaration(privateCheckpoint);
try { await fixtureRunner(server, { allowedModels: [privateCheckpoint], profileId: "tinker-default", provider: "tinker", }).run(fixtureRunInput(declaration), async () => undefined); throw new Error("Expected invalid adapter output"); } catch (error) { expect(error).toMatchObject({ diagnostic: expect.objectContaining({ checkpointRevision: `sha256:${createHash("sha256").update(privateCheckpoint).digest("hex")}`, }), }); expect(JSON.stringify(error)).not.toContain(privateCheckpoint); expect(String(error)).not.toContain(privateCheckpoint); } }, 15_000);
test("does not begin adapter-backed provider egress until read-only prefetch is exhausted", async () => { const privateCheckpoint = "tinker://prefetch-boundary-fixture"; const server = await startServer((_request, response) => respondWithOutput(response, validOutput("After prefetch"))); process.env.THOUGHTSTREAM_TEST_API_KEY = "fixture-secret"; const declaration = await compiledAdapterDeclaration(privateCheckpoint, { atprotoPrefetch: true }); let prefetchStarted!: () => void; let releasePrefetch!: () => void; const started = new Promise<void>((resolve) => { prefetchStarted = resolve; }); const released = new Promise<void>((resolve) => { releasePrefetch = resolve; }); let upstreamRequests = 0; const base = fixtureRunInput(declaration); const input = { ...base, event: { ...base.event, id: "evt_adapter_prefetch", type: "stream.thought.source.atproto.commit", source: "jetstream:test", sourceKind: "jetstream" as const, externalId: "did:plc:test/app.bsky.feed.post/one", rootEventId: "evt_adapter_prefetch", payload: { atUri: "at://did:plc:test/app.bsky.feed.post/one", cid: "bafyfixture", collection: "app.bsky.feed.post", rkey: "one", operation: "create", }, }, }; const run = fixtureRunner(server, { allowedModels: [privateCheckpoint], profileId: "tinker-default", provider: "tinker", providerFetchImpl: async (...args) => { upstreamRequests += 1; return fetch(...args); }, fetchImpl: async () => { prefetchStarted(); await released; return new Response("missing", { status: 404 }); }, }).run(input, async () => undefined);
await started; expect(upstreamRequests).toBe(0); releasePrefetch(); await expect(run).resolves.toMatchObject({ summary: "After prefetch" }); expect(upstreamRequests).toBe(1); }, 15_000);
test("uses the broker-authorized model identity instead of a provider-returned field", async () => { const untrustedModelField = "provider-secret-fragment"; const server = await startServer((_request, response) => { response.writeHead(200, { "content-type": "text/event-stream" }); response.write(`data: ${JSON.stringify(chunk({ role: "assistant", content: JSON.stringify(validOutput("Bound model identity")) }, null, untrustedModelField))}\n\n`); response.write(`data: ${JSON.stringify(chunk({}, "stop", untrustedModelField))}\n\n`); response.end("data: [DONE]\n\n"); }); process.env.THOUGHTSTREAM_TEST_API_KEY = "fixture-secret";
const output = await fixtureRunner(server).run( fixtureRunInput(fixtureDeclaration()), async () => undefined, );
expect(output.model).toEqual({ provider: "openai-compatible", id: "fixture-model" }); expect(JSON.stringify(output)).not.toContain(untrustedModelField); }, 15_000);
test("executes read-only enrichment in the trusted parent and persists metadata only", async () => { let receivedBody: Record<string, unknown> | undefined; const server = await startServer((request, response) => { let body = ""; request.on("data", (part) => { body += part; }); request.on("end", () => { receivedBody = JSON.parse(body) as Record<string, unknown>; respondWithOutput(response, validOutput("Fetched the canonical post")); }); }); process.env.THOUGHTSTREAM_TEST_API_KEY = "fixture-secret"; const declaration = fixtureDeclaration({ tools: ["atproto.fetch-markdown"], acceptedPrivacy: ["public-source"], compiledEventTypes: ["stream.thought.source.atproto.commit"], }); const event: ThoughtEvent = { ...fixtureRunInput(declaration).event, id: "evt_atproto_pi", type: "stream.thought.source.atproto.commit", source: "jetstream:test", sourceKind: "jetstream", externalId: "did:plc:alice/app.bsky.feed.post/post-one", idempotencyKey: "atproto-pi", actor: "did:plc:alice", rootEventId: "evt_atproto_pi", privacy: "public-source", payload: { atUri: "at://did:plc:alice/app.bsky.feed.post/post-one", collection: "app.bsky.feed.post", operation: "create", record: { text: "fixture" }, }, }; const output = await fixtureRunner(server, { fetchImpl: async (input) => input.toString().startsWith("https://atproto.md/") ? new Response("# Canonical post\n\nFetched body", { headers: { "content-type": "text/markdown" } }) : new Response(JSON.stringify({ posts: [] }), { headers: { "content-type": "application/json" } }), }).run( { runId: "run_atproto_pi", declaration, event, context: buildContextPacket(declaration, event) }, async () => undefined, );
expect(output.enrichments).toEqual([expect.objectContaining({ tool: "atproto.fetch-markdown", status: "succeeded", requestKeys: ["target"], argumentsRedacted: true, resultPresent: true, resultSha256: expect.any(String), })]); expect(JSON.stringify(output.enrichments)).not.toContain("post-one"); expect(JSON.stringify(output.enrichments)).not.toContain("Fetched body"); expect(receivedBody?.tools).toBeUndefined(); expect(JSON.stringify(receivedBody)).toContain("## Pre-fetched read-only evidence"); expect(JSON.stringify(receivedBody)).toContain("Fetched body"); }, 15_000);
test("does not download or forward images unless the trusted model profile declares image input", async () => { let receivedBody: Record<string, unknown> | undefined; const fetched: string[] = []; const imageUrl = "https://cdn.bsky.app/img/feed_thumbnail/plain/did:plc:alice/image@jpeg"; const server = await startServer((request, response) => { let body = ""; request.on("data", (part) => { body += part; }); request.on("end", () => { receivedBody = JSON.parse(body) as Record<string, unknown>; respondWithOutput(response, validOutput("Inspected text evidence only")); }); }); process.env.THOUGHTSTREAM_TEST_API_KEY = "fixture-secret"; const declaration = fixtureDeclaration({ tools: ["atproto.fetch-markdown", "web.download-image"], acceptedPrivacy: ["public-source"], compiledEventTypes: ["stream.thought.source.atproto.commit"], }); const event: ThoughtEvent = { ...fixtureRunInput(declaration).event, id: "evt_text_only_image", type: "stream.thought.source.atproto.commit", source: "jetstream:test", sourceKind: "jetstream", externalId: "did:plc:alice/app.bsky.feed.post/post-image", idempotencyKey: "text-only-image", actor: "did:plc:alice", rootEventId: "evt_text_only_image", privacy: "public-source", payload: { atUri: "at://did:plc:alice/app.bsky.feed.post/post-image", collection: "app.bsky.feed.post", operation: "create", record: { text: "fixture" }, }, }; const output = await fixtureRunner(server, { fetchImpl: async (input) => { const url = input.toString(); fetched.push(url); if (url.startsWith("https://atproto.md/")) { return new Response(`# Post\n\n`, { headers: { "content-type": "text/markdown" } }); } return new Response(JSON.stringify({ posts: [] }), { headers: { "content-type": "application/json" } }); }, }).run( { runId: "run_text_only_image", declaration, event, context: buildContextPacket(declaration, event) }, async () => undefined, );
expect(output.summary).toBe("Inspected text evidence only"); expect(fetched).not.toContain(imageUrl); expect(JSON.stringify(receivedBody)).not.toContain("image_url"); expect(JSON.stringify(receivedBody)).toContain("trusted model profile is text-only"); expect(output.enrichments?.some((entry) => entry.tool === "web.download-image")).toBe(false); }, 15_000);
test("treats provider reasoning as non-authoritative and validates only the final text", async () => { const privateThinking = "SECRET_REASONING: preserve no provider trace text"; const outputText = JSON.stringify(validOutput("Strict final output")); const server = await startServer((_request, response) => { response.writeHead(200, { "content-type": "text/event-stream" }); response.write(`data: ${JSON.stringify(chunk({ role: "assistant", reasoning_content: privateThinking }, null))}\n\n`); response.write(`data: ${JSON.stringify(chunk({ content: outputText }, null))}\n\n`); response.write(`data: ${JSON.stringify(chunk({}, "stop"))}\n\n`); response.end("data: [DONE]\n\n"); }); process.env.THOUGHTSTREAM_TEST_API_KEY = "fixture-secret"; const traces: Array<{ kind: string; data: unknown }> = [];
const output = await fixtureRunner(server).run( fixtureRunInput(fixtureDeclaration()), async (trace) => { traces.push(trace); }, ); expect(output.summary).toBe("Strict final output"); expect(JSON.stringify(traces)).not.toContain(privateThinking); expect(JSON.stringify(traces)).not.toContain(outputText); }, 15_000);
test("maps the declaration thinking level to the Tinker Qwen chat-template flag", async () => { const requestBodies: Array<Record<string, unknown>> = []; const server = await startServer((request, response) => { let body = ""; request.on("data", (part) => { body += part; }); request.on("end", () => { requestBodies.push(JSON.parse(body) as Record<string, unknown>); respondWithOutput(response, validOutput("Reasoning configuration observed")); }); }); process.env.THOUGHTSTREAM_TEST_API_KEY = "fixture-secret"; const runner = fixtureRunner(server, { allowedModels: ["thinkingmachines/Inkling-Small"], profileId: "tinker-default", provider: "tinker", }); const base = { provider: "tinker" as const, providerProfile: "tinker-default", model: "thinkingmachines/Inkling-Small", };
await runner.run(fixtureRunInput(fixtureDeclaration({ ...base, thinkingLevel: "medium" })), async () => undefined); await runner.run(fixtureRunInput(fixtureDeclaration(base)), async () => undefined);
expect(requestBodies[0]).toMatchObject({ chat_template_kwargs: { enable_thinking: true, preserve_thinking: true }, }); expect(requestBodies[1]).toMatchObject({ chat_template_kwargs: { enable_thinking: false, preserve_thinking: true }, }); expect(requestBodies[0]).not.toHaveProperty("reasoning_effort"); expect(requestBodies[1]).not.toHaveProperty("reasoning_effort"); }, 15_000);
test("wraps one bounded Telegram conversation text part in the trusted observation envelope", async () => { let receivedBody: Record<string, unknown> | undefined; const server = await startServer((request, response) => { let body = ""; request.on("data", (part) => { body += part; }); request.on("end", () => { receivedBody = JSON.parse(body) as Record<string, unknown>; response.writeHead(200, { "content-type": "text/event-stream" }); respondWithText(response, " A small model can just answer now. ", false); }); }); process.env.THOUGHTSTREAM_TEST_API_KEY = "fixture-secret"; const declaration = fixtureDeclaration({ outputMode: "conversation-text", contextStrategy: "telegram-conversation", outputEventType: "stream.thought.derived.message.observation", emit: ["stream.thought.derived.message.observation"], acceptedPrivacy: ["sensitive"], compiledEventTypes: ["stream.thought.source.telegram.message"], eventTypes: ["stream.thought.source.telegram.message"], tools: [], });
const output = await fixtureRunner(server).run( fixtureRunInput(declaration), async () => undefined, );
expect(output).toMatchObject({ summary: "A small model can just answer now.", tags: ["conversation"], importance: "normal", confidence: 0.5, }); expect(receivedBody?.response_format).toBeUndefined(); expect(JSON.stringify(receivedBody)).toContain("Reply with concise visible text"); }, 15_000);
test("wraps one plain-text prefix summary in the trusted compaction envelope", async () => { let receivedBody: Record<string, unknown> | undefined; const server = await startServer((request, response) => { let body = ""; request.on("data", (part) => { body += part; }); request.on("end", () => { receivedBody = JSON.parse(body) as Record<string, unknown>; response.writeHead(200, { "content-type": "text/event-stream" }); respondWithText(response, " The user corrected the repeated identity introduction. The implementation review remains open. ", false); }); }); process.env.THOUGHTSTREAM_TEST_API_KEY = "fixture-secret"; const declaration = fixtureDeclaration({ role: "compactor", outputMode: "compaction-text", outputContract: { ...CONVERSATION_COMPACTION_OUTPUT_CONTRACT.identity }, contextStrategy: "telegram-compaction", outputEventType: "stream.thought.derived.conversation.compaction", emit: ["stream.thought.derived.conversation.compaction"], acceptedPrivacy: ["sensitive"], compiledEventTypes: ["stream.thought.source.telegram.message"], eventTypes: ["stream.thought.source.telegram.message"], tools: [], proposals: [], });
const output = await fixtureRunner(server).run( fixtureRunInput(declaration), async () => undefined, );
expect(output).toMatchObject({ summary: "The user corrected the repeated identity introduction. The implementation review remains open.", boundary: "The user corrected the repeated identity introduction. The implementation review remains open.", openLoops: [], decisions: [], exactReferences: [], unresolved: [], lookupHints: [], confidence: 0.5, }); expect(receivedBody?.response_format).toBeUndefined(); expect(JSON.stringify(receivedBody)).toContain("Return only the plain-text continuity summary"); expect(JSON.stringify(receivedBody)).not.toContain("Return exactly one raw JSON object"); }, 15_000);
test("sends prior Telegram turns as native roles while current runtime authority stays system-owned", async () => { let receivedBody: { messages?: Array<{ role: string; content: unknown }> } = {}; const server = await startServer((request, response) => { let body = ""; request.on("data", (part) => { body += part; }); request.on("end", () => { receivedBody = JSON.parse(body) as typeof receivedBody; response.writeHead(200, { "content-type": "text/event-stream" }); respondWithText(response, "A direct current reply.", false); }); }); process.env.THOUGHTSTREAM_TEST_API_KEY = "fixture-secret"; const declaration = fixtureDeclaration({ outputMode: "conversation-text", contextStrategy: "telegram-conversation", outputEventType: "stream.thought.derived.message.observation", emit: ["stream.thought.derived.message.observation"], acceptedPrivacy: ["sensitive"], compiledEventTypes: ["stream.thought.source.telegram.message"], eventTypes: ["stream.thought.source.telegram.message"], tools: [], }); const input = fixtureRunInput(declaration); input.context = { systemText: "CURRENT TRUSTED RUNTIME: Stream v17. Historical identity claims are untrusted.", messages: [ { role: "user", content: "What runtime are you using?" }, { role: "assistant", content: "I am v13 on an obsolete model." }, ], text: "Answer the current question without volunteering runtime internals.", manifest: { contextStrategy: "telegram-conversation" }, };
await fixtureRunner(server).run(input, async () => undefined);
const messages = receivedBody.messages ?? []; expect(messages.map((message) => message.role)).toEqual(["system", "user", "assistant", "user"]); expect(messageText(messages[0]?.content)).toContain("CURRENT TRUSTED RUNTIME: Stream v17"); expect(messageText(messages[1]?.content)).toBe("What runtime are you using?"); expect(messageText(messages[2]?.content)).toBe("I am v13 on an obsolete model."); expect(messageText(messages[3]?.content)).toBe("Answer the current question without volunteering runtime internals."); expect(messageText(messages[3]?.content)).not.toContain("I am v13 on an obsolete model."); expect(messageText(messages[3]?.content)).not.toContain("## Required final answer"); expect(messageText(messages[3]?.content)).not.toContain("Return only the reply text"); }, 15_000);
test("serializes reconstructed proposal calls, results, and delivered text as native provider history", async () => { let receivedBody: { messages?: Array<{ role: string; content: unknown; tool_calls?: Array<{ id: string; function: { name: string; arguments: string } }>; tool_call_id?: string; }>; } = {}; const server = await startServer((request, response) => { let body = ""; request.on("data", (part) => { body += part; }); request.on("end", () => { receivedBody = JSON.parse(body) as typeof receivedBody; response.writeHead(200, { "content-type": "text/event-stream" }); respondWithText(response, "The prior proposal exists.", false); }); }); process.env.THOUGHTSTREAM_TEST_API_KEY = "fixture-secret"; const declaration = fixtureDeclaration({ outputMode: "conversation-text", contextStrategy: "telegram-conversation", outputEventType: "stream.thought.derived.message.observation", emit: ["stream.thought.derived.message.observation"], acceptedPrivacy: ["sensitive"], compiledEventTypes: ["stream.thought.source.telegram.message"], eventTypes: ["stream.thought.source.telegram.message"], tools: [], }); const input = fixtureRunInput(declaration); input.context = { systemText: "CURRENT TRUSTED RUNTIME", messages: [ { role: "user", content: "Please remember concise replies." }, { role: "assistant", content: "", toolCalls: [{ id: "history_0123456789abcdef", name: "request_memory_change", arguments: { operation: "append", proposed_text: "Use concise replies.", reason: "Explicit user preference.", evidence_event_ids: ["evt-evidence"], }, }], }, { role: "toolResult", toolCallId: "history_0123456789abcdef", toolName: "request_memory_change", content: "Tool executed successfully. A durable proposal was created for trusted human review. It has not been approved or applied.", isError: false, }, { role: "assistant", content: "I saved that as a memory suggestion." }, ], text: "Did that proposal actually exist?", manifest: { contextStrategy: "telegram-conversation" }, };
await fixtureRunner(server).run(input, async () => undefined);
const messages = receivedBody.messages ?? []; expect(messages.map((message) => message.role)).toEqual([ "system", "user", "assistant", "tool", "assistant", "user", ]); const call = messages[2]?.tool_calls?.[0]; expect(call).toMatchObject({ id: "history_0123456789abcdef", function: { name: "request_memory_change" }, }); expect(JSON.parse(call!.function.arguments)).toEqual({ operation: "append", proposed_text: "Use concise replies.", reason: "Explicit user preference.", evidence_event_ids: ["evt-evidence"], }); expect(messages[3]).toMatchObject({ role: "tool", tool_call_id: "history_0123456789abcdef", content: "Tool executed successfully. A durable proposal was created for trusted human review. It has not been approved or applied.", }); expect(messageText(messages[4]?.content)).toBe("I saved that as a memory suggestion."); expect(messageText(messages[5]?.content)).toBe("Did that proposal actually exist?"); expect(messageText(messages[5]?.content)).not.toContain("## Required final answer"); }, 15_000);
test("rejects oversized conversation text without retaining it", async () => { const oversized = `PRIVATE_CONVERSATION_TEXT_${"x".repeat(4_096)}`; const server = await startServer((_request, response) => { response.writeHead(200, { "content-type": "text/event-stream" }); respondWithText(response, oversized, false); }); process.env.THOUGHTSTREAM_TEST_API_KEY = "fixture-secret"; const declaration = fixtureDeclaration({ outputMode: "conversation-text", contextStrategy: "telegram-conversation", outputEventType: "stream.thought.derived.message.observation", emit: ["stream.thought.derived.message.observation"], tools: [], });
try { await fixtureRunner(server).run(fixtureRunInput(declaration), async () => undefined); throw new Error("Expected oversized conversation text to fail"); } catch (error) { expect(error).toMatchObject({ message: "Pi final output rejected: expected one bounded conversation text part", diagnostic: expect.objectContaining({ reason: "final-text-too-large", outputMode: "conversation-text", textChars: oversized.length, }), }); expect(JSON.stringify(error)).not.toContain("PRIVATE_CONVERSATION_TEXT"); } }, 15_000);
test("rejects prose, fences, multiple objects, truncation, and schema drift without retaining model text", async () => { const valid = JSON.stringify(validOutput("A valid-looking object")); const cases = [ { label: "surrounding prose", text: `Here is the result:\n${valid}`, reason: "invalid-json" }, { label: "markdown fence", text: `\`\`\`json\n${valid}\n\`\`\``, reason: "invalid-json" }, { label: "multiple objects", text: `${valid}\n${valid}`, reason: "invalid-json" }, { label: "truncated second attempt", text: `${valid}\n{\"summary\":\"unfinished`, reason: "invalid-json" }, { label: "schema drift", text: valid.slice(0, -1) + ',"unexpected":"field"}', reason: "output-contract-invalid" }, ]; let responseIndex = 0; const server = await startServer((_request, response) => { const current = cases[responseIndex++]!; response.writeHead(200, { "content-type": "text/event-stream", "x-tinker-checkpoint": "checkpoint-strict-test", }); respondWithText(response, current.text, false); }); process.env.THOUGHTSTREAM_TEST_API_KEY = "fixture-secret"; const declaration = fixtureDeclaration();
for (const current of cases) { try { await fixtureRunner(server).run(fixtureRunInput(declaration), async () => undefined); throw new Error(`Expected ${current.label} to fail`); } catch (error) { expect(error).toMatchObject({ message: "Pi final output rejected: expected one strict JSON object", diagnostic: expect.objectContaining({ code: "invalid-final-output", stage: "final-output-validation", reason: current.reason, textParts: 1, textChars: current.text.length, thinkingRedacted: true, provider: "openai-compatible", providerProfile: "fixture", model: "fixture-model", checkpointRevision: "checkpoint-strict-test", }), }); const serialized = JSON.stringify(error); expect(serialized).not.toContain(current.text); expect(serialized).not.toContain("A valid-looking object"); expect(serialized).not.toContain("attemptedTranscript"); } } }, 30_000);
test("fails closed without advancing progress when the worker artifact is unavailable", async () => { const server = await startServer((_request, response) => respondWithOutput(response, validOutput("unused"))); process.env.THOUGHTSTREAM_TEST_API_KEY = "fixture-secret"; await expect(fixtureRunner(server, { workerBundlePath: "/missing/thoughtstream-worker.cjs", }).run(fixtureRunInput(fixtureDeclaration()), async () => undefined)).rejects.toMatchObject({ advanceProgress: false, diagnostic: expect.objectContaining({ code: "sandbox-unavailable", stage: "sandbox-execution" }), }); });
test("captures both fixed native proposal tools in one request and synthesizes a tool-only acknowledgment", async () => { let requestCount = 0; let requestBody: Record<string, unknown> = {}; const server = await startServer((request, response) => { requestCount += 1; let body = ""; request.on("data", (part) => { body += part; }); request.on("end", () => { requestBody = JSON.parse(body) as Record<string, unknown>; respondWithToolCalls(response, [ { id: "call-memory", name: "request_memory_change", arguments: { operation: "append", proposed_text: "## Preference\n\nUse compact answers.", reason: "The user stated a durable response preference.", evidence_event_ids: ["evt-evidence"], }, }, { id: "call-correction", name: "submit_correction", arguments: { target_output: "evt-prior-output", replacement: "The corrected concise reply.", reason: "The prior reply misstated the fact.", evidence_event_ids: ["evt-evidence"], }, }, ]); }); }); process.env.THOUGHTSTREAM_TEST_API_KEY = "fixture-secret"; const input = proposalRunInput(); const traces: Array<{ kind: string; data: unknown }> = [];
const output = await fixtureRunner(server).run(input, async (trace) => { traces.push(trace); });
expect(requestCount).toBe(1); expect(requestBody).toMatchObject({ tool_choice: "auto", tools: [ { type: "function", function: { name: "request_memory_change" } }, { type: "function", function: { name: "submit_correction" } }, ], }); expect(requestBody).not.toHaveProperty("response_format"); expect(JSON.stringify(requestBody)).toContain("latest"); expect(JSON.stringify(requestBody)).toContain("private inert suggestions for review"); expect(JSON.stringify(requestBody)).not.toContain("most recent delivered output: evt-prior-output"); expect(JSON.stringify(requestBody)).toContain("The current user turn includes no image input"); expect(JSON.stringify(requestBody)).not.toContain("correction_target_output"); expect(output.summary).toBe("I saved those as memory and correction suggestions."); expect(output.proposals).toEqual([ expect.objectContaining({ toolCallId: "call-memory", kind: "memory-change" }), expect.objectContaining({ toolCallId: "call-correction", kind: "self-correction" }), ]); expect(JSON.stringify(traces)).not.toContain("Use compact answers"); expect(JSON.stringify(traces)).not.toContain("corrected concise reply"); }, 15_000);
test("captures a complete /focus declaration as one inert native proposal", async () => { let requestBody: Record<string, unknown> = {}; const declaration = focusDeclarationFixture(); const server = await startServer((request, response) => { let body = ""; request.on("data", (part) => { body += part; }); request.on("end", () => { requestBody = JSON.parse(body) as Record<string, unknown>; respondWithToolCalls(response, [{ id: "call-focus", name: "propose_focus", arguments: { declaration }, }]); }); }); process.env.THOUGHTSTREAM_TEST_API_KEY = "fixture-secret";
const output = await fixtureRunner(server).run(proposalRunInput(true), async () => undefined);
expect(requestBody).toMatchObject({ tool_choice: "auto", tools: expect.arrayContaining([ expect.objectContaining({ type: "function", function: expect.objectContaining({ name: "propose_focus" }), }), ]), }); expect(JSON.stringify(requestBody)).toContain("validated `/focus` command"); expect(output.summary).toBe("I saved that as a focus proposal."); expect(output.proposals).toEqual([{ toolCallId: "call-focus", kind: "focus-declaration", arguments: { declaration }, }]); }, 15_000);
test("fails a description-bearing /focus turn closed when the model returns text without a focus proposal", async () => { let requestCount = 0; const server = await startServer((_request, response) => { requestCount += 1; respondWithText(response, "I can help with that.", false); }); process.env.THOUGHTSTREAM_TEST_API_KEY = "fixture-secret"; const input = proposalRunInput(true); input.event.type = "stream.thought.source.telegram.message"; input.event.source = "telegram:fixture"; input.event.sourceKind = "telegram"; input.event.payload = { text: "/focus AI news", accountUsername: "TheStreamBot", chatId: "123456789" }; input.context.text = "/focus AI news";
await expect(fixtureRunner(server).run(input, async () => undefined)).rejects.toMatchObject({ advanceProgress: true, diagnostic: expect.objectContaining({ reason: "focus-proposal-required" }), }); expect(requestCount).toBe(1); }, 15_000);
test("answers exact Telegram help deterministically without provider, sandbox, or proposal execution", async () => { let requestCount = 0; const server = await startServer((_request, response) => { requestCount += 1; respondWithText(response, "provider must not be called", false); }); process.env.THOUGHTSTREAM_TEST_API_KEY = "fixture-secret"; const input = proposalRunInput(); input.event.type = "stream.thought.source.telegram.message"; input.event.source = "telegram:fixture"; input.event.sourceKind = "telegram"; input.event.payload = { text: "/help", accountUsername: "CameronStreamBot", chatId: "123456789", }; input.context.text = "/help"; const traces: Array<{ kind: string; data: unknown }> = [];
const output = await fixtureRunner(server).run(input, async (trace) => { traces.push(trace); });
expect(requestCount).toBe(0); expect(output.summary).toContain("/correct <replacement>"); expect(output.summary).toContain("https://thought.stream/inspector/"); expect(output.model).toEqual({ provider: "trusted-parent", id: "stream-help-v3" }); expect(traces).toEqual([{ kind: "telegram.command.completed", data: { command: "help", version: "stream-help-v3", providerRequests: 0 }, }]); });
test("answers bare /focus deterministically with usage and no provider request", async () => { let requestCount = 0; const server = await startServer((_request, response) => { requestCount += 1; respondWithText(response, "provider must not be called", false); }); process.env.THOUGHTSTREAM_TEST_API_KEY = "fixture-secret"; const input = proposalRunInput(true); input.event.type = "stream.thought.source.telegram.message"; input.event.source = "telegram:fixture"; input.event.sourceKind = "telegram"; input.event.payload = { text: "/focus", accountUsername: "CameronStreamBot", chatId: "123456789", }; input.context.text = "/focus"; const traces: Array<{ kind: string; data: unknown }> = [];
const output = await fixtureRunner(server).run(input, async (trace) => { traces.push(trace); });
expect(requestCount).toBe(0); expect(output.summary).toContain("/focus <description>"); expect(output.summary).toContain("private inert proposal"); expect(output.model).toEqual({ provider: "trusted-parent", id: "stream-focus-help-v1" }); expect(traces).toEqual([{ kind: "telegram.command.completed", data: { command: "focus", version: "stream-focus-help-v1", providerRequests: 0 }, }]); });
test("preserves visible text alongside one validated native correction suggestion", async () => { const server = await startServer((_request, response) => respondWithToolCalls(response, [{ id: "call-correction-text", name: "submit_correction", arguments: { target_output: "evt-prior-output", replacement: "The corrected concise reply.", reason: "The prior reply misstated the fact.", evidence_event_ids: ["evt-evidence"], }, }], "You're right about that.")); process.env.THOUGHTSTREAM_TEST_API_KEY = "fixture-secret";
const output = await fixtureRunner(server).run(proposalRunInput(), async () => undefined);
expect(output.summary).toBe("You're right about that."); expect(output.proposals).toHaveLength(1); expect(output.proposals?.[0]).toMatchObject({ kind: "self-correction", toolCallId: "call-correction-text" }); }, 15_000);
test("accepts the conversational latest shorthand for the most recent correction target", async () => { const server = await startServer((_request, response) => respondWithToolCalls(response, [{ id: "call-correction-latest", name: "submit_correction", arguments: { target_output: "latest", replacement: "The corrected concise reply.", reason: "The prior reply misstated the fact.", evidence_event_ids: ["evt-evidence"], }, }], "I can correct that.")); process.env.THOUGHTSTREAM_TEST_API_KEY = "fixture-secret";
const output = await fixtureRunner(server).run(proposalRunInput(), async () => undefined);
expect(output.summary).toBe("I can correct that."); expect(output.proposals?.[0]).toMatchObject({ kind: "self-correction", arguments: { target_output: "latest" }, }); }, 15_000);
test("fails closed on unknown, duplicate, and oversized proposal calls without a second provider request", async () => { const cases = [ [{ id: "unknown", name: "unknown_mutation", arguments: {} }], [ { id: "first", name: "request_memory_change", arguments: { operation: "append", proposed_text: "One", reason: "One", evidence_event_ids: ["evt-evidence"] }, }, { id: "second", name: "request_memory_change", arguments: { operation: "append", proposed_text: "Two", reason: "Two", evidence_event_ids: ["evt-evidence"] }, }, ], [{ id: "oversized", name: "request_memory_change", arguments: { operation: "append", proposed_text: "x".repeat(32_769), reason: "Too large", evidence_event_ids: ["evt-evidence"] }, }], ]; for (const calls of cases) { let requestCount = 0; const server = await startServer((_request, response) => { requestCount += 1; respondWithToolCalls(response, calls); }); process.env.THOUGHTSTREAM_TEST_API_KEY = "fixture-secret"; await expect(fixtureRunner(server).run(proposalRunInput(), async () => undefined)).rejects.toMatchObject({ advanceProgress: false, }); expect(requestCount).toBe(1); } }, 20_000);
test("retains a content-dark classification when a correction targets evidence outside its snapshot", async () => { let requestCount = 0; const server = await startServer((_request, response) => { requestCount += 1; respondWithToolCalls(response, [{ id: "outside-target", name: "submit_correction", arguments: { target_output: "evt-not-admitted", replacement: "A private replacement that must not enter diagnostics.", reason: "A private reason that must not enter diagnostics.", evidence_event_ids: ["evt-evidence"], }, }]); }); process.env.THOUGHTSTREAM_TEST_API_KEY = "fixture-secret"; const traces: Array<{ kind: string; data: unknown }> = [];
try { await fixtureRunner(server).run(proposalRunInput(), async (trace) => { traces.push(trace); }); throw new Error("Expected proposal capability failure"); } catch (error) { expect(error).toMatchObject({ advanceProgress: false, diagnostic: expect.objectContaining({ code: "sandbox-worker-model-failed", stage: "sandbox-execution", workerReason: "proposal-tool-error", proposalFailureCode: "target-outside-snapshot", }), }); const durable = JSON.stringify({ error, traces }); expect(durable).not.toContain("private replacement"); expect(durable).not.toContain("private reason"); } expect(requestCount).toBe(1); expect(traces.map((trace) => trace.kind)).toContain("pi.tool_execution_end"); }, 15_000);
test("classifies provider timeout without exposing upstream or credential text", async () => { const server = await startServer(() => undefined); process.env.THOUGHTSTREAM_TEST_API_KEY = "fixture-secret"; const declaration = fixtureDeclaration({ timeoutMs: 50 }); const failure = fixtureRunner(server).run(fixtureRunInput(declaration), async () => undefined); await expect(failure).rejects.toMatchObject({ advanceProgress: false, diagnostic: expect.objectContaining({ code: "timeout", stage: "provider" }), }); await expect(failure).rejects.not.toThrow(/fixture-secret/); }, 10_000);
test("classifies deterministic provider 4xx rejection as terminal without retaining its body", async () => { const privateBody = "PRIVATE_UPSTREAM_VALIDATION_DETAIL"; const server = await startServer((_request, response) => { response.writeHead(400, { "content-type": "application/json" }); response.end(JSON.stringify({ error: { message: privateBody } })); }); process.env.THOUGHTSTREAM_TEST_API_KEY = "fixture-secret";
try { await fixtureRunner(server).run(fixtureRunInput(fixtureDeclaration()), async () => undefined); throw new Error("Expected provider rejection"); } catch (error) { expect(error).toMatchObject({ advanceProgress: true, diagnostic: expect.objectContaining({ code: "provider-client-error", stage: "provider" }), }); expect(JSON.stringify(error)).not.toContain(privateBody); } }, 15_000);});
function fixtureRunner( server: http.Server, options: { fetchImpl?: typeof fetch; workerBundlePath?: string; allowedModels?: string[]; profileId?: string; provider?: "tinker" | "openai-compatible"; providerFetchImpl?: typeof fetch; } = {},): PiAgentRunner { return new PiAgentRunner({ providerProfiles: staticProviderProfileResolver([fixtureProfile( server, options.allowedModels, options.profileId, options.provider, )]), workerBundlePath: options.workerBundlePath ?? workerBundlePath, ...(options.fetchImpl ? { fetchImpl: options.fetchImpl } : {}), // These tool responses are synthetic; no public DNS lookup is needed. ...(options.fetchImpl ? { resolveHostname: async () => ["93.184.216.34"] } : {}), ...(options.providerFetchImpl ? { providerFetchImpl: options.providerFetchImpl } : {}), });}
function fixtureProfile( server: http.Server, allowedModels: string[] = ["fixture-model"], id = "fixture", provider: "tinker" | "openai-compatible" = "openai-compatible",): ProviderProfile { return { id, provider, baseUrl: serverBaseUrl(server), route: "/chat/completions", apiKeyEnv: "THOUGHTSTREAM_TEST_API_KEY", allowedModels: new Set(allowedModels), imageInputModels: new Set(), jsonObjectResponseFormat: true, jsonSchemaResponseFormat: false, requestTimeoutMs: 5_000, maxRequestBytes: 2 * 1024 * 1024, maxResponseBytes: 1_000_000, };}
async function compiledAdapterDeclaration( privateCheckpoint: string, options: { atprotoPrefetch?: boolean } = {},): Promise<ThoughtAgentDeclaration> { const root = await fs.mkdtemp(path.join(os.tmpdir(), "thoughtstream-pi-adapter-")); temporaryRoots.push(root); await fs.mkdir(path.join(root, "adapters", "releases"), { recursive: true }); await fs.mkdir(path.join(root, "agents"), { recursive: true }); await fs.mkdir(path.join(root, "prompts"), { recursive: true }); const releaseManifest = [ "schemaVersion: 1", "id: fixture-adapter", "version: 1", "description: Fixture learned adapter", "releasedAt: 2026-07-25T20:00:00.000Z", "providerProfile: tinker-default", "baseModel: Qwen/Qwen3.5-4B", "checkpoint: { env: THOUGHTSTREAM_TEST_ADAPTER_MODEL }", `dataset: { id: fixture-dataset, sha256: ${"c".repeat(64)} }`, `evals: [{ id: fixture-eval, sha256: ${"d".repeat(64)} }]`, "capabilities: [fixture-task]", "privacyClass: private", "exportClass: forbidden", ].join("\n"); await fs.writeFile(path.join(root, "adapters", "releases", "fixture.yaml"), `${releaseManifest}\n`); const manifestSha256 = createHash("sha256").update(canonicalJson(YAML.parse(releaseManifest))).digest("hex"); await fs.writeFile(path.join(root, "adapters", "deployment.yaml"), [ "schemaVersion: 1", "generation: 1", "selections:", " - declaration: { id: pi-test, version: 1 }", ` release: { id: fixture-adapter, version: 1, manifestSha256: ${manifestSha256} }`, " state: active", "", ].join("\n")); await fs.writeFile(path.join(root, "prompts", "fixture.md"), "Read the fixture.\n"); await fs.writeFile(path.join(root, "agents", "fixture.yaml"), [ "id: pi-test", "version: 1", "description: Exercises the Pi adapter", options.atprotoPrefetch ? "subscribe: { types: [stream.thought.source.atproto.commit], sources: ['*'], privacy: [private] }" : "subscribe: { types: [stream.thought.source.file.changed], sources: ['*'], privacy: [private] }", "context: { maxEvents: 1, maxChars: 10000, strategy: single-event }", "runner:", " kind: pi", " profile: tinker-default", " adapter: { id: fixture-adapter, version: 1 }", " maxOutputTokens: 500", " timeoutMs: 5000", "accounting:", " leaseMs: 10000", " reservation: { inputTokens: 5000, outputTokens: 500 }", " limits: [{ window: hour, maxCalls: 10, maxInputTokens: 50000, maxOutputTokens: 5000 }]", "prompt: prompts/fixture.md", "emit: [stream.thought.derived.document.read]", options.atprotoPrefetch ? "policy: { tools: [atproto.fetch-markdown], externalActions: false }" : "policy: { tools: [], externalActions: false }", "enabled: true", "", ].join("\n")); process.env.THOUGHTSTREAM_TINKER_ALLOWED_MODELS = `Qwen/Qwen3.5-4B,${privateCheckpoint}`; const [declaration] = await loadAgentDeclarations(path.join(root, "agents"), { THOUGHTSTREAM_TEST_ADAPTER_MODEL: privateCheckpoint, THOUGHTSTREAM_TINKER_ALLOWED_MODELS: process.env.THOUGHTSTREAM_TINKER_ALLOWED_MODELS, }); if (!declaration) throw new Error("Missing compiled adapter declaration"); return declaration;}
function fixtureRunInput(declaration: ThoughtAgentDeclaration) { const event: ThoughtEvent = { id: "evt_failure_test", sourceSequence: 1, type: "stream.thought.source.file.changed", schemaVersion: 1, source: "filesystem:test", sourceKind: "filesystem", externalId: "doc_test", idempotencyKey: "failure-test", occurredAt: "2026-07-14T00:00:00.000Z", observedAt: "2026-07-14T00:00:00.000Z", actor: "filesystem:test", rootEventId: "evt_failure_test", correlationId: "scan_test", privacy: "private", payload: { path: "fixture.md" }, payloadHash: "hash", createdByRuntime: "test", }; return { runId: "run_failure_test", declaration, event, context: buildContextPacket(declaration, event) };}
function proposalRunInput(includeFocus = false) { const proposals: ThoughtAgentDeclaration["proposals"] = [ "memory-change", "self-correction", ...(includeFocus ? ["focus-declaration" as const] : []), ]; const declaration = fixtureDeclaration({ acceptedPrivacy: ["sensitive"], outputEventType: "stream.thought.derived.message.observation", emit: ["stream.thought.derived.message.observation"], outputMode: "conversation-text", contextStrategy: "telegram-conversation", contextDocumentSubscriptions: [ { source: "filesystem:telegram-agent-context", paths: ["identity.md", "memory.md"], required: true }, ], contextDocumentMaxChars: 64_000, proposals, }); const input = fixtureRunInput(declaration); input.context = { text: "A bounded private conversation transcript.", manifest: { contextSnapshot: { id: "snapshot-proposal-test" }, proposalCapabilities: { enabled: proposals, evidenceEventIds: ["evt-evidence"], memoryTarget: { source: "filesystem:telegram-agent-context", documentId: "doc-memory", path: "memory.md", versionId: "version-memory", sha256: "a".repeat(64), contentType: "text/markdown", }, correctionTargets: [{ runId: "run-prior", outputEventId: "evt-prior-output", deliveryReceiptEventId: "evt-prior-delivery", sourceRootEventId: "evt-prior-root", outputContract: { id: "observation", version: 1, sha256: "b".repeat(64) }, }], }, }, }; return input;}
function focusDeclarationFixture() { return { id: "ai-news", version: 1, name: "AI news", description: "Notice public AI news that changes the user's model of agents and training.", parentFocusIds: ["news"], scope: { summary: "Public AI research, product, and infrastructure developments.", includeTags: ["ai-news"], excludeTags: [], }, objective: { statement: "Produce concise evidence-grounded observations about meaningful AI developments." }, subscriptions: { eventTypes: ["stream.thought.source.rss.item", "stream.thought.source.atproto.commit"], sourcePatterns: ["rss:*", "jetstream:*"], privacy: ["public-source" as const], replay: "now" as const, }, budgets: { inference: { maxCallsPerDay: 24, maxInputTokensPerDay: 480_000, maxOutputTokensPerDay: 12_288, maxCostMicrousdPerDay: 2_400_000, }, attention: { maxDeliveriesPerDay: 6 }, }, permissions: { requestedReadTools: [], requestedExternalActions: [] }, retirementRule: "Retire or revise after 100 eligible events without accepted value.", };}
function fixtureDeclaration(overrides: Partial<ThoughtAgentDeclaration> = {}): ThoughtAgentDeclaration { return { id: "pi-test", version: 1, name: "Pi test", description: "Exercises the Pi adapter", mode: "pi", provider: "openai-compatible", providerProfile: "fixture", model: "fixture-model", eventTypes: ["*"], compiledEventTypes: ["stream.thought.source.file.changed"], sourcePatterns: ["*"], acceptedPrivacy: ["private"], outputEventType: "stream.thought.derived.document.read", emit: ["stream.thought.derived.document.read"], promptRef: "fixture", systemPrompt: "Read the fixture.", enabled: true, maxEvents: 1, maxInputChars: 10_000, maxOutputTokens: 500, timeoutMs: 5_000, tools: [], externalActions: false, ...overrides, };}
async function startServer(handler: http.RequestListener): Promise<http.Server> { const server = http.createServer(handler); servers.push(server); await new Promise<void>((resolve) => server.listen(0, "127.0.0.1", resolve)); return server;}
function serverBaseUrl(server: http.Server): string { const address = server.address(); if (!address || typeof address === "string") throw new Error("Missing test server address"); return `http://127.0.0.1:${address.port}/v1`;}
function validOutput(summary: string): Record<string, unknown> { return { summary, tags: ["fixture"], importance: "normal", confidence: 0.8 };}
function messageText(content: unknown): string { if (typeof content === "string") return content; if (!Array.isArray(content)) return ""; return content.map((part) => ( part && typeof part === "object" && typeof (part as { text?: unknown }).text === "string" ? String((part as { text: string }).text) : "" )).join("");}
function respondWithOutput( response: http.ServerResponse, output: Record<string, unknown>, usage?: { promptTokens: number; completionTokens: number },): void { response.writeHead(200, { "content-type": "text/event-stream" }); response.write(`data: ${JSON.stringify(chunk({ role: "assistant", content: JSON.stringify(output) }, null))}\n\n`); response.write(`data: ${JSON.stringify(chunk({}, "stop"))}\n\n`); if (usage) { response.write(`data: ${JSON.stringify({ id: "chatcmpl-fixture", object: "chat.completion.chunk", created: 1, model: "fixture-model", choices: [], usage: { prompt_tokens: usage.promptTokens, completion_tokens: usage.completionTokens, total_tokens: usage.promptTokens + usage.completionTokens, }, })}\n\n`); } response.end("data: [DONE]\n\n");}
function respondWithToolCalls( response: http.ServerResponse, calls: Array<{ id: string; name: string; arguments: Record<string, unknown> }>, text?: string,): void { response.writeHead(200, { "content-type": "text/event-stream" }); response.write(`data: ${JSON.stringify(chunk({ role: "assistant", ...(text ? { content: text } : {}), tool_calls: calls.map((call, index) => ({ index, id: call.id, type: "function", function: { name: call.name, arguments: JSON.stringify(call.arguments) }, })), }, null))}\n\n`); response.write(`data: ${JSON.stringify(chunk({}, "tool_calls"))}\n\n`); response.end("data: [DONE]\n\n");}
function respondWithText(response: http.ServerResponse, text: string, writeHead = true): void { if (writeHead) response.writeHead(200, { "content-type": "text/event-stream" }); response.write(`data: ${JSON.stringify(chunk({ role: "assistant", content: text }, null))}\n\n`); response.write(`data: ${JSON.stringify(chunk({}, "stop"))}\n\n`); response.end("data: [DONE]\n\n");}
function chunk( delta: Record<string, unknown>, finishReason: string | null, model = "fixture-model",): Record<string, unknown> { return { id: "chatcmpl-fixture", object: "chat.completion.chunk", created: 1, model, choices: [{ index: 0, delta, finish_reason: finishReason }], };}