Something went wrong. Try again.
A local-first event pipeline for independent agents, built on Jazz.
Something went wrong. Try again.
50 kB · 1082 lines
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083import fs from "node:fs/promises";import { afterEach, describe, expect, test, vi } from "vitest";import { startInspectorServer } from "../src/web/inspector.js";import type { JazzThoughtStore } from "../src/jazz/store.js";import type { AgentRun } from "../src/store/types.js";import { temporaryProject, testStore } from "./helpers.js";
const stores: JazzThoughtStore[] = [];const roots: string[] = [];const servers: import("node:http").Server[] = [];
afterEach(async () => { await Promise.all(servers.splice(0).map((server) => new Promise<void>((resolve) => { server.closeAllConnections(); server.close(() => resolve()); }))); await Promise.all(stores.splice(0).map((store) => store.close())); await Promise.all(roots.splice(0).map((root) => fs.rm(root, { recursive: true, force: true })));});
describe("thought stream inspector", () => { test("serves local read-only activity and rejects mutations", async () => { const project = await temporaryProject(); roots.push(project); const store = await testStore(project); stores.push(store); await store.appendEvent({ type: "stream.thought.source.rss.item", schemaVersion: 1, source: "rss:test", sourceKind: "rss", externalId: "item-1", idempotencyKey: "rss:item-1", occurredAt: "2026-07-14T00:00:00.000Z", actor: "rss:test", correlationId: "poll-1", privacy: "public-source", payload: { title: "Fixture" }, }); await store.upsertRun({ id: "run_contradictory", executionKey: "execution_fixture", triggerEventId: "event_fixture", agentId: "fixture-agent", agentVersion: 1, status: "completed", inputEventIds: [], outputEventIds: [], attempt: 1, provider: "deterministic", model: "deterministic", promptHash: "hash", contextManifest: {}, createdAt: "2026-07-14T00:00:00.000Z", completedAt: "2026-07-14T00:00:01.000Z", updatedAt: "2026-07-14T00:00:01.000Z", }); await store.appendEvent({ type: "stream.thought.connector.poll.started", schemaVersion: 1, source: "rss:test", sourceKind: "rss", externalId: "poll-1", idempotencyKey: "poll-1:started", occurredAt: "2026-07-14T00:00:00.000Z", actor: "rss:test", correlationId: "poll-1", privacy: "public-source", payload: { status: "started" }, }); await store.appendEvent({ type: "stream.thought.connector.failed", schemaVersion: 1, source: "rss:test", sourceKind: "rss", externalId: "poll-1", idempotencyKey: "poll-1:failed", occurredAt: "2026-07-14T00:00:02.000Z", actor: "rss:test", correlationId: "poll-1", privacy: "public-source", payload: { status: "failed", error: "fixture failure" }, }); await store.upsertSourceCursor({ id: "cursor:rss:test", source: "rss:test", cursor: { etag: "fixture-v1" }, lastSuccessAt: "2026-07-14T00:00:01.000Z", lastFailureAt: "2026-07-14T00:00:02.000Z", lastError: "fixture failure", updatedAt: "2026-07-14T00:00:02.000Z", }); const server = await startInspectorServer(store, { port: 0 }); servers.push(server); const address = server.address(); if (!address || typeof address === "string") throw new Error("Missing inspector address"); const base = `http://127.0.0.1:${address.port}`;
const page = await fetch(base); expect(page.status).toBe(200); const pageHtml = await page.text(); expect(pageHtml).toContain("<title>Stream</title>"); expect(pageHtml).toContain("<h1>Stream</h1>"); expect(pageHtml).toContain('src:url("assets/around-regular.woff2")'); expect(pageHtml).toContain('rel="manifest" href="manifest.webmanifest"'); expect(pageHtml).toContain('name="apple-mobile-web-app-title" content="Stream"'); expect(pageHtml).toContain("@media(max-width:640px)"); expect(pageHtml).toContain("text-transform:uppercase"); expect(pageHtml).toContain('class="primary-nav" aria-label="Stream destinations"'); expect(pageHtml).toContain('<button id="events-tab" class="active">Feed</button>'); expect(pageHtml).toContain('<button id="learn-tab">Learn</button>'); expect(pageHtml).toContain('<button id="review-tab">Review</button>'); expect(pageHtml).toContain('<button id="artifacts-tab">Artifacts</button>'); expect(pageHtml).toContain('<button id="system-tab">System</button>'); expect(pageHtml).toContain('<button id="filter-open" class="filter-nav" hidden>Filter'); expect(pageHtml).toContain('id="review-nav" class="context-nav"'); expect(pageHtml).toContain('id="system-nav" class="context-nav"'); expect(pageHtml).toContain(".primary-nav button { min-width:max-content; min-height:32px"); expect(pageHtml).toContain(".primary-nav button { min-height:36px; padding:5px 16px"); expect(pageHtml).toContain(".context-nav button { min-height:36px; padding:5px 15px }"); expect(pageHtml).not.toContain(".app-header { position:sticky"); expect(pageHtml).toContain(".primary-nav button.active { border-color:var(--selected-bg); background:var(--selected-bg); color:var(--selected-text) }"); expect(pageHtml).toContain(".telegram-inlay { border:1px solid var(--line); border-radius:20px; background:var(--surface-control)"); expect(pageHtml).toContain("overflow-x:auto"); expect(pageHtml).not.toContain("grid-template-columns:repeat(6,minmax(0,1fr))"); expect(pageHtml).toContain("const primaryResource=tabResource();"); expect(pageHtml).toContain("await loadResource(primaryResource)"); expect(pageHtml).toContain("for(const key of ['activity','course','artifacts','documents','proposals','reviews','system'])"); expect(pageHtml).toContain("tabResource()===key||(key==='system'&&state.tab==='events')"); expect(pageHtml).toContain("activity:{path:'api/activity',label:'recent activity'}"); expect(pageHtml).toContain("system:{path:'api/system',label:'system evidence'}"); expect(pageHtml).not.toContain("Promise.all([fetch('api/snapshot')"); expect(pageHtml).not.toContain("fetch('/api/snapshot')"); expect(pageHtml).toContain("response.headers.get('content-type')"); expect(pageHtml).toContain("Loading recent activity…"); expect(pageHtml).toContain("Loading '+label+'…"); expect(pageHtml).toContain("detail-retry"); expect(pageHtml).toContain('id="filter-bar"'); expect(pageHtml).toContain('id="filters-dialog" class="filter-dialog"'); expect(pageHtml).toContain('<h2 id="filter-title">Filter feed</h2>'); expect(pageHtml).toContain('aria-label="Back to list"'); expect(pageHtml).toContain('id="active-filters" class="filter-chips"'); expect(pageHtml).toContain(".primary-nav .filter-nav { margin-left:auto }"); expect(pageHtml).toContain(".filter-count[hidden] { display:none }"); expect(pageHtml).toContain('<h2 id="list-title" class="visually-hidden">Recent activity</h2>'); expect(pageHtml).not.toContain('class="section-head"'); expect(pageHtml).toContain("filterDialog.showModal()"); expect(pageHtml).toContain("max-height:min(78svh,720px)"); expect(pageHtml).not.toContain("var(--surface-alt)"); expect(pageHtml).toContain("All sources"); expect(pageHtml).toContain("is configured but has no durable activity yet."); expect(pageHtml).toContain("No agent processed this observation."); expect(pageHtml).toContain("Earlier processing may be outside this window."); expect(pageHtml).toContain("Some older observations are outside this bounded feed window."); expect(pageHtml).toContain("api/activity/evidence"); expect(pageHtml).toContain("loadActivityEvidence"); expect(pageHtml).toContain("processed||!processingComplete"); expect(pageHtml).toContain("Technical details"); expect(pageHtml).not.toContain("No consumer processed this event."); expect(pageHtml).toContain("No recent activity has been recorded."); expect(pageHtml).toContain("No artifacts are available."); expect(pageHtml).toContain("new EventSource('api/live')"); expect(pageHtml).toContain("scheduleActivityRefresh"); expect(pageHtml).toContain("Internal producers"); expect(pageHtml).toContain("Loading source activity…"); expect(pageHtml).toContain("id==='resident-letta-conversation'?'Stream'"); expect(pageHtml).not.toContain('id="status"'); expect(pageHtml).not.toContain("#status"); expect(pageHtml).not.toContain("updateStatus"); expect(pageHtml).not.toContain('id="count"'); expect(pageHtml).not.toContain("review enabled"); expect(pageHtml).toContain("main.detail-selected #list-pane"); expect(pageHtml).toContain("body.detail-selected .primary-nav,body.detail-selected .context-nav { display:none }"); expect(pageHtml).toContain("document.body.classList.add('detail-selected')"); expect(pageHtml).toContain("const routeLists={events:'feed',learn:'learn/post-training',suggestions:'review/suggestions',reviews:'review/comparisons',artifacts:'artifacts',documents:'documents',runs:'system/runs',sources:'system/sources'}"); expect(pageHtml).toContain('<button id="documents-tab">Documents</button>'); expect(pageHtml).toContain("document:{tab:'documents',path:'documents'}"); expect(pageHtml).toContain("event:{tab:'events',path:'feed/observations'}"); expect(pageHtml).toContain('id="course-chat-toggle"'); expect(pageHtml).toContain("Ask about this lesson"); expect(pageHtml).toContain("api/courses/post-training/questions"); expect(pageHtml).toContain("workshopBody(kind)"); expect(pageHtml).toContain("Scores come from the full functional Pi gate"); expect(pageHtml).not.toContain("Pi transfer gate passes"); expect(pageHtml).toContain('class="course-glossary"'); expect(pageHtml).toContain("thoughtstream-course-post-training-receipts"); expect(pageHtml).toContain("sessionStorage.setItem"); expect(pageHtml).toContain("Math.min(5000,Math.round(delay*1.45))"); const inlineScript = pageHtml.match(/<script>([\s\S]*?)<\/script>/)?.[1]; expect(inlineScript).toBeDefined(); expect(() => new Function(inlineScript!)).not.toThrow(); expect(pageHtml).toContain("history.pushState({streamInspector:true,hasInternalParent:true}"); expect(pageHtml).toContain("history.replaceState(history.state?.streamInspector===true"); expect(pageHtml).toContain("window.addEventListener('popstate',()=>applyRoute(parseRoute()))"); expect(pageHtml).toContain("function navigateBack(){if(history.state?.streamInspector===true&&history.state.hasInternalParent===true)history.back()"); expect(pageHtml).toContain(".observation-card .bluesky-inlay,.observation-card .x-inlay { border:0; border-radius:0; background:transparent; padding:0 }"); expect(pageHtml).toContain('class="item event-item" tabindex="0" data-kind="event"'); expect(pageHtml).toContain("feed-processing"); expect(pageHtml).toContain("renderSourceInlay"); expect(pageHtml).toContain("'x-post':renderXInlay"); expect(pageHtml).toContain("Replying to <a href=\"https://x.com/"); expect(pageHtml).toContain("media attachment"); expect(pageHtml).toContain("api/inlays/bluesky?"); expect(pageHtml).toContain("telegram-inlay"); expect(pageHtml).toContain("bluesky-inlay"); expect(pageHtml).toContain("id==='conceptualizer'?'Conceptualizer'"); expect(pageHtml).toContain(".item:hover,.item.active { background:transparent }"); expect(pageHtml).toContain(".item:hover .telegram-inlay"); expect(pageHtml).toContain("cursor:pointer; overflow-wrap:anywhere"); expect(pageHtml).toContain("font:15px/1.58 var(--font-body)"); expect(pageHtml).toContain('<script src="/assets/font-debug.js" defer></script>'); expect(pageHtml).not.toContain("border-bottom:1px solid var(--line); border-radius:0; margin:0; padding:20px 0"); expect(page.headers.get("content-security-policy")).toContain("frame-ancestors 'none'"); expect(page.headers.get("content-security-policy")).toContain("img-src 'self' data:"); expect(page.headers.get("content-security-policy")).toContain("worker-src 'self'"); expect(page.headers.get("content-security-policy")).toContain("https://fonts.googleapis.com"); expect(page.headers.get("content-security-policy")).toContain("https://fonts.gstatic.com"); expect(page.headers.get("content-security-policy")).not.toContain("img-src 'self' data: https:");
const font = await fetch(`${base}/assets/around-regular.woff2`); expect(font.status).toBe(200); expect(font.headers.get("content-type")).toBe("font/woff2"); expect(Buffer.from(await font.arrayBuffer()).subarray(0, 4).toString("ascii")).toBe("wOF2"); const fontDebug = await fetch(`${base}/assets/font-debug.js`); expect(fontDebug.headers.get("content-type")).toBe("text/javascript; charset=utf-8"); const fontDebugScript = await fontDebug.text(); expect(fontDebugScript).toContain("fonts.googleapis.com/css2?family="); expect(fontDebugScript).toContain("event.altKey && event.shiftKey"); expect(fontDebugScript).toContain('querySelectorAll("link[data-stream-font-debug]")'); expect(fontDebugScript).toContain("pendingLink?.remove()"); expect(() => new Function(fontDebugScript)).not.toThrow(); const manifest = await fetch(`${base}/manifest.webmanifest`); expect(manifest.headers.get("content-type")).toBe("application/manifest+json; charset=utf-8"); expect(await manifest.json()).toMatchObject({ name: "Stream", start_url: "./", display: "standalone" }); const icon = await fetch(`${base}/assets/stream-icon.svg`); expect(icon.headers.get("content-type")).toBe("image/svg+xml; charset=utf-8"); expect(await icon.text()).toContain(">S</text>"); const serviceWorker = await fetch(`${base}/sw.js`); expect(serviceWorker.headers.get("content-type")).toBe("text/javascript; charset=utf-8"); expect(await serviceWorker.text()).not.toContain("fetch");
const snapshot = await (await fetch(`${base}/api/snapshot`)).json() as { activity: { totalEvents: number }; evidenceContradictions: number; runs: Array<{ id: string; summary?: string; contextManifest?: unknown }>; runEvidence: Array<{ runId: string; consistent: boolean; issues: Array<{ code: string }> }>; sources: Array<{ source: string; status: string; recordCount: number; operationsStarted: number; operationFailures: number; inFlight: number; cursor: { etag?: string } }>; }; await fetch(`${base}/api/activity/evidence`); const rootEventsSpy = vi.spyOn(store, "listRecentRootEvents"); const recentRunsSpy = vi.spyOn(store, "listRecentRuns"); const activity = await (await fetch(`${base}/api/activity`)).json() as { totalEvents: number; processingPending: boolean }; expect(rootEventsSpy).not.toHaveBeenCalled(); expect(recentRunsSpy).not.toHaveBeenCalled(); rootEventsSpy.mockRestore(); recentRunsSpy.mockRestore(); const system = await (await fetch(`${base}/api/system`)).json() as { evidenceContradictions: number; sources: unknown[] }; expect(activity.totalEvents).toBe(3); expect(activity.processingPending).toBe(true); expect(system.evidenceContradictions).toBe(1); expect(system.sources).toHaveLength(1); expect(snapshot.activity.totalEvents).toBe(3); expect(snapshot.evidenceContradictions).toBe(1); expect(snapshot.runs[0]).toMatchObject({ id: "run_contradictory" }); expect(snapshot.runs[0]?.contextManifest).toBeUndefined(); expect(snapshot.runEvidence[0]).toMatchObject({ runId: "run_contradictory", consistent: false, issues: expect.arrayContaining([ expect.objectContaining({ code: "missing-terminal-event" }), expect.objectContaining({ code: "completed-without-output" }), ]), }); expect(snapshot.sources).toEqual([expect.objectContaining({ source: "rss:test", status: "failing", recordCount: 3, operationsStarted: 1, operationFailures: 1, inFlight: 0, cursor: { etag: "fixture-v1" }, })]); const runDetail = await (await fetch(`${base}/api/runs/run_contradictory`)).json() as { evidence: { consistent: boolean }; }; expect(runDetail.evidence.consistent).toBe(false); const sourceDetail = await (await fetch(`${base}/api/sources/${encodeURIComponent("rss:test")}`)).json() as { source: { source: string }; activity: { totalEvents: number; items: Array<{ presentation: { title: string; body: string } }> }; }; expect(sourceDetail).toMatchObject({ source: { source: "rss:test" }, activity: { totalEvents: 1, items: [{ presentation: { title: "Item", body: "Fixture" } }], }, }); const mutation = await fetch(`${base}/api/snapshot`, { method: "POST" }); expect(mutation.status).toBe(405); expect(await mutation.json()).toEqual({ error: "Inspector is read-only" }); });
test("projects direct Telegram messages and strong-reference Bluesky post inlays", async () => { const project = await temporaryProject(); roots.push(project); const store = await testStore(project); stores.push(store); const telegram = await store.appendEvent({ type: "stream.thought.source.telegram.message", schemaVersion: 1, source: "telegram:fixture", sourceKind: "telegram", externalId: "message-42", idempotencyKey: "telegram:message-42", occurredAt: "2026-08-11T01:00:00.000Z", actor: "user-fixture", correlationId: "chat-fixture", privacy: "sensitive", payload: { senderName: "Cameron", text: "Show the message directly.", editedAt: "2026-08-11T01:00:01.000Z", attachments: [{ kind: "image", status: "stored" }], }, }); await store.appendEvent({ type: "stream.thought.source.atproto.commit", schemaVersion: 1, source: "jetstream:cameron-atproto", sourceKind: "jetstream", externalId: "like-42", idempotencyKey: "atproto:like-42", occurredAt: "2026-08-11T01:01:00.000Z", actor: "did:plc:cameron", correlationId: "commit-fixture", privacy: "public-source", payload: { operation: "create", collection: "app.bsky.feed.like", record: { subject: { uri: "at://did:plc:author/app.bsky.feed.post/post42", cid: "bafyreifixture42" } }, }, }); const unsafe = await store.appendEvent({ type: "stream.thought.source.atproto.commit", schemaVersion: 1, source: "jetstream:cameron-atproto", sourceKind: "jetstream", externalId: "like-unsafe", idempotencyKey: "atproto:like-unsafe", occurredAt: "2026-08-11T01:02:00.000Z", actor: "did:plc:cameron", correlationId: "commit-fixture", privacy: "public-source", payload: { operation: "create", collection: "app.bsky.feed.like", record: { subject: { uri: "javascript:alert(1)", cid: "bafyreifixturebad" } }, }, }); const server = await startInspectorServer(store, { port: 0 }); servers.push(server); const address = server.address(); if (!address || typeof address === "string") throw new Error("Missing inspector address"); const base = `http://127.0.0.1:${address.port}`; const snapshot = await (await fetch(`${base}/api/snapshot`)).json() as { activity: { items: Array<{ id: string; type: string; presentation: Record<string, unknown> }> }; }; expect(snapshot.activity.items).toEqual(expect.arrayContaining([ expect.objectContaining({ type: "stream.thought.source.telegram.message", presentation: expect.objectContaining({ renderer: "telegram-message", body: "Show the message directly.", telegram: { senderName: "Cameron", edited: true, attachmentCount: 1 }, }), }), expect.objectContaining({ type: "stream.thought.source.atproto.commit", presentation: expect.objectContaining({ renderer: "bluesky-post", bluesky: { uri: "at://did:plc:author/app.bsky.feed.post/post42", cid: "bafyreifixture42" }, }), }), ])); const unsafeItem = snapshot.activity.items.find((item) => item.id === unsafe.event.id); expect(unsafeItem?.presentation).toMatchObject({ renderer: "generic" }); expect(unsafeItem?.presentation).not.toHaveProperty("url"); expect(unsafeItem?.presentation).not.toHaveProperty("bluesky"); const detail = await (await fetch(`${base}/api/events/${encodeURIComponent(telegram.event.id)}`)).json() as { presentation: Record<string, unknown>; }; expect(detail.presentation).toMatchObject({ renderer: "telegram-message", body: "Show the message directly." }); const invalid = await fetch(`${base}/api/inlays/bluesky?uri=${encodeURIComponent("https://example.com")}`); expect(invalid.status).toBe(400); const requestFetch = globalThis.fetch.bind(globalThis); let mediaFetches = 0; const appViewFetch = vi.spyOn(globalThis, "fetch").mockImplementation(async (input, init) => { const href = input instanceof Request ? input.url : String(input); if (href.startsWith("https://cdn.bsky.app/img/")) { mediaFetches += 1; expect(init?.redirect).toBe("error"); if (href.endsWith("/text")) return new Response("not an image", { status: 200, headers: { "content-type": "text/html" } }); if (href.endsWith("/bad-magic")) return new Response("not jpeg", { status: 200, headers: { "content-type": "image/jpeg" } }); if (href.endsWith("/oversized")) return new Response(Uint8Array.from([0xff, 0xd8, 0xff]), { status: 200, headers: { "content-type": "image/jpeg", "content-length": String(7 * 1024 * 1024 + 1) } }); return new Response(Uint8Array.from([0xff, 0xd8, 0xff, 0xd9]), { status: 200, headers: { "content-type": "image/jpeg", "content-length": "4" }, }); } if (!href.startsWith("https://public.api.bsky.app/xrpc/app.bsky.feed.getPosts")) { return requestFetch(input, init); } return new Response(JSON.stringify({ posts: [{ uri: "at://did:plc:author/app.bsky.feed.post/post42", cid: "bafyreifixture42", author: { handle: "author.example", displayName: "Author", avatar: "https://cdn.bsky.app/img/avatar/plain/did:plc:author/avatar@jpeg", }, record: { text: "The actual post body.", createdAt: "2026-08-11T01:00:00.000Z" }, embed: { $type: "app.bsky.embed.images#view", images: [{ thumb: "https://cdn.bsky.app/img/feed_thumbnail/plain/did:plc:author/photo@jpeg", fullsize: "https://cdn.bsky.app/img/feed_fullsize/plain/did:plc:author/photo@jpeg", alt: "A fixture image", }], }, replyCount: 2, repostCount: 3, likeCount: 5, }], }), { status: 200, headers: { "content-type": "application/json" } }); }); try { const query = new URLSearchParams({ uri: "at://did:plc:author/app.bsky.feed.post/post42", cid: "bafyreifixture42", }); const inlayResponse = await requestFetch(`${base}/api/inlays/bluesky?${query.toString()}`); expect(inlayResponse.status).toBe(200); const inlay = await inlayResponse.json() as { author: { avatar: string }; images: Array<{ thumb: string; fullsize: string; alt: string }>; }; expect(inlay).toMatchObject({ uri: "at://did:plc:author/app.bsky.feed.post/post42", cid: "bafyreifixture42", url: "https://bsky.app/profile/did:plc:author/post/post42", author: { handle: "author.example", displayName: "Author", avatar: expect.stringContaining("api/inlays/bluesky-media?url=https%3A%2F%2Fcdn.bsky.app%2Fimg%2Favatar"), }, text: "The actual post body.", images: [expect.objectContaining({ thumb: expect.stringContaining("api/inlays/bluesky-media?url=https%3A%2F%2Fcdn.bsky.app%2Fimg%2Ffeed_thumbnail"), fullsize: "https://cdn.bsky.app/img/feed_fullsize/plain/did:plc:author/photo@jpeg", alt: "A fixture image", })], counts: { replies: 2, reposts: 3, likes: 5 }, }); const media = await requestFetch(`${base}/${inlay.images[0]!.thumb}`); expect(media.status).toBe(200); expect(media.headers.get("content-type")).toBe("image/jpeg"); expect(Buffer.from(await media.arrayBuffer())).toEqual(Buffer.from([0xff, 0xd8, 0xff, 0xd9])); const cachedMedia = await requestFetch(`${base}/${inlay.images[0]!.thumb}`); expect(cachedMedia.status).toBe(200); expect(mediaFetches).toBe(1); const wrongMediaHost = await requestFetch(`${base}/api/inlays/bluesky-media?url=${encodeURIComponent("https://example.com/img/photo.jpg")}`); expect(wrongMediaHost.status).toBe(400); const wrongMediaPath = await requestFetch(`${base}/api/inlays/bluesky-media?url=${encodeURIComponent("https://cdn.bsky.app/not-images/photo.jpg")}`); expect(wrongMediaPath.status).toBe(400); const duplicateMedia = await requestFetch(`${base}/api/inlays/bluesky-media?url=${encodeURIComponent("https://cdn.bsky.app/img/photo.jpg")}&url=${encodeURIComponent("https://cdn.bsky.app/img/other.jpg")}`); expect(duplicateMedia.status).toBe(400); for (const unsafe of [ "http://cdn.bsky.app/img/photo.jpg", "https://user:password@cdn.bsky.app/img/photo.jpg", "https://cdn.bsky.app:8443/img/photo.jpg", "https://cdn.bsky.app/img/photo.jpg#fragment", ]) { expect((await requestFetch(`${base}/api/inlays/bluesky-media?url=${encodeURIComponent(unsafe)}`)).status).toBe(400); } for (const invalidMedia of ["text", "bad-magic", "oversized"]) { const source = `https://cdn.bsky.app/img/${invalidMedia}`; expect((await requestFetch(`${base}/api/inlays/bluesky-media?url=${encodeURIComponent(source)}`)).status).toBe(502); } query.set("cid", "bafyreibadversion"); const mismatch = await requestFetch(`${base}/api/inlays/bluesky?${query.toString()}`); expect(mismatch.status).toBe(502); } finally { appViewFetch.mockRestore(); } });
test("shows one broad-batch Stream run as secondary processing on every member observation", async () => { const project = await temporaryProject(); roots.push(project); const store = await testStore(project); stores.push(store); const first = (await store.appendEvent({ type: "stream.thought.source.rss.item", schemaVersion: 1, source: "rss:first", sourceKind: "rss", externalId: "first", idempotencyKey: "first", occurredAt: "2026-08-10T17:00:00.000Z", actor: "rss:first", correlationId: "poll:first", privacy: "public-source", payload: { title: "First source observation" }, })).event; const second = (await store.appendEvent({ type: "stream.thought.source.rss.item", schemaVersion: 1, source: "rss:second", sourceKind: "rss", externalId: "second", idempotencyKey: "second", occurredAt: "2026-08-10T17:00:01.000Z", actor: "rss:second", correlationId: "poll:second", privacy: "private", payload: { title: "Second source observation" }, })).event; const member = (event: typeof first) => ({ eventId: event.id, source: event.source, sourceSequence: event.sourceSequence, type: event.type, schemaVersion: event.schemaVersion, privacy: event.privacy, occurredAt: event.occurredAt, observedAt: event.observedAt, payloadHash: event.payloadHash, }); const batch = (await store.appendEvent({ type: "stream.thought.derived.event.batch", schemaVersion: 1, source: "batch:stream-activity", sourceKind: "system", externalId: "activity-window", idempotencyKey: "activity-window", occurredAt: "2026-08-10T17:10:00.000Z", actor: "batch:stream-activity-ten-minute", correlationId: "activity-window", privacy: "private", payload: { declaration: { id: "stream-activity-ten-minute", version: 1, fingerprint: "a".repeat(64) }, flushReason: "max-age", firstOccurredAt: first.occurredAt, lastOccurredAt: second.occurredAt, members: [member(first), member(second)], }, })).event; const output = (await store.appendEvent({ type: "stream.thought.derived.message.observation", schemaVersion: 1, source: "agent:resident-letta-conversation", sourceKind: "agent", externalId: "run_stream_window", idempotencyKey: "run_stream_window:output", occurredAt: "2026-08-10T17:10:02.000Z", actor: "resident-letta-conversation", rootEventId: batch.rootEventId, parentEventId: batch.id, correlationId: batch.correlationId, privacy: "private", payload: { runId: "run_stream_window", summary: "The two observations form one concrete pattern." }, traceId: "run_stream_window", })).event; await store.upsertRun({ id: "run_stream_window", executionKey: "execution:run_stream_window", triggerEventId: batch.id, agentId: "resident-letta-conversation", agentVersion: 5, status: "completed", inputEventIds: [batch.id], outputEventIds: [output.id], attempt: 1, provider: "letta-cloud", model: "fixture-model", promptHash: "fixture-prompt", contextManifest: {}, result: { summary: "The two observations form one concrete pattern.", tags: ["synthesis"], importance: "normal", confidence: 0.8 }, createdAt: "2026-08-10T17:10:00.000Z", completedAt: "2026-08-10T17:10:02.000Z", updatedAt: "2026-08-10T17:10:02.000Z", });
const server = await startInspectorServer(store, { port: 0 }); servers.push(server); const address = server.address(); if (!address || typeof address === "string") throw new Error("Missing inspector address"); const base = `http://127.0.0.1:${address.port}`; const snapshot = await (await fetch(`${base}/api/snapshot`)).json() as { activity: { items: Array<{ id: string; consumerRuns: Array<{ id: string }> }> }; }; expect(snapshot.activity.items.filter((item) => [first.id, second.id].includes(item.id))) .toEqual(expect.arrayContaining([ expect.objectContaining({ id: first.id, consumerRuns: [expect.objectContaining({ id: "run_stream_window" })] }), expect.objectContaining({ id: second.id, consumerRuns: [expect.objectContaining({ id: "run_stream_window" })] }), ])); const detail = await (await fetch(`${base}/api/events/${encodeURIComponent(first.id)}`)).json() as { agentActivities: Array<{ run: { id: string }; inputs: Array<{ id: string }> }>; }; expect(detail.agentActivities).toEqual([ expect.objectContaining({ run: expect.objectContaining({ id: "run_stream_window" }), inputs: [expect.objectContaining({ id: batch.id })], }), ]); });
test("preserves recent roots across a child-heavy source window without listing all runs", async () => { const project = await temporaryProject(); roots.push(project); const store = await testStore(project); stores.push(store); const root = (await store.appendEvent({ type: "stream.thought.source.rss.item", schemaVersion: 1, source: "rss:child-window", sourceKind: "rss", externalId: "root", idempotencyKey: "root", occurredAt: "2026-08-22T10:00:00.000Z", observedAt: "2026-08-22T10:00:00.000Z", actor: "rss:child-window", correlationId: "root", privacy: "public-source", payload: { title: "Root observation" }, })).event; const secondRoot = (await store.appendEvent({ type: "stream.thought.source.rss.item", schemaVersion: 1, source: "rss:child-window", sourceKind: "rss", externalId: "second-root", idempotencyKey: "second-root", occurredAt: "2026-08-22T10:00:30.000Z", observedAt: "2026-08-22T10:00:30.000Z", actor: "rss:child-window", correlationId: "second-root", privacy: "public-source", payload: { title: "Second root observation" }, })).event; for (let index = 0; index < 100; index += 1) { await store.appendEvent({ type: "stream.thought.source.rss.item", schemaVersion: 1, source: "rss:child-window", sourceKind: "rss", externalId: `child-${index}`, idempotencyKey: `child-${index}`, occurredAt: `2026-08-22T10:01:${String(index % 60).padStart(2, "0")}.000Z`, observedAt: `2026-08-22T10:${String(1 + Math.floor(index / 60)).padStart(2, "0")}:${String(index % 60).padStart(2, "0")}.000Z`, actor: "rss:child-window", rootEventId: root.id, parentEventId: root.id, correlationId: "root", privacy: "public-source", payload: { title: `Child ${index}` }, }); }
const listRunsSpy = vi.spyOn(store, "listRuns"); const recentRunsSpy = vi.spyOn(store, "listRecentRuns"); const server = await startInspectorServer(store, { port: 0 }); servers.push(server); const address = server.address(); if (!address || typeof address === "string") throw new Error("Missing inspector address");
const activity = await (await fetch(`http://127.0.0.1:${address.port}/api/activity/evidence`)).json() as { totalEvents: number; items: Array<{ id: string }>; rootWindowComplete: boolean; processingPending: boolean; }; expect(activity.totalEvents).toBe(102); expect(activity.rootWindowComplete).toBe(true); expect(activity.processingPending).toBe(false); expect(activity.items).toEqual([ expect.objectContaining({ id: secondRoot.id }), expect.objectContaining({ id: root.id }), ]); expect(listRunsSpy).not.toHaveBeenCalled(); expect(recentRunsSpy).toHaveBeenCalledWith(250); });
test("recovers a child-triggered run outside the recent run window and marks history incomplete", async () => { const project = await temporaryProject(); roots.push(project); const store = await testStore(project); stores.push(store); const root = (await store.appendEvent({ type: "stream.thought.source.rss.item", schemaVersion: 1, source: "rss:child-run", sourceKind: "rss", externalId: "root", idempotencyKey: "root", occurredAt: "2026-08-22T10:00:00.000Z", actor: "rss:child-run", correlationId: "root", privacy: "public-source", payload: { title: "Root with child processing" }, })).event; const child = (await store.appendEvent({ type: "stream.thought.source.rss.item", schemaVersion: 1, source: "rss:child-run", sourceKind: "rss", externalId: "child", idempotencyKey: "child", occurredAt: "2026-08-22T10:01:00.000Z", actor: "rss:child-run", rootEventId: root.id, parentEventId: root.id, correlationId: "root", privacy: "public-source", payload: { title: "Child activity" }, })).event; await store.upsertRun({ id: "run_child_activity", executionKey: "execution:run_child_activity", triggerEventId: child.id, agentId: "child-processor", agentVersion: 1, status: "completed", inputEventIds: [child.id], outputEventIds: [], attempt: 1, provider: "deterministic", model: "deterministic", promptHash: "child-prompt", contextManifest: {}, result: { summary: "Processed child activity" }, createdAt: "2026-08-22T10:01:01.000Z", completedAt: "2026-08-22T10:01:02.000Z", updatedAt: "2026-08-22T10:01:02.000Z", }); const fillerRuns: AgentRun[] = Array.from({ length: 250 }, (_, index) => ({ id: `run_filler_${index}`, executionKey: `execution:run_filler_${index}`, triggerEventId: `missing_event_${index}`, agentId: "filler", agentVersion: 1, status: "running", inputEventIds: [], outputEventIds: [], attempt: 1, provider: "deterministic", model: "deterministic", promptHash: `filler_${index}`, contextManifest: {}, createdAt: `2026-08-22T11:${String(Math.floor(index / 60) % 60).padStart(2, "0")}:${String(index % 60).padStart(2, "0")}.000Z`, updatedAt: "2026-08-22T12:00:00.000Z", })); vi.spyOn(store, "listRecentRuns").mockResolvedValue(fillerRuns);
const server = await startInspectorServer(store, { port: 0 }); servers.push(server); const address = server.address(); if (!address || typeof address === "string") throw new Error("Missing inspector address"); const activity = await (await fetch(`http://127.0.0.1:${address.port}/api/activity/evidence`)).json() as { items: Array<{ id: string; consumerRunsComplete: boolean; consumerRuns: Array<{ id: string }> }>; }; expect(activity.items).toEqual([ expect.objectContaining({ id: root.id, consumerRunsComplete: false, consumerRuns: [expect.objectContaining({ id: "run_child_activity" })], }), ]); });
test("streams Jazz activity signals without exposing event content in the stream", async () => { const project = await temporaryProject(); roots.push(project); const store = await testStore(project); stores.push(store); const server = await startInspectorServer(store, { port: 0 }); servers.push(server); const address = server.address(); if (!address || typeof address === "string") throw new Error("Missing inspector address"); const controller = new AbortController(); const response = await fetch(`http://127.0.0.1:${address.port}/api/live`, { signal: controller.signal }); expect(response.status).toBe(200); expect(response.headers.get("content-type")).toBe("text/event-stream; charset=utf-8"); expect(response.headers.get("cache-control")).toBe("no-store, no-transform"); if (!response.body) throw new Error("Inspector live stream has no body"); const frames = sseFrames(response.body.getReader()); expect(await frames.next()).toContain("retry: 2000"); expect(await frames.next()).toContain("event: change");
const sourceText = "must-not-enter-the-live-signal"; await store.appendEvent({ type: "stream.thought.source.rss.item", schemaVersion: 1, source: "rss:live-fixture", sourceKind: "rss", externalId: "live-item-1", idempotencyKey: "live-item-1", occurredAt: "2026-08-10T17:00:00.000Z", actor: "rss:live-fixture", correlationId: "live-item-1", privacy: "public-source", payload: { title: sourceText }, }); const change = await frames.next(); expect(change).toContain("event: change"); expect(change).toContain('"at"'); expect(change).not.toContain(sourceText); controller.abort(); });
test("puts the semantic agent result ahead of lifecycle machinery", async () => { const project = await temporaryProject(); roots.push(project); const store = await testStore(project); stores.push(store); const source = (await store.appendEvent({ type: "stream.thought.source.file.changed", schemaVersion: 1, source: "filesystem:test", sourceKind: "filesystem", externalId: "spec/brief.md", idempotencyKey: "file:changed:brief", occurredAt: "2026-07-14T00:00:00.000Z", actor: "filesystem:test", correlationId: "scan-1", privacy: "private", payload: { documentId: "document:brief", path: "spec/brief.md", versionId: "version:brief:2", previousVersionId: "version:brief:1", contentType: "text/markdown", sizeBytes: 32, mtimeMs: 1, diff: "+new requirement", identityConfidence: "path", }, })).event; const output = (await store.appendEvent({ type: "stream.thought.derived.document.structure", schemaVersion: 1, source: "agent:document-structure", sourceKind: "agent", externalId: "run_visible", idempotencyKey: "derived:run_visible", occurredAt: "2026-07-14T00:00:01.000Z", actor: "document-structure", correlationId: "run_visible", parentEventId: source.id, rootEventId: source.id, privacy: "private", payload: { runId: "run_visible", agentId: "document-structure", summary: "Changed: spec/brief.md", recommendation: { target: "charter", reason: "A specification document changed at spec/brief.md", proposedAction: "Re-evaluate Charter work affected by this specification change; do not edit files automatically.", }, }, })).event; await store.upsertRun({ id: "run_visible", executionKey: "execution_visible", triggerEventId: "event_visible", agentId: "document-structure", agentVersion: 1, status: "completed", inputEventIds: [source.id], outputEventIds: [output.id], attempt: 1, provider: "deterministic", model: "deterministic", promptHash: "hash", contextManifest: {}, result: { summary: "Changed: spec/brief.md", tags: ["filesystem", "md", "specification"], importance: "high", recommendation: { target: "charter", reason: "A specification document changed at spec/brief.md", proposedAction: "Re-evaluate Charter work affected by this specification change; do not edit files automatically.", }, }, createdAt: "2026-07-14T00:00:00.000Z", completedAt: "2026-07-14T00:00:02.000Z", updatedAt: "2026-07-14T00:00:02.000Z", }); const completed = (await store.appendEvent({ type: "stream.thought.agent.run.completed", schemaVersion: 1, source: "agent:document-structure", sourceKind: "agent", externalId: "run_visible", idempotencyKey: "run_visible:completed", occurredAt: "2026-07-14T00:00:02.000Z", actor: "document-structure", correlationId: "run_visible", parentEventId: output.id, rootEventId: source.id, privacy: "private", payload: { runId: "run_visible", agentId: "document-structure", agentVersion: 1, inputEventIds: [source.id], attempt: 1, status: "completed", outputEventId: output.id, }, })).event; const cursor = (await store.appendEvent({ type: "stream.thought.connector.cursor.advanced", schemaVersion: 1, source: "filesystem:test", sourceKind: "filesystem", externalId: "cursor:filesystem:test:1", idempotencyKey: "cursor:filesystem:test:1", occurredAt: "2026-07-14T00:00:03.000Z", actor: "filesystem:test", correlationId: "scan-1", privacy: "private", payload: { cursor: { version: 1 } }, })).event;
const server = await startInspectorServer(store, { port: 0 }); servers.push(server); const address = server.address(); if (!address || typeof address === "string") throw new Error("Missing inspector address"); const base = `http://127.0.0.1:${address.port}`;
const pageText = await (await fetch(base)).text(); expect(pageText).toContain("<h3>Processing</h3>"); expect(pageText).toContain("LLM inference:"); expect(pageText.indexOf("<h3>Processing</h3>")).toBeLessThan(pageText.indexOf("<h3>Raw payload</h3>"));
const snapshot = await (await fetch(`${base}/api/snapshot`)).json() as { activity: { items: Array<{ id: string; summary: string; descendantEventCount: number; consumerRuns: Array<{ kind: string; description: string; outputCount: number }>; }> }; }; expect(snapshot.activity.items).toHaveLength(1); expect(snapshot.activity.items[0]).toMatchObject({ id: source.id, summary: "spec/brief.md changed", descendantEventCount: 2, consumerRuns: [{ kind: "rule", outputCount: 1, description: "Classified spec/brief.md as a specification change (high importance). Proposed: Re-evaluate Charter work affected by this specification change; do not edit files automatically. The proposal was not executed.", }], }); expect(snapshot.activity.items.some((item) => item.id === completed.id)).toBe(false); expect(snapshot.activity.items.some((item) => item.id === cursor.id)).toBe(false);
const detail = await (await fetch(`${base}/api/events/${encodeURIComponent(completed.id)}`)).json() as { agentActivities: Array<{ run: { result?: { summary?: string; recommendation?: { proposedAction?: string } } }; inputs: Array<{ id: string }>; outputs: Array<{ id: string; type: string }>; }>; }; expect(detail.agentActivities).toHaveLength(1); expect(detail.agentActivities[0]?.run.result?.summary).toBe("Changed: spec/brief.md"); expect(detail.agentActivities[0]?.run.result?.recommendation?.proposedAction).toContain("do not edit files automatically"); expect(detail.agentActivities[0]?.inputs.map((event) => event.id)).toEqual([source.id]); expect(detail.agentActivities[0]?.outputs).toEqual([ expect.objectContaining({ id: output.id, type: "stream.thought.derived.document.structure" }), ]); });
test("shows loaded adapter catalog selection without exposing the private checkpoint", async () => { const project = await temporaryProject(); roots.push(project); const store = await testStore(project); stores.push(store); const privateCheckpoint = "tinker://must-never-reach-inspector"; const digest = "a".repeat(64); const modelAdapter = { id: "julia-adapter", version: 1, description: "Julia transformation adapter", releasedAt: "2026-07-25T20:00:00.000Z", manifestSha256: "b".repeat(64), checkpointSelector: { kind: "env", reference: "THOUGHTSTREAM_JULIA_ADAPTER_CHECKPOINT" }, providerProfile: "tinker-default", baseModel: "Qwen/Qwen3.5-35B-A3B-Base", dataset: { id: "julia-dataset", sha256: "c".repeat(64) }, evals: [{ id: "julia-eval", sha256: "d".repeat(64) }], capabilities: ["julia-dict-transform"], privacyClass: "private", exportClass: "restricted", binding: { checkpointReferenceSha256: "e".repeat(64) }, }; await store.upsertAgent({ id: "conceptualizer", version: 1, enabled: true, spec: { id: "conceptualizer", version: 1, adapterCatalogDigest: digest, adapterCatalogGeneration: 7, modelAdapter, }, specHash: "fixture-spec-hash", updatedAt: "2026-07-25T20:00:00.000Z", });
const server = await startInspectorServer(store, { port: 0 }); servers.push(server); const address = server.address(); if (!address || typeof address === "string") throw new Error("Missing inspector address"); const snapshot = await (await fetch(`http://127.0.0.1:${address.port}/api/snapshot`)).json() as { adapterInventory: { catalogs: Array<{ digest: string; generation: number }>; selections: Array<{ id: string; version: number; enabled: boolean; adapter: { release: Record<string, unknown> } }>; }; };
expect(snapshot.adapterInventory.catalogs).toEqual([{ digest, generation: 7 }]); expect(snapshot.adapterInventory.selections).toEqual([expect.objectContaining({ id: "conceptualizer", version: 1, enabled: true, adapter: expect.objectContaining({ release: modelAdapter }), })]); const serialized = JSON.stringify(snapshot.adapterInventory); expect(serialized).not.toContain(privateCheckpoint); expect(serialized).not.toContain("/home/"); });
test("refuses a non-loopback bind", async () => { const project = await temporaryProject(); roots.push(project); const store = await testStore(project); stores.push(store); await expect(startInspectorServer(store, { host: "0.0.0.0", port: 0 })).rejects.toThrow("loopback"); });});
function sseFrames(reader: ReadableStreamDefaultReader<Uint8Array>): { next(): Promise<string> } { const decoder = new TextDecoder(); let buffer = ""; return { next: async () => { const deadline = Date.now() + 3_000; while (Date.now() < deadline) { const boundary = buffer.indexOf("\n\n"); if (boundary >= 0) { const frame = buffer.slice(0, boundary); buffer = buffer.slice(boundary + 2); return frame; } const remaining = deadline - Date.now(); const next = await Promise.race([ reader.read(), new Promise<never>((_, reject) => setTimeout(() => reject(new Error("Timed out waiting for inspector SSE frame")), remaining)), ]); if (next.done) throw new Error("Inspector SSE stream ended early"); buffer += decoder.decode(next.value, { stream: true }); } throw new Error("Timed out waiting for inspector SSE frame"); }, };}