Something went wrong. Try again.
A local-first event pipeline for independent agents, built on Jazz.
Something went wrong. Try again.
11 kB · 263 lines
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264import fs from "node:fs/promises";import path from "node:path";import { afterEach, describe, expect, test } from "vitest";import { createRunTools, fetchAtprotoMarkdownUriDocument, fetchBskyMarkdownDocument,} from "../src/agents/tools.js";import type { ThoughtEvent } from "../src/events/types.js";import { temporaryProject } from "./helpers.js";
const roots: string[] = [];
afterEach(async () => { await Promise.all(roots.splice(0).map((root) => fs.rm(root, { recursive: true, force: true })));});
describe("Pi enrichment tools", () => { test("fetches one fixed-host bsky.md social view from a feed-post AT URI", async () => { const requests: Array<{ url: string; redirect: RequestRedirect | undefined }> = []; const document = await fetchBskyMarkdownDocument({ atUri: "at://did:plc:alice/app.bsky.feed.post/post-one", fetchImpl: async (input, init) => { requests.push({ url: input.toString(), redirect: init?.redirect }); return new Response("# Social post\n", { headers: { "content-type": "text/markdown; charset=utf-8" }, }); }, resolveHostname: async (hostname) => { expect(hostname).toBe("bsky-md.noz.am"); return ["8.8.8.8"]; }, });
expect(requests).toEqual([{ url: "https://bsky-md.noz.am/profile/did%3Aplc%3Aalice/post/post-one", redirect: "error", }]); expect(document.markdown).toBe("# Social post\n"); expect(document.details).toMatchObject({ atUri: "at://did:plc:alice/app.bsky.feed.post/post-one", endpoint: "https://bsky-md.noz.am/profile/did%3Aplc%3Aalice/post/post-one", mediaType: "text/markdown", sizeBytes: 14, sha256: expect.any(String), }); });
test("fetches a non-Bluesky ATProto record without calling the Bluesky AppView", async () => { const requests: string[] = []; const document = await fetchAtprotoMarkdownUriDocument({ atUri: "at://did:plc:fixturecard/network.cosmik.card/card-fixture", fetchImpl: async (input) => { requests.push(input.toString()); return new Response("SYNTHETIC CARD MARKDOWN", { headers: { "content-type": "text/markdown" }, }); }, resolveHostname: async (hostname) => { expect(hostname).toBe("atproto.md"); return ["8.8.8.8"]; }, });
expect(requests).toEqual([ "https://atproto.md/at://did:plc:fixturecard/network.cosmik.card/card-fixture", ]); expect(document.markdown).toBe("SYNTHETIC CARD MARKDOWN"); expect(document.details).toMatchObject({ atUri: "at://did:plc:fixturecard/network.cosmik.card/card-fixture", imageResolution: { status: "unavailable", imageUrls: [], }, }); });
test("bounds Markdown DNS preflight with the same abort signal as the fetch", async () => { let fetchCalls = 0; await expect(fetchBskyMarkdownDocument({ atUri: "at://did:plc:alice/app.bsky.feed.post/post-one", fetchImpl: async () => { fetchCalls += 1; return new Response("must not fetch"); }, resolveHostname: async () => new Promise<string[]>(() => {}), signal: AbortSignal.timeout(10), })).rejects.toThrow("resolution aborted"); expect(fetchCalls).toBe(0); });
test("chains ATProto Markdown into a bounded image download and durable content-addressed artifact", async () => { const root = await temporaryProject(); roots.push(root); const image = Buffer.from([0x89, 0x50, 0x4e, 0x47, 0x0d, 0x0a]); const requests: string[] = []; const traces: Array<{ kind: string; data: unknown }> = []; const fetchImpl = async (input: string | URL | Request): Promise<Response> => { const url = input.toString(); requests.push(url); if (url === "https://atproto.md/at://did:plc:alice/app.bsky.feed.post/post-one") { return new Response("# Post\n\n**Images:** image", { headers: { "content-type": "text/markdown" }, }); } if (url.startsWith("https://public.api.bsky.app/xrpc/app.bsky.feed.getPosts?")) { return Response.json({ posts: [{ embed: { images: [{ fullsize: "https://cdn.example/post.png" }] } }] }); } if (url === "https://cdn.example/post.png") { return new Response(image, { headers: { "content-type": "image/png", "content-length": String(image.byteLength) }, }); } return new Response("missing", { status: 404 }); }; const toolSet = createRunTools({ event: atprotoEvent(), names: ["atproto.fetch-markdown", "web.download-image"], artifactRoot: path.join(root, "artifacts"), fetchImpl, resolveHostname: async () => ["8.8.8.8"], onTrace: async (trace) => { traces.push(trace); }, });
const markdown = toolSet.tools.find((tool) => tool.name === "fetch_atproto_markdown"); const imageTool = toolSet.tools.find((tool) => tool.name === "download_image"); if (!markdown || !imageTool) throw new Error("Missing fixture tools"); const markdownResult = await markdown.execute("call_markdown", { target: "event" }); const imageResult = await imageTool.execute("call_image", { url: "https://cdn.example/post.png" });
expect(markdownResult.content[0]).toMatchObject({ type: "text", text: expect.stringContaining("https://cdn.example/post.png") }); expect(imageResult.content).toEqual(expect.arrayContaining([ expect.objectContaining({ type: "image", mimeType: "image/png", data: image.toString("base64") }), ])); expect(toolSet.outcomes).toMatchObject([ { tool: "atproto.fetch-markdown", status: "succeeded", requestKeys: ["target"], argumentsRedacted: true, resultPresent: true, resultSha256: expect.any(String), }, { tool: "web.download-image", status: "succeeded", requestKeys: ["url"], argumentsRedacted: true, resultPresent: true, resultSha256: expect.any(String), }, ]); expect(JSON.stringify(toolSet.outcomes)).not.toContain("https://cdn.example/post.png"); const artifactPath = (imageResult.details as { artifactPath?: unknown }).artifactPath; expect(typeof artifactPath).toBe("string"); expect(await fs.readFile(path.join(root, "artifacts", String(artifactPath)))).toEqual(image); expect(requests).toEqual([ "https://atproto.md/at://did:plc:alice/app.bsky.feed.post/post-one", expect.stringMatching(/^https:\/\/public\.api\.bsky\.app\/xrpc\/app\.bsky\.feed\.getPosts\?/), "https://cdn.example/post.png", ]); expect(traces.map((trace) => trace.kind)).toEqual(["enrichment.succeeded", "enrichment.succeeded"]); const persistedTraceShape = JSON.stringify(traces); expect(persistedTraceShape).not.toContain("https://cdn.example/post.png"); expect(persistedTraceShape).not.toContain(image.toString("base64")); expect(persistedTraceShape).toContain("argumentsRedacted"); });
test("resolves a like subject and rejects image URLs not discovered in source evidence", async () => { const traces: string[] = []; const toolSet = createRunTools({ event: atprotoEvent({ collection: "app.bsky.feed.like", record: { subject: { uri: "at://did:plc:bob/app.bsky.feed.post/liked-one", cid: "bafy-liked" } }, }), names: ["atproto.fetch-markdown", "web.download-image"], fetchImpl: async (input) => new Response(`Fetched ${input.toString()}`, { headers: { "content-type": "text/markdown" }, }), resolveHostname: async () => ["8.8.8.8"], onTrace: async (trace) => { traces.push(trace.kind); }, }); const markdown = toolSet.tools.find((tool) => tool.name === "fetch_atproto_markdown"); const imageTool = toolSet.tools.find((tool) => tool.name === "download_image"); if (!markdown || !imageTool) throw new Error("Missing fixture tools");
const result = await markdown.execute("call_subject", { target: "subject" }); expect(result.content[0]).toMatchObject({ type: "text", text: "Fetched https://atproto.md/at://did:plc:bob/app.bsky.feed.post/liked-one", }); await expect(imageTool.execute("call_bad_image", { url: "https://localhost/private.png" })) .rejects.toThrow("not discovered"); expect(toolSet.outcomes.at(-1)).toMatchObject({ tool: "web.download-image", status: "failed", requestKeys: ["url"], argumentsRedacted: true, resultPresent: false, errorCode: "tool-execution-failed", }); expect(JSON.stringify(toolSet.outcomes)).not.toContain("localhost"); expect(JSON.stringify(toolSet.outcomes)).not.toContain("not discovered"); expect(traces).toEqual(["enrichment.succeeded", "enrichment.failed"]); });
test("blocks discovered image URLs whose host resolves to a non-public address", async () => { let imageFetches = 0; const toolSet = createRunTools({ event: atprotoEvent(), names: ["atproto.fetch-markdown", "web.download-image"], fetchImpl: async (input) => { if (input.toString().startsWith("https://atproto.md/")) { return new Response("", { headers: { "content-type": "text/markdown" }, }); } if (input.toString().startsWith("https://public.api.bsky.app/")) { return Response.json({ posts: [] }); } imageFetches += 1; return new Response(Buffer.from("private"), { headers: { "content-type": "image/png" } }); }, resolveHostname: async (hostname) => hostname === "images.example" ? ["127.0.0.1"] : ["8.8.8.8"], onTrace: async () => undefined, }); const markdown = toolSet.tools.find((tool) => tool.name === "fetch_atproto_markdown"); const imageTool = toolSet.tools.find((tool) => tool.name === "download_image"); if (!markdown || !imageTool) throw new Error("Missing fixture tools"); await markdown.execute("call_markdown", { target: "event" }); await expect(imageTool.execute("call_private", { url: "https://images.example/private.png" })) .rejects.toThrow("does not resolve exclusively to public addresses"); expect(imageFetches).toBe(0); });});
function atprotoEvent(payload: Record<string, unknown> = {}): ThoughtEvent { return { id: "evt_atproto_fixture", sourceSequence: 1, type: "stream.thought.source.atproto.commit", schemaVersion: 1, source: "jetstream:fixture", sourceKind: "jetstream", externalId: "did:plc:alice/app.bsky.feed.post/post-one", idempotencyKey: "atproto-fixture", occurredAt: "2026-07-14T00:00:00.000Z", observedAt: "2026-07-14T00:00:00.000Z", actor: "did:plc:alice", rootEventId: "evt_atproto_fixture", correlationId: "rev-fixture", privacy: "public-source", payload: { operation: "create", collection: "app.bsky.feed.post", atUri: "at://did:plc:alice/app.bsky.feed.post/post-one", record: { text: "fixture post" }, ...payload, }, payloadHash: "hash", createdByRuntime: "test", };}