import http, { type IncomingMessage, type ServerResponse } from "node:http"; import fs from "node:fs/promises"; import path from "node:path"; import type { JazzThoughtStore } from "../jazz/store.js"; import { buildRecentRootActivity, buildRecentRootObservations, buildRecentSourceActivity, describeRunResult, presentActivityEvent, type RootActivityProjection, } from "../projections/activity.js"; import { buildSourceHealth } from "../projections/source-health.js"; import { auditRunEvidence } from "../agents/evidence.js"; import type { AgentRun } from "../store/types.js"; import type { ThoughtEvent } from "../events/types.js"; import { ReviewDecisionConflictError, projectReviewQueue, recordBrowserReviewDecision, } from "../review/review.js"; import { ReviewCapabilityVerifier } from "../review/web-capability.js"; import { getArtifactBody, getArtifactCatalog, getArtifactContent } from "../artifacts/catalog.js"; import { POST_TRAINING_COURSE } from "../courses/post-training.js"; import { COURSE_QUESTION_EVENT_TYPE, COURSE_TUTOR_AGENT_ID, CourseQuestionConflictError, CourseRevisionConflictError, appendPostTrainingCourseQuestion, courseQuestionRequestSchema, parseCourseQuestionEvent, } from "../courses/questions.js"; import { CourseChatCapabilityVerifier } from "../courses/web-capability.js"; import { requireCompletedObservationOutput } from "../agents/output-lineage.js"; import { CORRECTION_PROPOSAL_EVENT_TYPE, MEMORY_MATERIALIZATION_FAILED_EVENT_TYPE, MEMORY_MATERIALIZED_EVENT_TYPE, MEMORY_PROPOSAL_EVENT_TYPE, PROPOSAL_DECISION_EVENT_TYPE, correctionProposalPayloadSchema, memoryMaterializationFailedPayloadSchema, memoryMaterializedPayloadSchema, memoryProposalPayloadSchema, proposalDecisionPayloadSchema, } from "../agent-proposals/contracts.js"; import { materializeMemoryDecision } from "../agent-proposals/memory-materializer.js"; import { recordProposalDecision } from "../agent-proposals/review.js"; import { z } from "zod"; import { FONT_DEBUG_ASSET_PATH, FONT_DEBUG_SCRIPT } from "./font-debug.js"; import { createDocumentRequestSchema, decideProposalRequestSchema, requestProposalRequestSchema, saveEditRequestSchema, selectContextRequestSchema, } from "../workbench/contracts.js"; import type { WorkbenchRunnerRegistry } from "../workbench/runners.js"; import { StaleBaseError, WorkbenchConflictError, WorkbenchNotFoundError, WorkbenchRunnerError, WorkbenchWorkflow, } from "../workbench/workflow.js"; const RUN_TERMINAL_EVENT_TYPES = [ "stream.thought.agent.run.completed", "stream.thought.agent.run.failed", "stream.thought.agent.run.blocked", "stream.thought.agent.run.abandoned", ]; const BLUESKY_POST_URI_PATTERN = /^at:\/\/[^/]{3,256}\/app\.bsky\.feed\.post\/[a-zA-Z0-9._~:-]{1,256}$/; const SAFE_CID_PATTERN = /^[a-zA-Z0-9]{8,128}$/; const BLUESKY_INLAY_CACHE_TTL_MS = 10 * 60_000; const BLUESKY_MEDIA_CACHE_TTL_MS = 10 * 60_000; const BLUESKY_MEDIA_CACHE_MAX_BYTES = 24 * 1024 * 1024; const BLUESKY_MEDIA_CACHE_MAX_ENTRIES = 96; const BLUESKY_MEDIA_MAX_BYTES = 7 * 1024 * 1024; const BLUESKY_MEDIA_CONTENT_TYPES = new Set(["image/jpeg", "image/png", "image/webp", "image/gif"]); const STREAM_APP_MANIFEST = `${JSON.stringify({ name: "Stream", short_name: "Stream", id: "./", start_url: "./", scope: "./", display: "standalone", background_color: "#0a0a0a", theme_color: "#0a0a0a", icons: [{ src: "assets/stream-icon.svg", sizes: "any", type: "image/svg+xml", purpose: "any maskable" }], })}\n`; const STREAM_APP_ICON = `S`; const STREAM_SERVICE_WORKER = "self.addEventListener('install',()=>self.skipWaiting());self.addEventListener('activate',event=>event.waitUntil(self.clients.claim()));\n"; interface BlueskyPostInlay { uri: string; cid: string; url: string; author: { displayName: string; handle: string; avatar?: string; }; text: string; createdAt?: string; images: Array<{ thumb: string; fullsize: string; alt: string }>; external?: { uri: string; title: string; description: string; thumb?: string }; quote?: { uri: string; url: string; author: { displayName: string; handle: string; avatar?: string }; text: string; }; counts: { replies: number; reposts: number; likes: number }; } const blueskyPostViewSchema = z.object({ uri: z.string(), cid: z.string(), author: z.object({ handle: z.string(), displayName: z.string().optional(), avatar: z.string().url().optional(), }).passthrough(), record: z.object({ text: z.string().optional(), createdAt: z.string().optional(), }).passthrough(), embed: z.unknown().optional(), replyCount: z.number().int().nonnegative().optional(), repostCount: z.number().int().nonnegative().optional(), likeCount: z.number().int().nonnegative().optional(), }).passthrough(); const blueskyPostResponseSchema = z.object({ posts: z.array(blueskyPostViewSchema) }); const blueskyInlayCache = new Map(); const blueskyMediaCache = new Map(); let blueskyMediaCacheBytes = 0; let aroundFontCache: Buffer | undefined; export interface InspectorServerOptions { host?: string; port?: number; reviewCapability?: Buffer | undefined; reviewVerifier?: ReviewCapabilityVerifier | undefined; courseChatCapability?: Buffer | undefined; courseChatVerifier?: CourseChatCapabilityVerifier | undefined; agentContextRoot?: string | undefined; proposalActor?: string | undefined; workbenchRunners?: WorkbenchRunnerRegistry | undefined; } interface InspectorActivityCache { observations(fresh?: boolean): Promise; evidence(fresh?: boolean): Promise; } function createInspectorActivityCache(store: JazzThoughtStore): InspectorActivityCache { let observationValue: RootActivityProjection | undefined; let evidenceValue: RootActivityProjection | undefined; let observationUpdatedAt = 0; let evidenceUpdatedAt = 0; let observationBuild: Promise | undefined; let evidenceBuild: Promise | undefined; const stillFresh = (updatedAt: number) => Date.now() - updatedAt < 5_000; const observations = async (fresh = false): Promise => { if (observationValue && (!fresh || stillFresh(observationUpdatedAt))) return observationValue; if (observationBuild) return observationBuild; observationBuild = buildRecentRootObservations(store).then((value) => { observationValue = value; observationUpdatedAt = Date.now(); return value; }).finally(() => { observationBuild = undefined; }); return observationBuild; }; const evidence = async (fresh = false): Promise => { if (evidenceValue && (!fresh || stillFresh(evidenceUpdatedAt))) return evidenceValue; if (evidenceBuild) return evidenceBuild; evidenceBuild = buildInspectorActivity(store).then((value) => { evidenceValue = value; evidenceUpdatedAt = Date.now(); return value; }).finally(() => { evidenceBuild = undefined; }); return evidenceBuild; }; void observations(true).then(() => evidence(true)).catch(() => undefined); return { observations, evidence }; } export async function startInspectorServer( store: JazzThoughtStore, options: InspectorServerOptions = {}, ): Promise { const host = options.host ?? "127.0.0.1"; if (!isLoopback(host)) throw new Error("thought stream inspector may only bind to a loopback address"); const port = options.port ?? 4317; const reviewVerifier = options.reviewVerifier ?? (options.reviewCapability ? new ReviewCapabilityVerifier(options.reviewCapability) : undefined); const courseChatVerifier = options.courseChatVerifier ?? (options.courseChatCapability ? new CourseChatCapabilityVerifier(options.courseChatCapability) : undefined); const proposalWritesEnabled = Boolean(reviewVerifier && options.agentContextRoot); const activityCache = createInspectorActivityCache(store); const workbench = new WorkbenchWorkflow(store, { actor: options.proposalActor ?? "operator:cameron", ...(options.workbenchRunners ? { runners: options.workbenchRunners } : {}), }); const server = http.createServer((request, response) => { void handleRequest(store, activityCache, request, response, reviewVerifier, courseChatVerifier, proposalWritesEnabled ? options.agentContextRoot : undefined, options.proposalActor ?? "operator:cameron", workbench).catch((error) => { sendJson(response, 500, { error: error instanceof Error ? error.message : String(error) }); }); }); await new Promise((resolve, reject) => { server.once("error", reject); server.listen(port, host, () => { server.off("error", reject); resolve(); }); }); return server; } async function handleRequest( store: JazzThoughtStore, activityCache: InspectorActivityCache, request: IncomingMessage, response: ServerResponse, reviewVerifier?: ReviewCapabilityVerifier, courseChatVerifier?: CourseChatCapabilityVerifier, agentContextRoot?: string, proposalActor = "operator:cameron", workbench: WorkbenchWorkflow = new WorkbenchWorkflow(store, { actor: proposalActor }), ): Promise { const url = new URL(request.url ?? "/", "http://127.0.0.1"); if (request.method === "POST" && url.pathname.startsWith("/api/workbench/")) { await handleWorkbenchWrite(workbench, request, response, url, reviewVerifier); return; } if (request.method === "POST" && url.pathname === "/api/courses/post-training/questions" && url.search === "") { if (!courseChatVerifier) { request.resume(); sendJson(response, 405, { error: "Course questions are read-only" }); return; } let body: Buffer; try { body = await readBody(request, 4_096); } catch { request.resume(); sendJson(response, 400, { error: "Course question is invalid" }); return; } if (!courseChatVerifier.verify(request.headers, "POST", url.pathname, body)) { sendJson(response, 403, { error: "Course question could not be verified" }); return; } try { const input = courseQuestionRequestSchema.parse(JSON.parse(body.toString("utf8"))); const appended = await appendPostTrainingCourseQuestion(store, input, { actor: proposalActor }); sendJson(response, 202, { eventId: appended.event.id, inserted: appended.inserted, statusPath: `api/courses/post-training/questions/${encodeURIComponent(appended.event.id)}`, }); } catch (error) { if (error instanceof CourseRevisionConflictError) { sendJson(response, 409, { error: "The course changed; reload before asking" }); } else if (error instanceof CourseQuestionConflictError) { sendJson(response, 409, { error: "Question identity changed; create a new question" }); } else { sendJson(response, 400, { error: "Course question is invalid" }); } } return; } const proposalDecisionMatch = url.pathname.match(/^\/api\/proposals\/([^/]+)\/decisions$/); if (request.method === "POST" && proposalDecisionMatch) { if (!reviewVerifier || !agentContextRoot) { request.resume(); sendJson(response, 405, { error: "Proposal review is read-only" }); return; } let body: Buffer; try { body = await readBody(request, 34_816); } catch { request.resume(); sendJson(response, 400, { error: "Proposal decision is invalid" }); return; } if (!reviewVerifier.verify(request.headers, "POST", url.pathname, body)) { sendJson(response, 403, { error: "Proposal decision could not be verified" }); return; } try { const eventId = decodeURIComponent(proposalDecisionMatch[1]!); const input = proposalDecisionRequestSchema.parse(JSON.parse(body.toString("utf8"))); const result = await recordProposalDecision(store, { proposalEventId: eventId, ...input, actor: proposalActor }); const materialization = result.proposal.type === MEMORY_PROPOSAL_EVENT_TYPE && input.disposition !== "reject" ? await materializeMemoryDecision(store, result.decision.id, { contextRoot: agentContextRoot }) : undefined; sendJson(response, 201, { decisionEventId: result.decision.id, disposition: input.disposition, ...(result.projectedJudgment ? { judgmentEventId: result.projectedJudgment.id, receipt: "Correction accepted as a private quality judgment" } : {}), ...(materialization ? { materializationEventId: materialization.event.id, materializationStatus: materialization.status, receipt: materialization.status === "materialized" ? "Stream memory updated" : "Stream memory was not changed" } : {}), }); } catch (error) { const message = error instanceof Error ? error.message : "Proposal decision is invalid"; if (message.includes("already has a human decision") || message.includes("submission id conflicts")) { sendJson(response, 409, { error: "Proposal decision changed; reload before submitting" }); } else { sendJson(response, 400, { error: "Proposal decision is invalid" }); } } return; } const reviewDecisionMatch = url.pathname.match(/^\/api\/reviews\/([^/]+)\/decisions$/); if (request.method === "POST" && reviewDecisionMatch) { if (!reviewVerifier) { request.resume(); sendJson(response, 405, { error: "Inspector is read-only" }); return; } let body: Buffer; try { body = await readBody(request, 98_304); } catch { request.resume(); sendJson(response, 400, { error: "Review request is invalid" }); return; } if (!reviewVerifier.verify(request.headers, "POST", url.pathname, body)) { sendJson(response, 403, { error: "Review request could not be verified" }); return; } try { const itemId = decodeURIComponent(reviewDecisionMatch[1]!); const decision = await recordBrowserReviewDecision(store, itemId, JSON.parse(body.toString("utf8"))); sendJson(response, 201, { decisionEventId: decision.id, active: true }); } catch (error) { if (error instanceof ReviewDecisionConflictError) { sendJson(response, 409, { error: "Review decision changed; reload before submitting" }); } else { sendJson(response, 400, { error: "Review decision is invalid" }); } } return; } if (request.method !== "GET" && request.method !== "HEAD") { sendJson(response, 405, { error: "Inspector is read-only" }); return; } if (url.pathname === "/") { send(response, 200, "text/html; charset=utf-8", renderInspectorHtml()); return; } if (url.pathname === "/manifest.webmanifest") { send(response, 200, "application/manifest+json; charset=utf-8", STREAM_APP_MANIFEST); return; } if (url.pathname === "/assets/stream-icon.svg") { send(response, 200, "image/svg+xml; charset=utf-8", STREAM_APP_ICON); return; } if (url.pathname === FONT_DEBUG_ASSET_PATH) { send(response, 200, "text/javascript; charset=utf-8", FONT_DEBUG_SCRIPT); return; } if (url.pathname === "/sw.js") { send(response, 200, "text/javascript; charset=utf-8", STREAM_SERVICE_WORKER); return; } if (url.pathname === "/assets/around-regular.woff2") { try { aroundFontCache ??= await loadAroundFont(); sendBytes(response, 200, "font/woff2", aroundFontCache, "public, max-age=31536000, immutable"); } catch { sendJson(response, 502, { error: "Display font is unavailable" }); } return; } if (url.pathname === "/api/inlays/bluesky-media") { const mediaReferences = url.searchParams.getAll("url"); const source = mediaReferences.length === 1 && [...url.searchParams.keys()].every((key) => key === "url") ? safeBlueskyMediaUrl(mediaReferences[0]) : undefined; if (!source) { sendJson(response, 400, { error: "Bluesky media reference is invalid" }); return; } try { const media = await loadBlueskyMedia(source); sendBytes(response, 200, media.contentType, media.bytes, "private, max-age=600"); } catch { sendJson(response, 502, { error: "Bluesky media is unavailable" }); } return; } if (url.pathname === "/api/inlays/bluesky") { const uri = url.searchParams.get("uri"); const cid = url.searchParams.get("cid") ?? undefined; if (!uri || !BLUESKY_POST_URI_PATTERN.test(uri) || (cid !== undefined && !SAFE_CID_PATTERN.test(cid))) { sendJson(response, 400, { error: "Bluesky post reference is invalid" }); return; } try { sendJson(response, 200, await loadBlueskyPostInlay(uri, cid)); } catch { sendJson(response, 502, { error: "Bluesky post is unavailable" }); } return; } if (url.pathname === "/api/activity") { sendJson(response, 200, await activityCache.observations(url.searchParams.get("fresh") === "1")); return; } if (url.pathname === "/api/activity/evidence") { sendJson(response, 200, await activityCache.evidence(url.searchParams.get("fresh") === "1")); return; } if (url.pathname === "/api/system") { sendJson(response, 200, await buildInspectorSystem(store)); return; } if (url.pathname === "/api/snapshot") { const runs = await store.listRuns(); const [activity, system] = await Promise.all([ buildInspectorActivity(store), buildInspectorSystem(store, runs), ]); sendJson(response, 200, { activity, ...system }); return; } if (url.pathname === "/api/live") { streamInspectorActivity(store, request, response); return; } if (url.pathname === "/api/proposals") { sendJson(response, 200, await projectProposalQueue(store)); return; } if (url.pathname === "/api/reviews") { sendJson(response, 200, await projectReviewQueue(store)); return; } if (url.pathname === "/api/session") { sendJson(response, 200, { reviewWriteEnabled: false, courseChatEnabled: false }); return; } if (url.pathname === "/api/courses/post-training") { sendJson(response, 200, POST_TRAINING_COURSE); return; } const courseQuestionMatch = url.pathname.match(/^\/api\/courses\/post-training\/questions\/([^/]+)$/); if (courseQuestionMatch) { const eventId = decodeURIComponent(courseQuestionMatch[1]!); sendJson(response, 200, await postTrainingQuestionStatus(store, eventId)); return; } if (url.pathname === "/api/artifacts") { sendJson(response, 200, await getArtifactCatalog(store)); return; } if (url.pathname === "/api/workbench/documents") { sendJson(response, 200, { items: await workbench.listDocuments(), runners: workbench.listRunners() }); return; } if (url.pathname === "/api/workbench/candidates") { const documentId = url.searchParams.get("documentId"); if (!documentId || [...url.searchParams.keys()].some((key) => key !== "documentId")) { return sendJson(response, 400, { error: "Candidate lookup requires exactly one documentId" }); } const candidates = await workbench.listCandidates(documentId); if (!candidates) return sendJson(response, 404, { error: "Working document not found" }); sendJson(response, 200, candidates); return; } const workbenchDocumentMatch = url.pathname.match(/^\/api\/workbench\/documents\/([^/]+)$/); if (workbenchDocumentMatch) { const documentId = decodeURIComponent(workbenchDocumentMatch[1]!); const detail = await workbench.getDocument(documentId); if (!detail) return sendJson(response, 404, { error: "Working document not found" }); sendJson(response, 200, detail); return; } const sourceMatch = url.pathname.match(/^\/api\/sources\/([^/]+)$/); if (sourceMatch) { const sourceId = decodeURIComponent(sourceMatch[1]!); const [runs, sources] = await Promise.all([store.listRuns(), buildSourceHealth(store)]); const source = sources.find((candidate) => candidate.source === sourceId); if (!source) return sendJson(response, 404, { error: `Source not found: ${sourceId}` }); sendJson(response, 200, { source, activity: await buildRecentSourceActivity(store, runs, sourceId) }); return; } const artifactContentMatch = url.pathname.match(/^\/api\/artifacts\/([^/]+)\/content$/); if (artifactContentMatch) { const id = decodeURIComponent(artifactContentMatch[1]!); const detail = await getArtifactContent(store, id); if (!detail) return sendJson(response, 404, { error: `Artifact not found: ${id}` }); if (!detail.payload.mediaType.startsWith("image/")) return sendJson(response, 404, { error: "Artifact has no binary inspection route" }); response.writeHead(200, { "content-type": detail.payload.mediaType, "content-length": detail.bytes.byteLength, "x-content-type-options": "nosniff", "cache-control": "private, no-store" }); response.end(detail.bytes); return; } const artifactMatch = url.pathname.match(/^\/api\/artifacts\/([^/]+)$/); if (artifactMatch) { const id = decodeURIComponent(artifactMatch[1]!); const detail = await getArtifactBody(store, id); if (!detail) return sendJson(response, 404, { error: `Artifact not found: ${id}` }); const catalog = await getArtifactCatalog(store); const entry = catalog.entries.find((candidate) => candidate.eventId === id); sendJson(response, 200, { event: detail.event, ...(detail.text !== undefined ? { text: detail.text } : {}), ...(detail.text !== undefined && detail.mediaType === "text/markdown" ? { renderedHtml: renderArtifactMarkdown(detail.text) } : {}), ...(detail.contentPath ? { contentPath: detail.contentPath } : {}), bodySha256: detail.bodySha256, byteCount: detail.byteCount, mediaType: detail.mediaType, catalogEntry: entry, }); return; } const eventMatch = url.pathname.match(/^\/api\/events\/([^/]+)$/); if (eventMatch) { const id = decodeURIComponent(eventMatch[1]!); const event = await store.getEvent(id); if (!event) return sendJson(response, 404, { error: `Event not found: ${id}` }); const [children, runs, parent, root] = await Promise.all([ store.listChildEvents(event.id), store.listRuns(), event.parentEventId ? store.getEvent(event.parentEventId) : undefined, event.rootEventId === event.id ? event : store.getEvent(event.rootEventId), ]); const runInputs = await store.getEvents([...new Set(runs.flatMap((run) => run.inputEventIds))]); const runInputsById = new Map(runInputs.map((input) => [input.id, input])); const relevantRuns = runs.filter((run) => eventBelongsToRun(event, run, runInputsById)); sendJson(response, 200, { event, presentation: presentActivityEvent(event), parent, children, root, agentActivities: await Promise.all(relevantRuns.map((run) => expandRun(store, run))), }); return; } const runMatch = url.pathname.match(/^\/api\/runs\/([^/]+)$/); if (runMatch) { const id = decodeURIComponent(runMatch[1]!); const run = await store.getRun(id); if (!run) return sendJson(response, 404, { error: `Run not found: ${id}` }); const [activity, evidenceEvents] = await Promise.all([ expandRun(store, run), loadRunEvidenceEvents(store, [run]), ]); const evidence = auditRunEvidence([run], evidenceEvents)[0]; sendJson(response, 200, { ...activity, evidence }); return; } sendJson(response, 404, { error: "Not found" }); } const WORKBENCH_WRITE_BODY_LIMIT = 98_304; /** * Trusted workbench mutations. Every route requires the same body-bound Review * capability as the proposal decision route; without a verifier the inspector * stays read-only (405). Request bodies are strict zod schemas and every * mutation carries a client-generated request/submission id so retries after * an uncertain outcome converge instead of duplicating effects. */ async function handleWorkbenchWrite( workbench: WorkbenchWorkflow, request: IncomingMessage, response: ServerResponse, url: URL, reviewVerifier?: ReviewCapabilityVerifier, ): Promise { const createMatch = url.pathname === "/api/workbench/documents"; const versionMatch = url.pathname.match(/^\/api\/workbench\/documents\/([^/]+)\/versions$/); const selectionMatch = url.pathname.match(/^\/api\/workbench\/documents\/([^/]+)\/selections$/); const proposalMatch = url.pathname.match(/^\/api\/workbench\/documents\/([^/]+)\/proposals$/); const decisionMatch = url.pathname.match(/^\/api\/workbench\/proposals\/([^/]+)\/decisions$/); if (url.search !== "" || (!createMatch && !versionMatch && !selectionMatch && !proposalMatch && !decisionMatch)) { request.resume(); sendJson(response, 404, { error: "Not found" }); return; } if (!reviewVerifier) { request.resume(); sendJson(response, 405, { error: "Workbench is read-only" }); return; } let body: Buffer; try { body = await readBody(request, WORKBENCH_WRITE_BODY_LIMIT); } catch { request.resume(); sendJson(response, 400, { error: "Workbench request is invalid" }); return; } if (!reviewVerifier.verify(request.headers, "POST", url.pathname, body)) { sendJson(response, 403, { error: "Workbench request could not be verified" }); return; } try { const parsed: unknown = JSON.parse(body.toString("utf8")); if (createMatch) { const input = createDocumentRequestSchema.parse(parsed); const head = await workbench.createWorkingDocument(input); sendJson(response, 201, { documentId: head.documentId, versionId: head.versionId, sha256: head.sha256 }); return; } if (versionMatch) { const input = saveEditRequestSchema.parse(parsed); const head = await workbench.saveOperatorEdit({ documentId: decodeURIComponent(versionMatch[1]!), ...input }); sendJson(response, 201, { documentId: head.documentId, versionId: head.versionId, sha256: head.sha256 }); return; } if (selectionMatch) { const input = selectContextRequestSchema.parse(parsed); const selection = await workbench.selectContext({ documentId: decodeURIComponent(selectionMatch[1]!), ...input }); sendJson(response, 201, { selectionId: selection.selectionId, snapshotVersionId: selection.snapshotVersionId, eventId: selection.eventId }); return; } if (proposalMatch) { const input = requestProposalRequestSchema.parse(parsed); const result = await workbench.requestProposal({ documentId: decodeURIComponent(proposalMatch[1]!), ...input }); sendJson(response, 201, result); return; } const input = decideProposalRequestSchema.parse(parsed); const result = await workbench.decideProposal({ proposalEventId: decodeURIComponent(decisionMatch![1]!), ...input }); sendJson(response, 201, result); } catch (error) { if (error instanceof StaleBaseError) { sendJson(response, 409, { error: "This document changed since you opened it", stale: true, headVersionId: error.headVersionId }); } else if (error instanceof WorkbenchNotFoundError) { sendJson(response, 404, { error: error.message }); } else if (error instanceof WorkbenchConflictError) { sendJson(response, 409, { error: error.message }); } else if (error instanceof WorkbenchRunnerError) { sendJson(response, 502, { error: "The proposal runner failed; no proposal was recorded", runId: error.runId }); } else if (error instanceof z.ZodError) { sendJson(response, 400, { error: "Workbench request is invalid" }); } else { throw error; } } } async function postTrainingQuestionStatus(store: JazzThoughtStore, eventId: string): Promise> { const event = await store.getEvent(eventId); if (!event || event.type !== COURSE_QUESTION_EVENT_TYPE || !parseCourseQuestionEvent(event)) { return { status: "not-found" }; } const runs = (await store.getRunsForTriggerEvents([event.id])) .filter((run) => run.agentId === COURSE_TUTOR_AGENT_ID) .sort((left, right) => right.attempt - left.attempt || right.createdAt.localeCompare(left.createdAt)); const run = runs[0]; if (!run) return { status: "pending", eventId: event.id }; if (!["completed", "failed", "blocked", "abandoned", "skipped"].includes(run.status)) { return { status: "running", eventId: event.id, runId: run.id }; } if (run.status !== "completed" || run.outputEventIds.length !== 1) { return { status: "failed", eventId: event.id, runId: run.id }; } const output = await store.getEvent(run.outputEventIds[0]!); if (!output) return { status: "failed", eventId: event.id, runId: run.id }; try { const lineage = await requireCompletedObservationOutput(store, output); if (lineage.run.id !== run.id || lineage.run.agentId !== COURSE_TUTOR_AGENT_ID || lineage.trigger.id !== event.id) { return { status: "failed", eventId: event.id, runId: run.id }; } } catch { return { status: "failed", eventId: event.id, runId: run.id }; } const answer = output.payload.summary; if (typeof answer !== "string" || answer.length === 0 || answer.length > 2_000) { return { status: "failed", eventId: event.id, runId: run.id }; } return { status: "completed", eventId: event.id, runId: run.id, outputEventId: output.id, answer, }; } const proposalDecisionRequestSchema = z.object({ disposition: z.enum(["accept", "edit", "reject"]), replacementText: z.string().min(1).max(32_768).optional(), submissionId: z.string().min(1).max(500), }).strict().superRefine((value, context) => { if ((value.disposition === "edit") !== (value.replacementText !== undefined)) { context.addIssue({ code: "custom", path: ["replacementText"], message: "Only edit decisions require replacement text" }); } }); export interface ProposalQueueItem { eventId: string; kind: "Memory suggestion" | "Proposed correction"; proposedText: string; reason: string; operation: string; originalText?: string | undefined; observedAt: string; status: "Awaiting review" | "Accepted" | "Edited and accepted" | "Rejected" | "Materialized" | "Failed"; decisionEventId?: string | undefined; receiptEventId?: string | undefined; failureReason?: string | undefined; technical: Record; } export async function projectProposalQueue(store: JazzThoughtStore): Promise<{ count: number; items: ProposalQueueItem[] }> { const [proposals, decisions, receipts] = await Promise.all([ store.listEvents({ types: [MEMORY_PROPOSAL_EVENT_TYPE, CORRECTION_PROPOSAL_EVENT_TYPE] }), store.listEvents({ types: [PROPOSAL_DECISION_EVENT_TYPE] }), store.listEvents({ types: [MEMORY_MATERIALIZED_EVENT_TYPE, MEMORY_MATERIALIZATION_FAILED_EVENT_TYPE] }), ]); const items: ProposalQueueItem[] = []; for (const proposal of proposals) { const decision = decisions.find((event) => event.payload.proposalEventId === proposal.id); const receipt = decision && receipts.find((event) => event.payload.decisionEventId === decision.id); const decisionPayload = decision ? proposalDecisionPayloadSchema.parse(decision.payload) : undefined; let status: ProposalQueueItem["status"] = "Awaiting review"; if (receipt?.type === MEMORY_MATERIALIZED_EVENT_TYPE) status = "Materialized"; else if (receipt?.type === MEMORY_MATERIALIZATION_FAILED_EVENT_TYPE) status = "Failed"; else if (decisionPayload?.disposition === "reject") status = "Rejected"; else if (decisionPayload?.disposition === "edit") status = "Edited and accepted"; else if (decisionPayload?.disposition === "accept") status = "Accepted"; if (proposal.type === MEMORY_PROPOSAL_EVENT_TYPE) { const payload = memoryProposalPayloadSchema.parse(proposal.payload); const failed = receipt?.type === MEMORY_MATERIALIZATION_FAILED_EVENT_TYPE ? memoryMaterializationFailedPayloadSchema.parse(receipt.payload) : undefined; if (receipt?.type === MEMORY_MATERIALIZED_EVENT_TYPE) memoryMaterializedPayloadSchema.parse(receipt.payload); items.push({ eventId: proposal.id, kind: "Memory suggestion", proposedText: decisionPayload?.replacementText ?? payload.proposedText, reason: payload.reason, operation: payload.operation === "append" ? "Append to Stream memory" : "Replace Stream memory", observedAt: proposal.observedAt, status, ...(decision ? { decisionEventId: decision.id } : {}), ...(receipt ? { receiptEventId: receipt.id } : {}), ...(failed ? { failureReason: humanMaterializationFailure(failed.reasonCode) } : {}), technical: { proposalEventId: proposal.id, agentId: payload.proposer.agentId, agentVersion: payload.proposer.agentVersion, runId: payload.proposer.runId, contextSnapshotId: payload.proposer.contextSnapshotId, targetDocumentId: payload.target.documentId, targetVersionId: payload.target.versionId, targetSha256: payload.target.sha256, proposedTextSha256: payload.proposedTextSha256, }, }); continue; } const payload = correctionProposalPayloadSchema.parse(proposal.payload); const target = await store.getEvent(payload.target.outputEventId); const originalText = typeof target?.payload.structuredOutput === "object" && target.payload.structuredOutput && !Array.isArray(target.payload.structuredOutput) && typeof (target.payload.structuredOutput as Record).summary === "string" ? String((target.payload.structuredOutput as Record).summary) : typeof target?.payload.summary === "string" ? target.payload.summary : "Original delivered reply unavailable"; items.push({ eventId: proposal.id, kind: "Proposed correction", proposedText: decisionPayload?.replacementText ?? payload.replacementText, reason: payload.reason, operation: "Correct the delivered reply", originalText, observedAt: proposal.observedAt, status, ...(decision ? { decisionEventId: decision.id } : {}), technical: { proposalEventId: proposal.id, agentId: payload.proposer.agentId, agentVersion: payload.proposer.agentVersion, runId: payload.proposer.runId, targetRunId: payload.target.runId, targetOutputEventId: payload.target.outputEventId, deliveryReceiptEventId: payload.target.deliveryReceiptEventId, outputContractId: payload.target.outputContract.id, outputContractVersion: payload.target.outputContract.version, outputContractSha256: payload.target.outputContract.sha256, replacementTextSha256: payload.replacementTextSha256, }, }); } items.sort((left, right) => right.observedAt.localeCompare(left.observedAt) || right.eventId.localeCompare(left.eventId)); return { count: items.length, items }; } function humanMaterializationFailure(code: string): string { return ({ "stale-base": "Stream memory changed after this suggestion was created", "symlink-refused": "The memory target did not pass the safe-file check", "frontmatter-invalid": "The approved text did not preserve Stream memory identity", "document-identity-changed": "The approved text changed Stream memory identity", "write-failed": "The trusted memory update could not be completed", } as Record)[code] ?? "The trusted memory update failed its safety checks"; } function adapterSelectionFromSpec(spec: Record): { digest: string; generation: number; release: Record } | undefined { const digest = spec.adapterCatalogDigest; const generation = spec.adapterCatalogGeneration; const release = spec.modelAdapter; if (typeof digest !== "string" || !/^[a-f0-9]{64}$/.test(digest) || !Number.isSafeInteger(generation) || Number(generation) < 1) return undefined; if (!release || typeof release !== "object" || Array.isArray(release)) return undefined; return { digest, generation: Number(generation), release: release as Record }; } export function renderArtifactMarkdown(markdown: string): string { let text = markdown.replace(/\r\n?/g, "\n"); if (text.startsWith("---\n")) { const frontmatterEnd = text.indexOf("\n---\n", 4); if (frontmatterEnd >= 0) text = text.slice(frontmatterEnd + 5); } const lines = text.split("\n"); let html = ""; let paragraph: string[] = []; let list: "ul" | "ol" | undefined; let code = false; let codeLines: string[] = []; const flushParagraph = (): void => { if (paragraph.length === 0) return; html += `

${renderInlineMarkdown(paragraph.join(" "))}

`; paragraph = []; }; const closeList = (): void => { if (!list) return; html += ``; list = undefined; }; const flushCode = (): void => { html += `
${escapeArtifactHtml(codeLines.join("\n"))}
`; codeLines = []; }; for (const line of lines) { if (line.startsWith("```")) { flushParagraph(); closeList(); if (code) flushCode(); code = !code; continue; } if (code) { codeLines.push(line); continue; } const heading = line.match(/^(#{1,6})\s+(.+)$/); const unordered = line.match(/^\s*[-*+]\s+(.+)$/); const ordered = line.match(/^\s*\d+[.)]\s+(.+)$/); const quote = line.match(/^>\s?(.*)$/); if (heading) { flushParagraph(); closeList(); const level = heading[1]!.length; html += `${renderInlineMarkdown(heading[2]!)}`; continue; } if (unordered || ordered) { flushParagraph(); const next = unordered ? "ul" : "ol"; if (list !== next) { closeList(); list = next; html += `<${list}>`; } html += `
  • ${renderInlineMarkdown((unordered ?? ordered)![1]!)}
  • `; continue; } if (quote) { flushParagraph(); closeList(); html += `
    ${renderInlineMarkdown(quote[1]!)}
    `; continue; } if (!line.trim()) { flushParagraph(); closeList(); continue; } paragraph.push(line.trim()); } flushParagraph(); closeList(); if (code) flushCode(); return html; } function renderInlineMarkdown(value: string): string { return escapeArtifactHtml(value) .replace(/`([^`]+)`/g, "$1") .replace(/\*\*\*([^*]+)\*\*\*/g, "$1") .replace(/\*\*([^*]+)\*\*/g, "$1") .replace(/\*([^*]+)\*/g, "$1") .replace(/\[([^\]]+)\]\((https?:\/\/[^\s)]+)\)/g, '$1'); } function escapeArtifactHtml(value: string): string { return value.replace(/[&<>"']/g, (character) => ({ "&": "&", "<": "<", ">": ">", '"': """, "'": "'", })[character]!); } async function loadBlueskyPostInlay(uri: string, expectedCid?: string): Promise { const cacheKey = `${uri}\u0000${expectedCid ?? "current"}`; const cached = blueskyInlayCache.get(cacheKey); if (cached && cached.expiresAt > Date.now()) return cached.value; blueskyInlayCache.delete(cacheKey); const endpoint = new URL("https://public.api.bsky.app/xrpc/app.bsky.feed.getPosts"); endpoint.searchParams.append("uris", uri); const response = await fetch(endpoint, { headers: { accept: "application/json" }, signal: AbortSignal.timeout(4_000), }); if (!response.ok) throw new Error(`Bluesky AppView returned ${response.status}`); const parsed = blueskyPostResponseSchema.parse(await response.json()); const post = parsed.posts.find((candidate) => candidate.uri === uri); if (!post || (expectedCid && post.cid !== expectedCid)) throw new Error("Bluesky post did not match the observed strong reference"); const embed = blueskyEmbed(post.embed); const value: BlueskyPostInlay = { uri: post.uri, cid: post.cid, url: atUriToBlueskyUrl(post.uri), author: blueskyAuthor(post.author), text: post.record.text ?? "", ...(post.record.createdAt ? { createdAt: post.record.createdAt } : {}), images: embed.images, ...(embed.external ? { external: embed.external } : {}), ...(embed.quote ? { quote: embed.quote } : {}), counts: { replies: post.replyCount ?? 0, reposts: post.repostCount ?? 0, likes: post.likeCount ?? 0, }, }; blueskyInlayCache.set(cacheKey, { expiresAt: Date.now() + BLUESKY_INLAY_CACHE_TTL_MS, value }); while (blueskyInlayCache.size > 512) { const oldest = blueskyInlayCache.keys().next().value as string | undefined; if (!oldest) break; blueskyInlayCache.delete(oldest); } return value; } async function loadAroundFont(): Promise { const fontPath = path.resolve(process.cwd(), "public/assets/around-regular.woff2"); const stat = await fs.lstat(fontPath); if (!stat.isFile() || stat.isSymbolicLink()) throw new Error("Around font is not a regular file"); const bytes = await fs.readFile(fontPath); if (bytes.length < 4 || bytes.length > 65_536 || bytes.subarray(0, 4).toString("ascii") !== "wOF2") { throw new Error("Around font is invalid"); } return bytes; } async function loadBlueskyMedia(source: string): Promise<{ contentType: string; bytes: Buffer }> { const cached = blueskyMediaCache.get(source); if (cached && cached.expiresAt > Date.now()) return cached; if (cached) removeBlueskyMediaCacheEntry(source, cached); const response = await fetch(source, { headers: { accept: "image/webp,image/png,image/jpeg,image/gif" }, redirect: "error", signal: AbortSignal.timeout(5_000), }); if (!response.ok || !response.body) throw new Error("Bluesky CDN request failed"); const contentType = (response.headers.get("content-type") ?? "").split(";", 1)[0]!.trim().toLowerCase(); if (!BLUESKY_MEDIA_CONTENT_TYPES.has(contentType)) { await response.body.cancel(); throw new Error("Bluesky CDN media type is unsupported"); } const declaredLength = response.headers.get("content-length"); if (declaredLength && (!/^\d+$/.test(declaredLength) || Number(declaredLength) > BLUESKY_MEDIA_MAX_BYTES)) { await response.body.cancel(); throw new Error("Bluesky CDN media is oversized"); } const reader = response.body.getReader(); const chunks: Buffer[] = []; let totalBytes = 0; while (true) { const chunk = await reader.read(); if (chunk.done) break; const bytes = Buffer.from(chunk.value); totalBytes += bytes.length; if (totalBytes > BLUESKY_MEDIA_MAX_BYTES) { await reader.cancel(); throw new Error("Bluesky CDN media is oversized"); } chunks.push(bytes); } const bytes = Buffer.concat(chunks, totalBytes); if (!matchesImageMagic(bytes, contentType)) throw new Error("Bluesky CDN media bytes are invalid"); const entry = { expiresAt: Date.now() + BLUESKY_MEDIA_CACHE_TTL_MS, contentType, bytes }; blueskyMediaCache.set(source, entry); blueskyMediaCacheBytes += bytes.length; trimBlueskyMediaCache(); return entry; } function trimBlueskyMediaCache(): void { const now = Date.now(); for (const [key, entry] of blueskyMediaCache) { if (entry.expiresAt <= now) removeBlueskyMediaCacheEntry(key, entry); } while (blueskyMediaCache.size > BLUESKY_MEDIA_CACHE_MAX_ENTRIES || blueskyMediaCacheBytes > BLUESKY_MEDIA_CACHE_MAX_BYTES) { const oldest = blueskyMediaCache.entries().next().value as [string, { expiresAt: number; contentType: string; bytes: Buffer }] | undefined; if (!oldest) break; removeBlueskyMediaCacheEntry(oldest[0], oldest[1]); } } function removeBlueskyMediaCacheEntry( key: string, entry: { expiresAt: number; contentType: string; bytes: Buffer }, ): void { if (!blueskyMediaCache.delete(key)) return; blueskyMediaCacheBytes = Math.max(0, blueskyMediaCacheBytes - entry.bytes.length); } function matchesImageMagic(bytes: Buffer, contentType: string): boolean { if (contentType === "image/jpeg") return bytes.length >= 3 && bytes[0] === 0xff && bytes[1] === 0xd8 && bytes[2] === 0xff; if (contentType === "image/png") return bytes.length >= 8 && bytes.subarray(0, 8).equals(Buffer.from([0x89, 0x50, 0x4e, 0x47, 0x0d, 0x0a, 0x1a, 0x0a])); if (contentType === "image/webp") return bytes.length >= 12 && bytes.subarray(0, 4).toString("ascii") === "RIFF" && bytes.subarray(8, 12).toString("ascii") === "WEBP"; if (contentType === "image/gif") return bytes.length >= 6 && ["GIF87a", "GIF89a"].includes(bytes.subarray(0, 6).toString("ascii")); return false; } function blueskyEmbed(value: unknown): Pick { const embed = objectValue(value); const type = stringValue(embed?.$type); const media = type === "app.bsky.embed.recordWithMedia#view" ? objectValue(embed?.media) : embed; const images = Array.isArray(media?.images) ? media.images.flatMap((candidate) => { const image = objectValue(candidate); const thumb = blueskyMediaProxyPath(image?.thumb); const fullsize = safeHttpsUrl(image?.fullsize); if (!thumb || !fullsize) return []; return [{ thumb, fullsize, alt: stringValue(image?.alt) ?? "" }]; }) : []; const externalView = objectValue(media?.external); const externalUri = safeHttpUrl(externalView?.uri); const externalThumb = blueskyMediaProxyPath(externalView?.thumb); const external = externalUri ? { uri: externalUri, title: stringValue(externalView?.title) ?? externalUri, description: stringValue(externalView?.description) ?? "", ...(externalThumb ? { thumb: externalThumb } : {}), } : undefined; const recordContainer = type === "app.bsky.embed.recordWithMedia#view" ? objectValue(embed?.record) : embed; const recordView = objectValue(recordContainer?.record); const quoteAuthor = objectValue(recordView?.author); const quoteValue = objectValue(recordView?.value); const quoteUri = stringValue(recordView?.uri); const quote = quoteUri && BLUESKY_POST_URI_PATTERN.test(quoteUri) && quoteAuthor && quoteValue ? { uri: quoteUri, url: atUriToBlueskyUrl(quoteUri), author: blueskyAuthor(quoteAuthor), text: stringValue(quoteValue.text) ?? "", } : undefined; return { images, ...(external ? { external } : {}), ...(quote ? { quote } : {}) }; } function blueskyAuthor(value: Record): BlueskyPostInlay["author"] { const handle = stringValue(value.handle) ?? "unknown.handle"; const avatar = blueskyMediaProxyPath(value.avatar); return { displayName: stringValue(value.displayName)?.trim() || handle, handle, ...(avatar ? { avatar } : {}), }; } function objectValue(value: unknown): Record | undefined { return value !== null && typeof value === "object" && !Array.isArray(value) ? value as Record : undefined; } function stringValue(value: unknown): string | undefined { return typeof value === "string" ? value : undefined; } function safeHttpsUrl(value: unknown): string | undefined { const url = safeHttpUrl(value); return url?.startsWith("https://") ? url : undefined; } function blueskyMediaProxyPath(value: unknown): string | undefined { const source = safeBlueskyMediaUrl(value); return source ? `api/inlays/bluesky-media?url=${encodeURIComponent(source)}` : undefined; } function safeBlueskyMediaUrl(value: unknown): string | undefined { if (typeof value !== "string") return undefined; try { const url = new URL(value); if (url.protocol !== "https:" || url.hostname !== "cdn.bsky.app" || url.port !== "" || url.username || url.password) return undefined; if (!url.pathname.startsWith("/img/") || url.hash) return undefined; return url.toString(); } catch { return undefined; } } function safeHttpUrl(value: unknown): string | undefined { if (typeof value !== "string") return undefined; try { const url = new URL(value); return url.protocol === "https:" || url.protocol === "http:" ? url.toString() : undefined; } catch { return undefined; } } function atUriToBlueskyUrl(uri: string): string { const match = uri.match(BLUESKY_POST_URI_PATTERN); if (!match) return uri; const parts = uri.slice(5).split("/"); return `https://bsky.app/profile/${blueskyActorPathSegment(parts[0]!)}/post/${encodeURIComponent(parts[2]!)}`; } function blueskyActorPathSegment(value: string): string { return encodeURIComponent(value).replaceAll("%3A", ":"); } export function renderInspectorHtml(): string { const html = ` Stream

    Stream

    Recent activity

    Loading recent activity…

    Filter feed

    `; return html; } function sendJson(response: ServerResponse, status: number, value: unknown): void { send(response, status, "application/json; charset=utf-8", JSON.stringify(value)); } function sendBytes(response: ServerResponse, status: number, contentType: string, body: Buffer, cacheControl: string): void { response.writeHead(status, { "content-type": contentType, "content-length": String(body.length), "cache-control": cacheControl, "x-content-type-options": "nosniff", }); response.end(body); } function streamInspectorActivity( store: JazzThoughtStore, request: IncomingMessage, response: ServerResponse, ): void { response.writeHead(200, { "content-type": "text/event-stream; charset=utf-8", "cache-control": "no-store, no-transform", connection: "keep-alive", "x-accel-buffering": "no", "x-content-type-options": "nosniff", }); if (request.method === "HEAD") { response.end(); return; } response.write("retry: 2000\n\n"); let timer: NodeJS.Timeout | undefined; const notify = () => { if (timer || response.writableEnded || response.destroyed) return; timer = setTimeout(() => { timer = undefined; if (!response.writableEnded && !response.destroyed) { response.write(`event: change\ndata: ${JSON.stringify({ at: new Date().toISOString() })}\n\n`); } }, 100); timer.unref(); }; const unsubscribe = store.subscribeInspectorActivity(notify); const heartbeat = setInterval(() => { if (!response.writableEnded && !response.destroyed) response.write(": keepalive\n\n"); }, 15_000); heartbeat.unref(); const close = () => { request.off("aborted", close); response.off("close", close); if (timer) clearTimeout(timer); clearInterval(heartbeat); unsubscribe(); }; request.once("aborted", close); response.once("close", close); } function send(response: ServerResponse, status: number, contentType: string, body: string): void { response.writeHead(status, { "content-type": contentType, "cache-control": "no-store", "content-security-policy": "default-src 'self'; script-src 'self' 'unsafe-inline'; style-src 'unsafe-inline' https://fonts.googleapis.com; connect-src 'self'; img-src 'self' data:; font-src 'self' https://fonts.gstatic.com; worker-src 'self'; frame-ancestors 'none'; base-uri 'none'", "x-content-type-options": "nosniff", "x-frame-options": "DENY", }); response.end(body); } async function readBody(request: IncomingMessage, maxBytes: number): Promise { const declared = request.headers["content-length"]; if (typeof declared === "string" && (!/^\d+$/.test(declared) || Number(declared) > maxBytes)) { throw new Error("Request body too large"); } 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 > maxBytes) throw new Error("Request body too large"); parts.push(buffer); } if (bytes === 0) throw new Error("Request body is empty"); return Buffer.concat(parts); } function isLoopback(host: string): boolean { return host === "127.0.0.1" || host === "::1" || host === "localhost"; } function inspectorRunSummary(run: AgentRun) { return { id: run.id, agentId: run.agentId, agentVersion: run.agentVersion, status: run.status, provider: run.provider, model: run.model, createdAt: run.createdAt, ...(run.startedAt ? { startedAt: run.startedAt } : {}), ...(run.completedAt ? { completedAt: run.completedAt } : {}), ...(run.errorText ? { errorText: run.errorText } : {}), ...(typeof run.result?.summary === "string" ? { summary: run.result.summary } : {}), }; } function eventBelongsToRun(event: ThoughtEvent, run: AgentRun, inputsById: Map): boolean { const payloadRunId = typeof event.payload.runId === "string" ? event.payload.runId : undefined; return payloadRunId === run.id || run.inputEventIds.includes(event.id) || run.outputEventIds.includes(event.id) || run.inputEventIds.some((id) => batchReferencesEvent(inputsById.get(id), event.id)); } function batchReferencesEvent(input: ThoughtEvent | undefined, eventId: string): boolean { if (input?.type !== "stream.thought.derived.event.batch" || !Array.isArray(input.payload.members)) return false; return input.payload.members.some((reference) => ( reference !== null && typeof reference === "object" && !Array.isArray(reference) && (reference as Record).eventId === eventId )); } async function loadRunEvidenceEvents(store: JazzThoughtStore, runs: AgentRun[]): Promise { const terminalEvents = await store.listEvents({ types: RUN_TERMINAL_EVENT_TYPES }); const referencedOutputIds = [ ...runs.flatMap((run) => run.outputEventIds), ...terminalEvents.map((event) => event.payload.outputEventId).filter((id): id is string => typeof id === "string"), ]; const outputEvents = await store.getEvents(referencedOutputIds); return [...new Map([...terminalEvents, ...outputEvents].map((event) => [event.id, event])).values()]; } async function buildInspectorActivity(store: JazzThoughtStore) { return buildRecentRootActivity(store); } async function buildInspectorSystem( store: JazzThoughtStore, runs?: AgentRun[], ) { const effectiveRuns = runs ?? await store.listRuns(); const [agents, sources, evidenceEvents] = await Promise.all([ store.listAgents(), buildSourceHealth(store), loadRunEvidenceEvents(store, effectiveRuns), ]); const runEvidence = auditRunEvidence(effectiveRuns, evidenceEvents); const adapterSelections = agents .map((agent) => ({ id: agent.id, version: agent.version, enabled: agent.enabled, adapter: adapterSelectionFromSpec(agent.spec), })) .filter((agent) => agent.adapter !== undefined); const adapterCatalogs = [...new Map(adapterSelections.map((selection) => [ `${selection.adapter!.digest}:${selection.adapter!.generation}`, { digest: selection.adapter!.digest, generation: selection.adapter!.generation }, ])).values()]; return { runs: effectiveRuns.slice().reverse().map(inspectorRunSummary), agents, sources, adapterInventory: { catalogs: adapterCatalogs, selections: adapterSelections }, runEvidence, evidenceContradictions: runEvidence.filter((report) => !report.consistent).length, }; } async function expandRun(store: JazzThoughtStore, run: AgentRun): Promise<{ run: AgentRun; kind: "rule" | "model"; description: string; inputs: ThoughtEvent[]; outputs: ThoughtEvent[]; trace: Awaited>; }> { const [inputs, outputs, trace] = await Promise.all([ Promise.all(run.inputEventIds.map((id) => store.getEvent(id))), Promise.all(run.outputEventIds.map((id) => store.getEvent(id))), store.listTrace(run.id), ]); const presentInputs = inputs.filter((event): event is ThoughtEvent => Boolean(event)); return { run, kind: run.provider === "deterministic" && run.model === "deterministic" ? "rule" : "model", description: describeRunResult(run, presentInputs), inputs: presentInputs, outputs: outputs.filter((event): event is ThoughtEvent => Boolean(event)), trace, }; }