import 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((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("Stream"); expect(pageHtml).toContain("

Stream

"); 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(''); expect(pageHtml).toContain(''); expect(pageHtml).toContain(''); expect(pageHtml).toContain(''); expect(pageHtml).toContain(''); expect(pageHtml).toContain(''); 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(/'); 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"); 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 }> }; }; 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; }; 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("

Processing

"); expect(pageText).toContain("LLM inference:"); expect(pageText.indexOf("

Processing

")).toBeLessThan(pageText.indexOf("

Raw payload

")); 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 } }>; }; }; 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): { next(): Promise } { 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((_, 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"); }, }; }