Something went wrong. Try again.
A local-first event pipeline for independent agents, built on Jazz.
Something went wrong. Try again.
21 kB · 564 lines
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565import { createHash } from "node:crypto";import { BlockList, isIP } from "node:net";import { resolve4, resolve6 } from "node:dns/promises";import type { AgentTool } from "@earendil-works/pi-agent-core";import { Type } from "typebox";import type { JsonObject } from "../core/json.js";import type { ThoughtEvent } from "../events/types.js";import type { EnrichmentOutcome, RunnerTrace } from "./types.js";import { storeBlob } from "../artifacts/blob.js";
export const AGENT_TOOL_NAMES = ["atproto.fetch-markdown", "web.download-image"] as const;export type AgentToolName = (typeof AGENT_TOOL_NAMES)[number];
const MARKDOWN_MAX_BYTES = 256 * 1024;const BSKY_MARKDOWN_MAX_BYTES = 256 * 1024;const APPVIEW_MAX_BYTES = 512 * 1024;const IMAGE_MAX_BYTES = 8 * 1024 * 1024;const MAX_REDIRECTS = 3;const nonPublicAddresses = createNonPublicBlockList();
type FetchLike = (input: string | URL | Request, init?: RequestInit) => Promise<Response>;type ResolveHostname = (hostname: string) => Promise<string[]>;
export interface AtprotoMarkdownDocument { markdown: string; details: JsonObject;}
export interface AtprotoMarkdownFetchOptions { event: ThoughtEvent; target: "event" | "subject"; allowedImages?: Set<string> | undefined; fetchImpl?: FetchLike | undefined; resolveHostname?: ResolveHostname | undefined; signal?: AbortSignal | undefined;}
export interface AtprotoMarkdownUriFetchOptions { atUri: string; allowedImages?: Set<string> | undefined; fetchImpl?: FetchLike | undefined; resolveHostname?: ResolveHostname | undefined; signal?: AbortSignal | undefined;}
export interface BskyMarkdownDocument { markdown: string; details: JsonObject;}
export interface BskyMarkdownFetchOptions { atUri: string; fetchImpl?: FetchLike | undefined; resolveHostname?: ResolveHostname | undefined; signal?: AbortSignal | undefined;}
export interface RunToolOptions { event: ThoughtEvent; names: AgentToolName[]; artifactRoot?: string | undefined; fetchImpl?: FetchLike | undefined; resolveHostname?: ResolveHostname | undefined; onTrace: (trace: RunnerTrace) => Promise<void>;}
export interface RunToolSet { tools: AgentTool[]; outcomes: EnrichmentOutcome[];}
export function createRunTools(options: RunToolOptions): RunToolSet { const outcomes: EnrichmentOutcome[] = []; const allowedImages = new Set(extractImageUrls(options.event.payload)); const fetchImpl = options.fetchImpl ?? fetch; const resolveHostname = options.resolveHostname ?? resolvePublicAddresses; const tools = options.names.map((name) => { if (name === "atproto.fetch-markdown") { return createAtprotoMarkdownTool(options.event, allowedImages, outcomes, fetchImpl, resolveHostname, options.onTrace); } if (name === "web.download-image") { return createImageDownloadTool( allowedImages, outcomes, fetchImpl, resolveHostname, options.artifactRoot, options.onTrace, ); } return assertNever(name); }); return { tools, outcomes };}
const markdownParameters = Type.Object({ target: Type.Union([Type.Literal("event"), Type.Literal("subject")], { description: "Fetch the current AT record, or the liked/reposted subject referenced by it.", }),});
function createAtprotoMarkdownTool( event: ThoughtEvent, allowedImages: Set<string>, outcomes: EnrichmentOutcome[], fetchImpl: FetchLike, resolveHostname: ResolveHostname, onTrace: (trace: RunnerTrace) => Promise<void>,): AgentTool<typeof markdownParameters> { return { name: "fetch_atproto_markdown", label: "Fetch ATProto Markdown", description: "Fetch clean public Markdown for the current AT record or its referenced subject from atproto.md.", parameters: markdownParameters, executionMode: "sequential", execute: async (_toolCallId, params, signal) => { const request = { target: params.target }; try { const document = await fetchAtprotoMarkdownDocument({ event, target: params.target, allowedImages, fetchImpl, resolveHostname, ...(signal ? { signal } : {}), }); await recordOutcome(outcomes, onTrace, { tool: "atproto.fetch-markdown", status: "succeeded", request, result: document.details, }); return { content: [{ type: "text", text: document.markdown }], details: document.details, }; } catch (error) { await recordFailure(outcomes, onTrace, "atproto.fetch-markdown", request, error); throw error; } }, };}
export async function fetchAtprotoMarkdownDocument( options: AtprotoMarkdownFetchOptions,): Promise<AtprotoMarkdownDocument> { return fetchAtprotoMarkdownUriDocument({ atUri: resolveAtUri(options.event, options.target), allowedImages: options.allowedImages ?? new Set(extractImageUrls(options.event.payload)), ...(options.fetchImpl ? { fetchImpl: options.fetchImpl } : {}), ...(options.resolveHostname ? { resolveHostname: options.resolveHostname } : {}), ...(options.signal ? { signal: options.signal } : {}), });}
export async function fetchAtprotoMarkdownUriDocument( options: AtprotoMarkdownUriFetchOptions,): Promise<AtprotoMarkdownDocument> { if (!isAtUri(options.atUri)) throw new Error("ATProto Markdown requires a canonical AT URI"); const allowedImages = options.allowedImages ?? new Set<string>(); const fetchImpl = options.fetchImpl ?? fetch; const resolveHostname = options.resolveHostname ?? resolvePublicAddresses; const endpoint = new URL(`https://atproto.md/${options.atUri}`); await assertPublicHost(endpoint, resolveHostname, "ATProto Markdown endpoint", options.signal); const response = await fetchImpl(endpoint, { headers: { accept: "text/markdown" }, redirect: "error", ...(options.signal ? { signal: options.signal } : {}), }); if (!response.ok) throw new Error(`atproto.md returned HTTP ${response.status}`); const bytes = await readBoundedBody(response, MARKDOWN_MAX_BYTES, "ATProto Markdown"); const markdown = bytes.toString("utf8"); for (const url of extractMarkdownImageUrls(markdown)) allowedImages.add(url); const imageResolution = await resolveBlueskyImageUrls( options.atUri, fetchImpl, resolveHostname, options.signal, ); for (const url of imageResolution.imageUrls) allowedImages.add(url); const enrichedMarkdown = imageResolution.imageUrls.length > 0 ? `${markdown.trimEnd()}\n\n## Resolved image URLs\n${imageResolution.imageUrls.map((url) => `- ${url}`).join("\n")}\n` : markdown; return { markdown: enrichedMarkdown, details: { atUri: options.atUri, endpoint: endpoint.toString(), mediaType: response.headers.get("content-type")?.split(";", 1)[0] ?? "text/markdown", sizeBytes: bytes.byteLength, sha256: digest(bytes), imageUrls: [...allowedImages], imageResolution, }, };}
export async function fetchBskyMarkdownDocument( options: BskyMarkdownFetchOptions,): Promise<BskyMarkdownDocument> { const match = /^at:\/\/([^/]+)\/app\.bsky\.feed\.post\/([^/?#]+)$/.exec(options.atUri); if (!match) throw new Error("bsky.md requires a Bluesky feed-post AT URI"); const [, actor, rkey] = match; const endpoint = new URL( `https://bsky-md.noz.am/profile/${encodeURIComponent(actor!)}/post/${encodeURIComponent(rkey!)}`, ); const fetchImpl = options.fetchImpl ?? fetch; const resolveHostname = options.resolveHostname ?? resolvePublicAddresses; await assertPublicHost(endpoint, resolveHostname, "Bluesky Markdown endpoint", options.signal); const response = await fetchImpl(endpoint, { headers: { accept: "text/markdown" }, redirect: "error", ...(options.signal ? { signal: options.signal } : {}), }); if (!response.ok) throw new Error(`bsky.md returned HTTP ${response.status}`); const bytes = await readBoundedBody(response, BSKY_MARKDOWN_MAX_BYTES, "Bluesky Markdown"); return { markdown: bytes.toString("utf8"), details: { atUri: options.atUri, endpoint: endpoint.toString(), mediaType: response.headers.get("content-type")?.split(";", 1)[0] ?? "text/markdown", sizeBytes: bytes.byteLength, sha256: digest(bytes), }, };}
async function resolveBlueskyImageUrls( atUri: string, fetchImpl: FetchLike, resolveHostname: ResolveHostname, signal: AbortSignal | undefined,): Promise<{ status: "succeeded" | "unavailable"; endpoint?: string; imageUrls: string[]; currentCid?: string; error?: string;}> { if (!/^at:\/\/[^/]+\/app\.bsky\.feed\.post\/[^/?#]+$/.test(atUri)) { return { status: "unavailable", imageUrls: [], error: "Record is not a Bluesky feed post" }; } const endpoint = new URL("https://public.api.bsky.app/xrpc/app.bsky.feed.getPosts"); endpoint.searchParams.append("uris", atUri); try { await assertPublicHost(endpoint, resolveHostname, "Bluesky AppView endpoint", signal); const response = await fetchImpl(endpoint, { headers: { accept: "application/json" }, redirect: "error", ...(signal ? { signal } : {}), }); if (!response.ok) return { status: "unavailable", endpoint: endpoint.toString(), imageUrls: [], error: `Bluesky AppView returned HTTP ${response.status}` }; const bytes = await readBoundedBody(response, APPVIEW_MAX_BYTES, "Bluesky AppView response"); const payload = JSON.parse(bytes.toString("utf8")) as unknown; const posts = asRecord(payload)?.posts; const records = Array.isArray(posts) ? posts.map(asRecord).filter((post) => post !== undefined) : []; const current = records.find((post) => post.uri === atUri) ?? records[0]; const currentCid = typeof current?.cid === "string" ? current.cid : undefined; const imageUrls = records.length > 0 ? records.flatMap((post) => extractImageUrls(post.embed)) : []; return { status: "succeeded", endpoint: endpoint.toString(), imageUrls: [...new Set(imageUrls)], ...(currentCid ? { currentCid } : {}), }; } catch (error) { return { status: "unavailable", endpoint: endpoint.toString(), imageUrls: [], error: error instanceof Error ? error.message : String(error), }; }}
const imageParameters = Type.Object({ url: Type.String({ description: "An image URL returned by the ATProto Markdown tool or present in the source record." }),});
function createImageDownloadTool( allowedImages: Set<string>, outcomes: EnrichmentOutcome[], fetchImpl: FetchLike, resolveHostname: ResolveHostname, artifactRoot: string | undefined, onTrace: (trace: RunnerTrace) => Promise<void>,): AgentTool<typeof imageParameters> { return { name: "download_image", label: "Download referenced image", description: "Download one image already discovered in the current AT record or fetched Markdown and return it for inspection.", parameters: imageParameters, executionMode: "sequential", execute: async (_toolCallId, params, signal) => { const request = { url: params.url }; try { if (!allowedImages.has(params.url)) throw new Error("Image URL was not discovered in the current record or fetched Markdown"); const { response, finalUrl } = await fetchPublic(params.url, fetchImpl, resolveHostname, signal); if (!response.ok) throw new Error(`Image origin returned HTTP ${response.status}`); const mediaType = response.headers.get("content-type")?.split(";", 1)[0]?.trim().toLowerCase() ?? ""; if (!mediaType.startsWith("image/")) throw new Error(`Expected image content but received ${mediaType || "unknown content type"}`); const bytes = await readBoundedBody(response, IMAGE_MAX_BYTES, "Image"); const sha256 = digest(bytes); const artifactPath = artifactRoot ? await writeArtifact(artifactRoot, sha256, bytes) : undefined; const result = { sourceUrl: params.url, finalUrl, mediaType, sizeBytes: bytes.byteLength, sha256, ...(artifactPath ? { artifactPath } : {}), }; await recordOutcome(outcomes, onTrace, { tool: "web.download-image", status: "succeeded", request, result, }); return { content: [ { type: "text", text: JSON.stringify(result) }, { type: "image", data: bytes.toString("base64"), mimeType: mediaType }, ], details: result, }; } catch (error) { await recordFailure(outcomes, onTrace, "web.download-image", request, error); throw error; } }, };}
function resolveAtUri(event: ThoughtEvent, target: "event" | "subject"): string { const payload = event.payload as Record<string, unknown>; if (target === "event") { const atUri = payload.atUri; if (typeof atUri === "string" && isAtUri(atUri)) return atUri; throw new Error("Current event has no canonical AT URI"); } const record = asRecord(payload.record); const subject = asRecord(record?.subject); const uri = subject?.uri; if (typeof uri === "string" && isAtUri(uri)) return uri; throw new Error("Current AT record has no referenced subject URI");}
function isAtUri(value: string): boolean { return /^at:\/\/[^/]+\/[a-zA-Z0-9.-]+\/[^/?#]+$/.test(value);}
async function fetchPublic( initialUrl: string, fetchImpl: FetchLike, resolveHostname: ResolveHostname, signal: AbortSignal | undefined,): Promise<{ response: Response; finalUrl: string }> { let current = initialUrl; for (let redirects = 0; redirects <= MAX_REDIRECTS; redirects += 1) { const url = new URL(current); if (url.protocol !== "https:") throw new Error("Only HTTPS image URLs are allowed"); await assertPublicHost(url, resolveHostname, "Image URL", signal); const response = await fetchImpl(url, { redirect: "manual", headers: { accept: "image/*" }, ...(signal ? { signal } : {}), }); if (![301, 302, 303, 307, 308].includes(response.status)) return { response, finalUrl: url.toString() }; const location = response.headers.get("location"); if (!location) throw new Error("Image redirect omitted its destination"); current = new URL(location, url).toString(); } throw new Error(`Image exceeded ${MAX_REDIRECTS} redirects`);}
async function assertPublicHost( url: URL, resolveHostname: ResolveHostname, label: string, signal?: AbortSignal | undefined,): Promise<void> { if (url.protocol !== "https:") throw new Error(`${label} must use HTTPS`); const addresses = await withAbort(resolveHostname(url.hostname), signal, `${label} resolution aborted`); if (addresses.length === 0 || addresses.some(isPrivateAddress)) { throw new Error(`${label} does not resolve exclusively to public addresses`); }}
async function withAbort<T>(promise: Promise<T>, signal: AbortSignal | undefined, message: string): Promise<T> { if (!signal) return promise; if (signal.aborted) throw new Error(message); return new Promise<T>((resolve, reject) => { const abort = () => reject(new Error(message)); signal.addEventListener("abort", abort, { once: true }); promise.then( (value) => { signal.removeEventListener("abort", abort); resolve(value); }, (error) => { signal.removeEventListener("abort", abort); reject(error); }, ); });}
async function resolvePublicAddresses(hostname: string): Promise<string[]> { if (isIP(hostname)) return [hostname]; const [ipv4, ipv6] = await Promise.all([ resolve4(hostname).catch(() => []), resolve6(hostname).catch(() => []), ]); return [...ipv4, ...ipv6];}
function isPrivateAddress(address: string): boolean { const normalized = address.toLowerCase(); if (normalized.startsWith("::ffff:")) return isPrivateAddress(normalized.slice(7)); const family = isIP(address); if (family === 0) return true; return nonPublicAddresses.check(address, family === 4 ? "ipv4" : "ipv6");}
function createNonPublicBlockList(): BlockList { const list = new BlockList(); for (const [network, prefix] of [ ["0.0.0.0", 8], ["10.0.0.0", 8], ["100.64.0.0", 10], ["127.0.0.0", 8], ["169.254.0.0", 16], ["172.16.0.0", 12], ["192.0.0.0", 24], ["192.0.2.0", 24], ["192.168.0.0", 16], ["198.18.0.0", 15], ["198.51.100.0", 24], ["203.0.113.0", 24], ["224.0.0.0", 4], ["240.0.0.0", 4], ] as const) list.addSubnet(network, prefix, "ipv4"); for (const [network, prefix] of [ ["::", 128], ["::1", 128], ["64:ff9b::", 96], ["100::", 64], ["2001:db8::", 32], ["fc00::", 7], ["fe80::", 10], ["ff00::", 8], ] as const) list.addSubnet(network, prefix, "ipv6"); return list;}
async function readBoundedBody(response: Response, maxBytes: number, label: string): Promise<Buffer> { const announced = Number(response.headers.get("content-length")); if (Number.isFinite(announced) && announced > maxBytes) throw new Error(`${label} exceeds ${maxBytes} bytes`); if (!response.body) return Buffer.alloc(0); const reader = response.body.getReader(); const chunks: Buffer[] = []; let size = 0; while (true) { const { value, done } = await reader.read(); if (done) break; size += value.byteLength; if (size > maxBytes) { await reader.cancel(); throw new Error(`${label} exceeds ${maxBytes} bytes`); } chunks.push(Buffer.from(value)); } return Buffer.concat(chunks, size);}
async function writeArtifact(root: string, sha256: string, bytes: Buffer): Promise<string> { const blob = await storeBlob(root, bytes); if (blob.sha256 !== sha256) throw new Error("Downloaded image digest changed before blob storage"); return blob.relativePath;}
function extractImageUrls(value: unknown): string[] { const urls = new Set<string>(); walk(value, undefined, urls); return [...urls];}
function walk(value: unknown, key: string | undefined, urls: Set<string>): void { if (typeof value === "string") { if (key && /^(image|fullsize|thumb|avatar|banner)$/i.test(key) && isHttpsUrl(value)) urls.add(value); return; } if (Array.isArray(value)) { for (const item of value) walk(item, key, urls); return; } if (!value || typeof value !== "object") return; for (const [childKey, child] of Object.entries(value)) walk(child, childKey, urls);}
function extractMarkdownImageUrls(markdown: string): string[] { const urls = new Set<string>(); const pattern = /!\[[^\]]*]\(<?(https:\/\/[^)>\s]+)>?(?:\s+"[^"]*")?\)/g; for (const match of markdown.matchAll(pattern)) { const url = match[1]; if (url && isHttpsUrl(url)) urls.add(url); } return [...urls];}
function isHttpsUrl(value: string): boolean { try { return new URL(value).protocol === "https:"; } catch { return false; }}
interface RawEnrichmentOutcome { tool: string; status: "succeeded" | "failed"; request: Record<string, unknown>; result?: Record<string, unknown> | undefined; error?: string | undefined;}
async function recordOutcome( outcomes: EnrichmentOutcome[], onTrace: (trace: RunnerTrace) => Promise<void>, outcome: RawEnrichmentOutcome,): Promise<void> { const safeOutcome: EnrichmentOutcome = { tool: outcome.tool, status: outcome.status, requestKeys: Object.keys(outcome.request).sort(), argumentsRedacted: true, resultPresent: outcome.result !== undefined, ...(outcome.result ? { resultKeys: Object.keys(outcome.result).sort(), resultSha256: digest(Buffer.from(JSON.stringify(outcome.result))), } : {}), ...(outcome.error ? { errorCode: "tool-execution-failed" } : {}), }; outcomes.push(safeOutcome); await onTrace({ kind: `enrichment.${outcome.status}`, data: safeOutcome, });}
async function recordFailure( outcomes: EnrichmentOutcome[], onTrace: (trace: RunnerTrace) => Promise<void>, tool: string, request: Record<string, unknown>, error: unknown,): Promise<void> { await recordOutcome(outcomes, onTrace, { tool, status: "failed", request, error: error instanceof Error ? error.message : String(error), });}
function digest(value: Uint8Array): string { return createHash("sha256").update(value).digest("hex");}
function asRecord(value: unknown): Record<string, unknown> | undefined { return value && typeof value === "object" && !Array.isArray(value) ? value as Record<string, unknown> : undefined;}
function assertNever(value: never): never { throw new Error(`Unsupported agent tool: ${String(value)}`);}