Something went wrong. Try again.
A local-first event pipeline for independent agents, built on Jazz.
Something went wrong. Try again.
10 kB · 241 lines
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242/** * Synthetic workbench demo. NOT a deployment surface. * * Seeds a temporary Jazz store with harmless fixture events, starts the * loopback inspector with a generated Review capability, and starts a second * loopback "demo front" that plays the role of the authenticated proxy: it * answers `/inspector/api/session` with write access, signs workbench and * decision POSTs with the capability under a CSRF check, and proxies every * other `/inspector/*` read to the inspector. There is no authentication at * all, so it binds 127.0.0.1 only and destroys its temporary store on exit. */import http, { type IncomingMessage, type ServerResponse } from "node:http";import fs from "node:fs/promises";import os from "node:os";import path from "node:path";import { randomBytes, timingSafeEqual, createHash } from "node:crypto";import { JazzThoughtStore } from "../src/jazz/store.js";import { REVIEW_CSRF_HEADER, REVIEW_NONCE_HEADER, REVIEW_SIGNATURE_HEADER, REVIEW_TIMESTAMP_HEADER, signReviewRequest,} from "../src/review/web-capability.js";import { startInspectorServer } from "../src/web/inspector.js";
process.env.THOUGHTSTREAM_JAZZ_AUTO ??= "1";
const HOST = "127.0.0.1";const inspectorPort = Number(process.env.WORKBENCH_DEMO_INSPECTOR_PORT ?? "4317");const frontPort = Number(process.env.WORKBENCH_DEMO_PORT ?? "4327");const MAX_WRITE_BODY = 98_304;
const root = await fs.mkdtemp(path.join(os.tmpdir(), "thoughtstream-workbench-demo-"));const store = await JazzThoughtStore.open({ projectRoot: root, appId: `thoughtstream-workbench-demo-${randomBytes(4).toString("hex")}`, runtimeRevision: "workbench-demo" });
await seed(store);
const reviewCapability = randomBytes(32);const csrfToken = randomBytes(24).toString("base64url");const inspector = await startInspectorServer(store, { host: HOST, port: inspectorPort, reviewCapability, proposalActor: "operator:demo",});
const front = http.createServer((request, response) => { void handle(request, response).catch(() => { if (!response.headersSent) { response.writeHead(500, { "content-type": "text/plain; charset=utf-8" }); } response.end("demo front error\n"); });});await new Promise<void>((resolve, reject) => { front.once("error", reject); front.listen(frontPort, HOST, () => { front.off("error", reject); resolve(); });});
process.stdout.write([ "", "thought stream workbench demo — SYNTHETIC, NO AUTHENTICATION, LOOPBACK ONLY", ` open: http://${HOST}:${frontPort}/inspector/`, ` inspector: http://${HOST}:${inspectorPort} (loopback, signed writes only)`, ` store: ${root} (temporary; removed on exit)`, " data: a few synthetic fixture observations, no real sources", " Ctrl+C to stop.", "",].join("\n"));
for (const signal of ["SIGINT", "SIGTERM"] as const) { process.once(signal, () => { front.closeAllConnections(); inspector.closeAllConnections(); front.close(() => { inspector.close(() => { void store.close() .then(() => fs.rm(root, { recursive: true, force: true })) .then(() => process.exit(0), () => process.exit(1)); }); }); setTimeout(() => process.exit(1), 3_000).unref(); });}
async function handle(request: IncomingMessage, response: ServerResponse): Promise<void> { const url = new URL(request.url ?? "/", `http://${HOST}`); if (url.pathname === "/" || url.pathname === "/inspector") { request.resume(); response.writeHead(302, { location: "/inspector/" }); response.end(); return; } if (!url.pathname.startsWith("/inspector/")) { request.resume(); response.writeHead(404, { "content-type": "text/plain; charset=utf-8" }); response.end("not found\n"); return; } if (url.pathname === "/inspector/api/session") { request.resume(); response.writeHead(200, { "content-type": "application/json; charset=utf-8", "cache-control": "no-store" }); response.end(JSON.stringify({ reviewWriteEnabled: true, courseChatEnabled: false, csrfToken })); return; } const writeRoute = url.search === "" && request.method === "POST" && ( /^\/inspector\/api\/(?:reviews|proposals)\/[^/]+\/decisions$/.test(url.pathname) || url.pathname === "/inspector/api/workbench/documents" || /^\/inspector\/api\/workbench\/documents\/[^/]+\/(?:versions|selections|proposals)$/.test(url.pathname) || /^\/inspector\/api\/workbench\/proposals\/[^/]+\/decisions$/.test(url.pathname) ); const upstreamPath = `${url.pathname.slice("/inspector".length)}${url.search}`; if (writeRoute) { const presented = request.headers[REVIEW_CSRF_HEADER]; if (typeof presented !== "string" || !safeEqual(presented, csrfToken)) { request.resume(); response.writeHead(403, { "content-type": "application/json; charset=utf-8" }); response.end(JSON.stringify({ error: "Request could not be verified (CSRF)" })); return; } let body: Buffer; try { body = await readJsonBody(request); JSON.parse(body.toString("utf8")); } catch { request.resume(); response.writeHead(400, { "content-type": "application/json; charset=utf-8" }); response.end(JSON.stringify({ error: "Request is invalid" })); return; } const signature = signReviewRequest(reviewCapability, { method: "POST", path: upstreamPath, body }); forward(request, response, upstreamPath, body, { "content-type": "application/json", "content-length": String(body.length), [REVIEW_TIMESTAMP_HEADER]: signature.timestamp, [REVIEW_NONCE_HEADER]: signature.nonce, [REVIEW_SIGNATURE_HEADER]: signature.signature, }); return; } if (request.method !== "GET" && request.method !== "HEAD") { request.resume(); response.writeHead(405, { "content-type": "application/json; charset=utf-8", allow: "GET, HEAD" }); response.end(JSON.stringify({ error: "Method not allowed" })); return; } forward(request, response, upstreamPath);}
function forward(request: IncomingMessage, response: ServerResponse, upstreamPath: string, body?: Buffer, extra: Record<string, string> = {}): void { const headers: Record<string, string> = { host: `${HOST}:${inspectorPort}` }; for (const name of ["accept", "accept-language", "user-agent"]) { const value = request.headers[name]; if (typeof value === "string") headers[name] = value; } Object.assign(headers, extra); const upstream = http.request({ hostname: HOST, port: inspectorPort, path: upstreamPath, method: request.method, headers }, (upstreamResponse) => { const responseHeaders: Record<string, string | string[]> = {}; for (const [name, value] of Object.entries(upstreamResponse.headers)) { if (value === undefined || name === "connection" || name === "keep-alive" || name === "transfer-encoding" || name === "set-cookie") continue; responseHeaders[name] = value; } response.writeHead(upstreamResponse.statusCode ?? 502, responseHeaders); upstreamResponse.pipe(response); }); upstream.on("error", () => { if (!response.headersSent) response.writeHead(502, { "content-type": "text/plain; charset=utf-8" }); response.end("inspector unavailable\n"); }); if (body) upstream.end(body); else if (request.method === "GET" || request.method === "HEAD") upstream.end(); else request.pipe(upstream);}
async function readJsonBody(request: IncomingMessage): Promise<Buffer> { const contentType = (request.headers["content-type"] ?? "").toLowerCase().split(";", 1)[0]?.trim(); if (contentType !== "application/json") throw new Error("Expected JSON body"); const parts: Buffer[] = []; let bytes = 0; for await (const part of request) { const buffer = Buffer.isBuffer(part) ? part : Buffer.from(part); bytes += buffer.length; if (bytes > MAX_WRITE_BODY) throw new Error("JSON body too large"); parts.push(buffer); } if (bytes === 0) throw new Error("JSON body is empty"); return Buffer.concat(parts);}
function safeEqual(left: string, right: string): boolean { const digest = (value: string) => createHash("sha256").update(value, "utf8").digest(); return timingSafeEqual(digest(left), digest(right));}
async function seed(target: JazzThoughtStore): Promise<void> { const base = Date.parse("2026-09-18T09:00:00.000Z"); const at = (minutes: number) => new Date(base + minutes * 60_000).toISOString(); const rssItems = [ { id: "fixture-1", title: "Fixture: notes on deterministic runners", summary: "A synthetic article about keeping proposal runners deterministic so tests stay reproducible." }, { id: "fixture-2", title: "Fixture: bounded context snapshots", summary: "A synthetic article on recording exact event ids and payload hashes instead of mutable titles." }, { id: "fixture-3", title: "Fixture: append-only decisions", summary: "A synthetic article arguing that human decisions should be appended, never edited in place." }, ]; for (const [index, item] of rssItems.entries()) { await target.appendEvent({ type: "stream.thought.source.rss.item", schemaVersion: 1, source: "rss:fixture-demo", sourceKind: "rss", externalId: item.id, idempotencyKey: `rss:${item.id}`, occurredAt: at(index * 7), actor: "rss:fixture-demo", correlationId: "demo-seed", privacy: "public-source", payload: { title: item.title, summary: item.summary, url: `https://example.invalid/${item.id}` }, }); } const telegramMessages = [ "Synthetic note to self: try turning the runner article into a working document.", "Synthetic reminder: the demo store is temporary and contains no real data.", ]; for (const [index, text] of telegramMessages.entries()) { await target.appendEvent({ type: "stream.thought.source.telegram.message", schemaVersion: 1, source: "telegram:fixture-demo", sourceKind: "telegram", externalId: `message-${index + 1}`, idempotencyKey: `telegram:demo:${index + 1}`, occurredAt: at(30 + index * 5), actor: "telegram:fixture-demo", correlationId: "demo-seed", privacy: "sensitive", payload: { chatId: "demo-chat", messageId: String(index + 1), senderName: "Fixture operator", text }, }); }}