Something went wrong. Try again.
A local-first event pipeline for independent agents, built on Jazz.
Something went wrong. Try again.
49 kB · 1074 lines
TypeScript
1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075import http from "node:http";import { afterEach, describe, expect, test } from "vitest";import fs from "node:fs/promises";import os from "node:os";import path from "node:path";import { authenticatedProxyOptionsFromEnv, startAuthenticatedInspectorProxy,} from "../src/web/authenticated-proxy.js";import { InspectorOAuthAuth, OAuthCallbackQuarantineCapacityError, type BrowserSession, type OAuthProtocolClient,} from "../src/web/oauth-auth.js";import { SecureJsonStore } from "../src/web/secure-store.js";import { REVIEW_CSRF_HEADER, ReviewCapabilityVerifier,} from "../src/review/web-capability.js";import { COURSE_CHAT_CSRF_HEADER, CourseChatCapabilityVerifier,} from "../src/courses/web-capability.js";import { POST_TRAINING_COURSE } from "../src/courses/post-training.js";import { COURSE_QUESTION_EVENT_TYPE } from "../src/courses/questions.js";import { startInspectorServer } from "../src/web/inspector.js";import type { JazzThoughtStore } from "../src/jazz/store.js";import { testStore } from "./helpers.js";
const servers: http.Server[] = [];const roots: string[] = [];const stores: JazzThoughtStore[] = [];
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("authenticated inspector proxy", () => { test("gates every route and forwards only authenticated read requests without credentials", async () => { const upstreamRequests: Array<{ method: string | undefined; authorization: string | undefined; cookie: string | undefined; path: string | undefined; forwardedFor: string | undefined; forwardedProto: string | undefined; }> = []; const upstream = http.createServer((request, response) => { upstreamRequests.push({ method: request.method, authorization: request.headers.authorization, cookie: request.headers.cookie, path: request.url, forwardedFor: request.headers["x-forwarded-for"] as string | undefined, forwardedProto: request.headers["x-forwarded-proto"] as string | undefined, }); response.writeHead(200, { "content-type": "application/json", "set-cookie": "upstream-secret=forbidden", }); response.end(JSON.stringify({ private: true })); }); servers.push(upstream); const upstreamPort = await listen(upstream);
const proxy = await startAuthenticatedInspectorProxy({ port: 0, upstreamPort, username: "cameron", password: "correct-horse-battery-staple-private", basicFallbackEnabled: true, }); servers.push(proxy); const base = baseUrl(proxy);
const missing = await fetch(`${base}/inspector/api/private-object`); expect(missing.status).toBe(401); expect(missing.headers.get("www-authenticate")).toContain("thought stream"); expect(missing.headers.get("cache-control")).toBe("no-store"); expect(missing.headers.get("referrer-policy")).toBe("no-referrer"); expect(upstreamRequests).toHaveLength(0);
const wrong = await fetch(`${base}/inspector/api/private-object`, { headers: { authorization: basic("cameron", "wrong-password-with-enough-bytes") }, }); expect(wrong.status).toBe(401); expect(upstreamRequests).toHaveLength(0);
const authorized = await fetch(`${base}/inspector/api/private-object?detail=1`, { headers: { authorization: basic("cameron", "correct-horse-battery-staple-private"), cookie: "browser-secret=forbidden", "x-forwarded-for": "198.51.100.4", "x-forwarded-proto": "https", }, }); expect(authorized.status).toBe(200); expect(await authorized.json()).toEqual({ private: true }); expect(authorized.headers.get("set-cookie")).toBeNull(); expect(authorized.headers.get("x-frame-options")).toBe("DENY"); expect(upstreamRequests).toEqual([{ method: "GET", authorization: undefined, cookie: undefined, path: "/api/private-object?detail=1", forwardedFor: undefined, forwardedProto: undefined, }]);
const head = await fetch(`${base}/inspector/api/private-object`, { method: "HEAD", headers: { authorization: basic("cameron", "correct-horse-battery-staple-private") }, }); expect(head.status).toBe(200); expect(await head.text()).toBe(""); expect(upstreamRequests.at(-1)?.method).toBe("HEAD");
const mutation = await fetch(`${base}/inspector/api/private-object`, { method: "POST", headers: { authorization: basic("cameron", "correct-horse-battery-staple-private") }, }); expect(mutation.status).toBe(405); expect(mutation.headers.get("allow")).toBe("GET, HEAD"); expect(mutation.headers.get("content-security-policy")).toContain("frame-ancestors 'none'"); expect(mutation.headers.get("permissions-policy")).toContain("camera=()"); expect(upstreamRequests).toHaveLength(2); });
test("permits the authenticated inspector loader only on upstream HTML", async () => { const upstream = http.createServer((request, response) => { if (request.url === "/") { response.writeHead(200, { "content-type": "text/html; charset=utf-8" }); response.end("<!doctype html><script>fetch('api/snapshot')</script>"); return; } if (request.url?.startsWith("/api/inlays/bluesky-media?")) { response.writeHead(200, { "content-type": "image/webp", "cache-control": "private, max-age=600" }); response.end(Buffer.concat([Buffer.from("RIFF"), Buffer.alloc(4), Buffer.from("WEBP")])); return; } response.writeHead(200, { "content-type": "application/json; charset=utf-8" }); response.end("{}\n"); }); servers.push(upstream); const upstreamPort = await listen(upstream); const proxy = await startAuthenticatedInspectorProxy({ port: 0, upstreamPort, username: "cameron", password: "correct-horse-battery-staple-private", basicFallbackEnabled: true, }); servers.push(proxy); const headers = { authorization: basic("cameron", "correct-horse-battery-staple-private") };
const page = await fetch(`${baseUrl(proxy)}/inspector/`, { headers }); expect(await page.text()).toContain("fetch('api/snapshot')"); expect(page.headers.get("content-security-policy")).toContain("script-src 'self' 'unsafe-inline'"); 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")).not.toContain("img-src 'self' data: https:"); expect(page.headers.get("content-security-policy")).toContain("form-action 'none'");
const api = await fetch(`${baseUrl(proxy)}/inspector/api/snapshot`, { headers }); expect(await api.json()).toEqual({}); expect(api.headers.get("content-security-policy")).toContain("script-src 'none'");
const media = await fetch(`${baseUrl(proxy)}/inspector/api/inlays/bluesky-media?url=fixture`, { headers }); expect(media.status).toBe(200); expect(media.headers.get("content-type")).toBe("image/webp"); expect(media.headers.get("cache-control")).toBe("no-store"); expect(Buffer.from(await media.arrayBuffer()).subarray(0, 4).toString("ascii")).toBe("RIFF"); });
test("serves only allowlisted public pages and never falls through to the private upstream", async () => { let upstreamRequests = 0; const upstream = http.createServer((_request, response) => { upstreamRequests += 1; response.end("private"); }); servers.push(upstream); const upstreamPort = await listen(upstream); const proxy = await startAuthenticatedInspectorProxy({ port: 0, upstreamPort, username: "cameron", password: "correct-horse-battery-staple-private", basicFallbackEnabled: true, }); servers.push(proxy); const base = baseUrl(proxy);
const landing = await fetch(base); expect(landing.status).toBe(200); const landingHtml = await landing.text(); expect(landingHtml).toContain("<title>Stream</title>"); expect(landingHtml).toContain("<h1>Stream</h1>"); expect(landingHtml).not.toContain("A private feed"); expect(landingHtml).toContain('<form class="login-form" method="post" action="/oauth/login">'); expect(landingHtml).toContain('<button type="submit">Log in</button>'); expect(landingHtml).toMatch(/\.landing-main h1\{[^}]*text-align:center\}/); expect(landingHtml).toContain(".login-form{margin-top:28px;text-align:center}"); expect(landingHtml).toContain( ".landing-footer{padding:16px 0 0;font-size:12px;text-align:center}", ); expect(landingHtml).toContain('<script src="/assets/font-debug.js" defer></script>'); expect(landingHtml).not.toContain(".landing-footer{padding:16px 0 0;border-top:"); expect(landingHtml).toContain('href="https://tangled.org/@cameron.stream/thought-stream">code</a>'); expect(landingHtml).toContain("data:font/woff2;base64,"); expect(landingHtml).not.toContain('href="/docs"'); expect(landingHtml).not.toContain("private inspector"); expect(landing.headers.get("content-security-policy")).toContain("form-action 'self' https:"); expect(landing.headers.get("content-security-policy")).toContain("font-src data:"); expect(landing.headers.get("content-security-policy")).toContain("https://fonts.googleapis.com"); expect(landing.headers.get("content-security-policy")).toContain("https://fonts.gstatic.com"); const fontDebug = await fetch(`${base}/assets/font-debug.js`); expect(fontDebug.status).toBe(200); expect(fontDebug.headers.get("content-type")).toBe("text/javascript; charset=utf-8"); const fontDebugScript = await fontDebug.text(); expect(() => new Function(fontDebugScript)).not.toThrow(); const architecture = await fetch(`${base}/docs/architecture`); expect(architecture.status).toBe(200); expect(await architecture.text()).toContain("Connector cursors advance only after durable events"); expect(architecture.headers.get("content-security-policy")).toContain("form-action 'self'"); expect(architecture.headers.get("content-security-policy")).not.toContain("form-action 'self' https:"); const traversal = await fetch(`${base}/docs/..%2f..%2fetc%2fpasswd`); expect(traversal.status).toBe(404); const oldApi = await fetch(`${base}/api/private-object`); expect(oldApi.status).toBe(404); const unknown = await fetch(`${base}/anything`); expect(unknown.status).toBe(404); expect(upstreamRequests).toBe(0); });
test("completes an OAuth browser session while Basic remains available", async () => { const root = await fs.mkdtemp(path.join(os.tmpdir(), "thoughtstream-proxy-oauth-")); roots.push(root); const key = Buffer.alloc(32, 6); const browserSessions = new SecureJsonStore<BrowserSession>({ directory: root, name: "browser", key, maxEntries: 8, maxSerializedBytes: 32 * 1024, }); const flowStates = new SecureJsonStore<{ createdAt: number }>({ directory: root, name: "flow", key, maxEntries: 64, maxSerializedBytes: 32 * 1024, ttlMs: 15 * 60_000, }); const protocolState = "sdk-protocol-state-1234567890abcdef"; let applicationState: string | undefined; const stagedAttempts = new Map<string, string>(); let currentGeneration: number | undefined; let nextGeneration = 0; const protocol: OAuthProtocolClient = { clientMetadata: { client_id: "https://thought.stream/oauth/client-metadata.json", dpop_bound_access_tokens: true }, jwks: { keys: [{ kty: "EC", kid: "fixture" }] }, authorize: async (_handle, options) => { applicationState = options.state; return new URL("https://pds.example/authorize?request_uri=urn:ietf:params:oauth:request_uri:fixture"); }, callback: async (_params, attemptId) => { stagedAttempts.set(attemptId, "did:plc:allowed"); return { session: { did: "did:plc:allowed" }, state: applicationState ?? null }; }, expireCallbackAttempt: () => undefined, dropCallbackAttempt: async (attemptId: string) => { stagedAttempts.delete(attemptId); }, promoteCallbackSession: async (attemptId, did, guard) => { if (!guard() || stagedAttempts.get(attemptId) !== did) throw new Error("missing staged session"); stagedAttempts.delete(attemptId); currentGeneration = ++nextGeneration; return currentGeneration; }, isCurrentGeneration: (_did, generation) => currentGeneration === generation, discardPromotedSession: async (_did, generation) => { if (currentGeneration === generation) currentGeneration = undefined; }, restore: async (did, generation) => { if (currentGeneration !== generation) throw new Error("stale generation"); return { did }; }, }; const oauth = new InspectorOAuthAuth({ protocol, allowedDid: "did:plc:allowed", expectedHandle: "cameron.stream", browserSessions, flowStates, }); const upstream = http.createServer((_request, response) => response.end("private inspector")); servers.push(upstream); const upstreamPort = await listen(upstream); const proxy = await startAuthenticatedInspectorProxy({ port: 0, upstreamPort, username: "cameron", password: "correct-horse-battery-staple-private", oauth, basicFallbackEnabled: true, }); servers.push(proxy); const base = baseUrl(proxy);
const metadata = await fetch(`${base}/oauth/client-metadata.json`); expect(await metadata.json()).toMatchObject({ client_id: "https://thought.stream/oauth/client-metadata.json" }); const loginPageResponse = await fetch(`${base}/oauth/login`); expect(loginPageResponse.headers.get("content-security-policy")).toContain("form-action 'self' https:"); const protectedRoute = await fetch(`${base}/inspector`, { headers: { accept: "text/html" }, redirect: "manual" }); expect(protectedRoute.status).toBe(303); expect(protectedRoute.headers.get("location")).toBe("/oauth/login");
const login = await fetch(`${base}/oauth/login`, { method: "POST", redirect: "manual" }); expect(login.status).toBe(303); expect(login.headers.get("content-security-policy")).toContain("form-action 'self' https:"); const flowCookie = login.headers.get("set-cookie")!; const state = cookieFromHeader(flowCookie, "__Host-thoughtstream_oauth"); const callback = await fetch(`${base}/oauth/callback?state=${encodeURIComponent(protocolState)}&code=fixture`, { headers: { cookie: `__Host-thoughtstream_oauth=${state}` }, redirect: "manual", }); expect(callback.status).toBe(302); expect(callback.headers.get("location")).toBe("/inspector/"); const sessionId = cookieFromHeader(callback.headers.get("set-cookie")!, "__Host-thoughtstream_session");
const inspector = await fetch(`${base}/inspector/api/private`, { headers: { cookie: `__Host-thoughtstream_session=${sessionId}` }, }); expect(await inspector.text()).toBe("private inspector"); const logoutPage = await fetch(`${base}/oauth/logout`, { headers: { cookie: `__Host-thoughtstream_session=${sessionId}` }, }); const csrf = (await logoutPage.text()).match(/name="csrfToken" value="([^"]+)"/)?.[1]; expect(csrf).toBeDefined(); const logout = await fetch(`${base}/oauth/logout`, { method: "POST", headers: { cookie: `__Host-thoughtstream_session=${sessionId}`, "content-type": "application/x-www-form-urlencoded", }, body: new URLSearchParams({ csrfToken: csrf! }), redirect: "manual", }); expect(logout.status).toBe(303); expect(currentGeneration).toBeUndefined(); });
test("keeps Basic read-only and forwards one CSRF-checked OAuth Review mutation under a body-bound capability", async () => { const root = await fs.mkdtemp(path.join(os.tmpdir(), "thoughtstream-review-proxy-")); roots.push(root); const fixture = proxyOAuthFixture(root); const reviewCapability = Buffer.alloc(32, 23); const upstreamRequests: Array<{ method: string | undefined; path: string | undefined; headers: http.IncomingHttpHeaders; body: Buffer; }> = []; const upstream = http.createServer(async (request, response) => { const parts: Buffer[] = []; for await (const part of request) parts.push(Buffer.isBuffer(part) ? part : Buffer.from(part)); upstreamRequests.push({ method: request.method, path: request.url, headers: request.headers, body: Buffer.concat(parts), }); response.writeHead(201, { "content-type": "application/json" }); response.end('{"decisionEventId":"fixture"}'); }); servers.push(upstream); const upstreamPort = await listen(upstream); const proxy = await startAuthenticatedInspectorProxy({ port: 0, upstreamPort, username: "cameron", password: "correct-horse-battery-staple-private", oauth: fixture.oauth, basicFallbackEnabled: true, reviewCapability, }); servers.push(proxy); const base = baseUrl(proxy); const basicAuthorization = basic("cameron", "correct-horse-battery-staple-private"); const basicSession = await fetch(`${base}/inspector/api/session`, { headers: { authorization: basicAuthorization }, }); expect(await basicSession.json()).toEqual({ reviewWriteEnabled: false, courseChatEnabled: false }); const basicWrite = await fetch(`${base}/inspector/api/reviews/review%3Aitem/decisions`, { method: "POST", headers: { authorization: basicAuthorization, "content-type": "application/json" }, body: '{"disposition":"skip"}', }); expect(basicWrite.status).toBe(403); expect(upstreamRequests).toHaveLength(0);
const login = await fetch(`${base}/oauth/login`, { method: "POST", redirect: "manual" }); const applicationState = cookieFromHeader(login.headers.get("set-cookie")!, "__Host-thoughtstream_oauth"); const callback = await fetch(`${base}/oauth/callback?state=sdk-protocol-state-1234567890abcdef&code=fixture`, { headers: { cookie: `__Host-thoughtstream_oauth=${applicationState}` }, redirect: "manual", }); const sessionId = cookieFromHeader(callback.headers.get("set-cookie")!, "__Host-thoughtstream_session"); const cookie = `__Host-thoughtstream_session=${sessionId}`; const session = await fetch(`${base}/inspector/api/session`, { headers: { cookie } }); const sessionBody = await session.json() as { reviewWriteEnabled: boolean; csrfToken: string }; expect(sessionBody.reviewWriteEnabled).toBe(true); expect(sessionBody.csrfToken).toMatch(/^[A-Za-z0-9_-]+$/);
const body = JSON.stringify({ disposition: "skip", reasonCodes: [], responseTags: [], trainingEligible: false, submissionId: "submission-proxy-canary-0001", }); const missingCsrf = await fetch(`${base}/inspector/api/reviews/review%3Aitem/decisions`, { method: "POST", headers: { cookie, "content-type": "application/json" }, body, }); expect(missingCsrf.status).toBe(403); expect(upstreamRequests).toHaveLength(0); const written = await fetch(`${base}/inspector/api/reviews/review%3Aitem/decisions`, { method: "POST", headers: { cookie, "content-type": "application/json", [REVIEW_CSRF_HEADER]: sessionBody.csrfToken, }, body, }); expect(written.status).toBe(201); expect(upstreamRequests).toHaveLength(1); const forwarded = upstreamRequests[0]!; expect(forwarded.method).toBe("POST"); expect(forwarded.path).toBe("/api/reviews/review%3Aitem/decisions"); expect(forwarded.body.toString("utf8")).toBe(body); expect(forwarded.headers.authorization).toBeUndefined(); expect(forwarded.headers.cookie).toBeUndefined(); expect(forwarded.headers[REVIEW_CSRF_HEADER]).toBeUndefined(); const verifier = new ReviewCapabilityVerifier(reviewCapability); expect(verifier.verify(forwarded.headers, "POST", forwarded.path!, forwarded.body)).toBe(true); });
test("forwards one OAuth-only course question under its separate body-bound capability", async () => { const root = await fs.mkdtemp(path.join(os.tmpdir(), "thoughtstream-course-chat-proxy-")); roots.push(root); const fixture = proxyOAuthFixture(root); const courseChatCapability = Buffer.alloc(32, 29); const upstreamRequests: Array<{ method: string | undefined; path: string | undefined; headers: http.IncomingHttpHeaders; body: Buffer; }> = []; const upstream = http.createServer(async (request, response) => { const parts: Buffer[] = []; for await (const part of request) parts.push(Buffer.isBuffer(part) ? part : Buffer.from(part)); upstreamRequests.push({ method: request.method, path: request.url, headers: request.headers, body: Buffer.concat(parts) }); response.writeHead(202, { "content-type": "application/json" }); response.end('{"eventId":"evt_fixture"}'); }); servers.push(upstream); const upstreamPort = await listen(upstream); const proxy = await startAuthenticatedInspectorProxy({ port: 0, upstreamPort, username: "cameron", password: "correct-horse-battery-staple-private", oauth: fixture.oauth, basicFallbackEnabled: true, courseChatCapability, }); servers.push(proxy); const base = baseUrl(proxy); const basicAuthorization = basic("cameron", "correct-horse-battery-staple-private"); const basicSession = await fetch(`${base}/inspector/api/session`, { headers: { authorization: basicAuthorization } }); expect(await basicSession.json()).toEqual({ reviewWriteEnabled: false, courseChatEnabled: false }); const body = JSON.stringify({ requestId: "question-proxy-0001", courseRevision: "a".repeat(64), lessonId: "model-factory", question: "What does the learner own?", }); const basicWrite = await fetch(`${base}/inspector/api/courses/post-training/questions`, { method: "POST", headers: { authorization: basicAuthorization, "content-type": "application/json" }, body, }); expect(basicWrite.status).toBe(403); expect(upstreamRequests).toHaveLength(0);
const login = await fetch(`${base}/oauth/login`, { method: "POST", redirect: "manual" }); const applicationState = cookieFromHeader(login.headers.get("set-cookie")!, "__Host-thoughtstream_oauth"); const callback = await fetch(`${base}/oauth/callback?state=sdk-protocol-state-1234567890abcdef&code=fixture`, { headers: { cookie: `__Host-thoughtstream_oauth=${applicationState}` }, redirect: "manual", }); const sessionId = cookieFromHeader(callback.headers.get("set-cookie")!, "__Host-thoughtstream_session"); const cookie = `__Host-thoughtstream_session=${sessionId}`; const session = await fetch(`${base}/inspector/api/session`, { headers: { cookie } }); const sessionBody = await session.json() as { reviewWriteEnabled: boolean; courseChatEnabled: boolean; csrfToken: string }; expect(sessionBody).toMatchObject({ reviewWriteEnabled: false, courseChatEnabled: true });
const missingCsrf = await fetch(`${base}/inspector/api/courses/post-training/questions`, { method: "POST", headers: { cookie, "content-type": "application/json" }, body, }); expect(missingCsrf.status).toBe(403); const written = await fetch(`${base}/inspector/api/courses/post-training/questions`, { method: "POST", headers: { cookie, "content-type": "application/json", [COURSE_CHAT_CSRF_HEADER]: sessionBody.csrfToken, }, body, }); expect(written.status).toBe(202); expect(upstreamRequests).toHaveLength(1); const forwarded = upstreamRequests[0]!; expect(forwarded.path).toBe("/api/courses/post-training/questions"); expect(forwarded.body.toString("utf8")).toBe(body); expect(forwarded.headers.authorization).toBeUndefined(); expect(forwarded.headers.cookie).toBeUndefined(); expect(forwarded.headers[COURSE_CHAT_CSRF_HEADER]).toBeUndefined(); expect(new CourseChatCapabilityVerifier(courseChatCapability) .verify(forwarded.headers, "POST", forwarded.path!, forwarded.body)).toBe(true); });
test("guards every workbench write route like the Review decision route and keeps workbench reads readable", async () => { const root = await fs.mkdtemp(path.join(os.tmpdir(), "thoughtstream-workbench-proxy-")); roots.push(root); const fixture = proxyOAuthFixture(root); const reviewCapability = Buffer.alloc(32, 37); const upstreamRequests: Array<{ method: string | undefined; path: string | undefined; headers: http.IncomingHttpHeaders; body: Buffer }> = []; const upstream = http.createServer(async (request, response) => { const parts: Buffer[] = []; for await (const part of request) parts.push(Buffer.isBuffer(part) ? part : Buffer.from(part)); upstreamRequests.push({ method: request.method, path: request.url, headers: request.headers, body: Buffer.concat(parts) }); response.writeHead(request.method === "POST" ? 201 : 200, { "content-type": "application/json" }); response.end(request.method === "POST" ? '{"documentId":"fixture"}' : '{"items":[]}'); }); servers.push(upstream); const upstreamPort = await listen(upstream); const proxy = await startAuthenticatedInspectorProxy({ port: 0, upstreamPort, username: "cameron", password: "correct-horse-battery-staple-private", oauth: fixture.oauth, basicFallbackEnabled: true, reviewCapability, }); servers.push(proxy); const base = baseUrl(proxy); const basicAuthorization = basic("cameron", "correct-horse-battery-staple-private"); const writeRoutes = [ "/inspector/api/workbench/documents", "/inspector/api/workbench/documents/doc%3A1/versions", "/inspector/api/workbench/documents/doc%3A1/selections", "/inspector/api/workbench/documents/doc%3A1/proposals", "/inspector/api/workbench/proposals/evt_1/decisions", ]; const body = JSON.stringify({ originEventId: "evt_origin", title: "Notes", body: "# Notes", requestId: "req-proxy-00001" });
// Unauthenticated: rejected before any upstream contact, without disclosing the route. for (const route of writeRoutes) { const anonymous = await fetch(`${base}${route}`, { method: "POST", headers: { "content-type": "application/json" }, body }); expect(anonymous.status, route).toBe(401); } // Wrong owner class: Basic remains read-only on every workbench write. for (const route of writeRoutes) { const basicWrite = await fetch(`${base}${route}`, { method: "POST", headers: { authorization: basicAuthorization, "content-type": "application/json" }, body }); expect(basicWrite.status, route).toBe(403); } expect(upstreamRequests).toHaveLength(0); // Basic still reads the workbench projections through the ordinary GET path. const basicRead = await fetch(`${base}/inspector/api/workbench/documents`, { headers: { authorization: basicAuthorization } }); expect(basicRead.status).toBe(200); expect(upstreamRequests).toHaveLength(1); expect(upstreamRequests[0]).toMatchObject({ method: "GET", path: "/api/workbench/documents" }); upstreamRequests.length = 0;
const login = await fetch(`${base}/oauth/login`, { method: "POST", redirect: "manual" }); const applicationState = cookieFromHeader(login.headers.get("set-cookie")!, "__Host-thoughtstream_oauth"); const callback = await fetch(`${base}/oauth/callback?state=sdk-protocol-state-1234567890abcdef&code=fixture`, { headers: { cookie: `__Host-thoughtstream_oauth=${applicationState}` }, redirect: "manual", }); const sessionId = cookieFromHeader(callback.headers.get("set-cookie")!, "__Host-thoughtstream_session"); const cookie = `__Host-thoughtstream_session=${sessionId}`; const session = await fetch(`${base}/inspector/api/session`, { headers: { cookie } }); const sessionBody = await session.json() as { reviewWriteEnabled: boolean; csrfToken: string }; expect(sessionBody.reviewWriteEnabled).toBe(true);
// Cross-origin shape: an owner cookie without the session CSRF header never reaches the inspector. for (const route of writeRoutes) { const missingCsrf = await fetch(`${base}${route}`, { method: "POST", headers: { cookie, "content-type": "application/json" }, body }); expect(missingCsrf.status, route).toBe(403); const wrongCsrf = await fetch(`${base}${route}`, { method: "POST", headers: { cookie, "content-type": "application/json", [REVIEW_CSRF_HEADER]: "not-the-token" }, body }); expect(wrongCsrf.status, route).toBe(403); } expect(upstreamRequests).toHaveLength(0); // Query strings, oversized bodies, and non-JSON bodies are rejected locally. const withQuery = await fetch(`${base}/inspector/api/workbench/documents?x=1`, { method: "POST", headers: { cookie, "content-type": "application/json", [REVIEW_CSRF_HEADER]: sessionBody.csrfToken }, body }); expect(withQuery.status).toBe(405); const oversized = await fetch(`${base}/inspector/api/workbench/documents`, { method: "POST", headers: { cookie, "content-type": "application/json", [REVIEW_CSRF_HEADER]: sessionBody.csrfToken }, body: JSON.stringify({ body: "x".repeat(100_000) }) }); expect(oversized.status).toBe(400); const notJson = await fetch(`${base}/inspector/api/workbench/documents`, { method: "POST", headers: { cookie, "content-type": "text/plain", [REVIEW_CSRF_HEADER]: sessionBody.csrfToken }, body }); expect(notJson.status).toBe(400); expect(upstreamRequests).toHaveLength(0);
// A CSRF-checked owner write is forwarded once with a fresh body-bound signature and no browser credentials. for (const route of writeRoutes) { const written = await fetch(`${base}${route}`, { method: "POST", headers: { cookie, "content-type": "application/json", [REVIEW_CSRF_HEADER]: sessionBody.csrfToken }, body }); expect(written.status, route).toBe(201); } expect(upstreamRequests).toHaveLength(writeRoutes.length); const verifier = new ReviewCapabilityVerifier(reviewCapability); for (const [index, forwarded] of upstreamRequests.entries()) { expect(forwarded.method).toBe("POST"); expect(forwarded.path).toBe(writeRoutes[index]!.slice("/inspector".length)); expect(forwarded.body.toString("utf8")).toBe(body); expect(forwarded.headers.authorization).toBeUndefined(); expect(forwarded.headers.cookie).toBeUndefined(); expect(forwarded.headers[REVIEW_CSRF_HEADER]).toBeUndefined(); expect(verifier.verify(forwarded.headers, "POST", forwarded.path!, forwarded.body)).toBe(true); } });
test("carries an OAuth course question through the signed proxy into one canonical event", async () => { const root = await fs.mkdtemp(path.join(os.tmpdir(), "thoughtstream-course-chat-chain-")); roots.push(root); const fixture = proxyOAuthFixture(root); const store = await testStore(root); stores.push(store); const courseChatCapability = Buffer.alloc(32, 31); const inspector = await startInspectorServer(store, { port: 0, courseChatCapability }); servers.push(inspector); const inspectorAddress = inspector.address(); if (!inspectorAddress || typeof inspectorAddress === "string") throw new Error("Missing inspector address"); const proxy = await startAuthenticatedInspectorProxy({ port: 0, upstreamPort: inspectorAddress.port, username: "cameron", password: "correct-horse-battery-staple-private", oauth: fixture.oauth, basicFallbackEnabled: true, courseChatCapability, }); servers.push(proxy); const base = baseUrl(proxy); const login = await fetch(`${base}/oauth/login`, { method: "POST", redirect: "manual" }); const applicationState = cookieFromHeader(login.headers.get("set-cookie")!, "__Host-thoughtstream_oauth"); const callback = await fetch(`${base}/oauth/callback?state=sdk-protocol-state-1234567890abcdef&code=fixture`, { headers: { cookie: `__Host-thoughtstream_oauth=${applicationState}` }, redirect: "manual", }); const sessionId = cookieFromHeader(callback.headers.get("set-cookie")!, "__Host-thoughtstream_session"); const cookie = `__Host-thoughtstream_session=${sessionId}`; const session = await fetch(`${base}/inspector/api/session`, { headers: { cookie } }); const sessionBody = await session.json() as { courseChatEnabled: boolean; csrfToken: string }; expect(sessionBody.courseChatEnabled).toBe(true); const accepted = await fetch(`${base}/inspector/api/courses/post-training/questions`, { method: "POST", headers: { cookie, "content-type": "application/json", [COURSE_CHAT_CSRF_HEADER]: sessionBody.csrfToken, }, body: JSON.stringify({ requestId: "question-chain-0001", courseRevision: POST_TRAINING_COURSE.revision, lessonId: "model-factory", question: "Which factory stage owns the rollback pointer?", }), }); expect(accepted.status).toBe(202); const receipt = await accepted.json() as { eventId: string; statusPath: string }; expect(await store.getEvent(receipt.eventId)).toMatchObject({ id: receipt.eventId, type: COURSE_QUESTION_EVENT_TYPE, source: "web-course:post-training-model-factory", privacy: "sensitive", }); const status = await fetch(`${base}/inspector/${receipt.statusPath}`, { headers: { cookie } }); expect(status.status).toBe(200); expect(await status.json()).toEqual({ status: "pending", eventId: receipt.eventId }); });
test("keeps Basic break-glass independent and stops advertising it when disabled", async () => { const root = await fs.mkdtemp(path.join(os.tmpdir(), "thoughtstream-break-glass-")); roots.push(root); const fixture = proxyOAuthFixture(root); let upstreamRequests = 0; const upstream = http.createServer((_request, response) => { upstreamRequests += 1; response.end("private"); }); servers.push(upstream); const upstreamPort = await listen(upstream); const disabled = await startAuthenticatedInspectorProxy({ port: 0, upstreamPort, oauth: fixture.oauth, }); servers.push(disabled);
const rejectedBasic = await fetch(`${baseUrl(disabled)}/inspector/api/private`, { headers: { authorization: basic("cameron", "correct-horse-battery-staple-private") }, }); const rejectedMissing = await fetch(`${baseUrl(disabled)}/inspector/api/private`); expect(rejectedBasic.status).toBe(401); expect(rejectedMissing.status).toBe(401); expect(await rejectedBasic.text()).toBe("Authentication required.\n"); expect(await rejectedMissing.text()).toBe("Authentication required.\n"); expect(rejectedBasic.headers.get("www-authenticate")).toBeNull(); expect(rejectedMissing.headers.get("www-authenticate")).toBeNull(); expect(upstreamRequests).toBe(0);
const enabled = await startAuthenticatedInspectorProxy({ port: 0, upstreamPort, username: "cameron", password: "correct-horse-battery-staple-private", oauth: fixture.oauth, basicFallbackEnabled: true, }); servers.push(enabled); const acceptedBasic = await fetch(`${baseUrl(enabled)}/inspector/api/private`, { headers: { authorization: basic("cameron", "correct-horse-battery-staple-private") }, }); expect(acceptedBasic.status).toBe(200); expect(await acceptedBasic.text()).toBe("private"); expect(fixture.protocol.restoreCalls).toEqual([]); });
test("surfaces callback quarantine exhaustion as operator recycle required", async () => { const root = await fs.mkdtemp(path.join(os.tmpdir(), "thoughtstream-quarantine-capacity-")); roots.push(root); const fixture = proxyOAuthFixture(root); fixture.protocol.callbackError = new OAuthCallbackQuarantineCapacityError(8); const proxy = await startAuthenticatedInspectorProxy({ port: 0, oauth: fixture.oauth, }); servers.push(proxy); const base = baseUrl(proxy); const login = await fetch(`${base}/oauth/login`, { method: "POST", redirect: "manual" }); const state = cookieFromHeader(login.headers.get("set-cookie")!, "__Host-thoughtstream_oauth"); const callback = await fetch(`${base}/oauth/callback?state=sdk-protocol-state-1234567890abcdef&code=fixture`, { headers: { cookie: `__Host-thoughtstream_oauth=${state}` }, redirect: "manual", }); expect(callback.status).toBe(503); expect(callback.headers.get("retry-after")).toBe("60"); expect(callback.headers.get("set-cookie")).toBeNull(); expect(await callback.text()).toBe("OAuth callback capacity reached. Operator recycle required.\n");
fixture.protocol.callbackError = undefined; const retry = await fetch(`${base}/oauth/callback?state=sdk-protocol-state-1234567890abcdef&code=fixture`, { headers: { cookie: `__Host-thoughtstream_oauth=${state}` }, redirect: "manual", }); expect(retry.status).toBe(302); expect(retry.headers.get("location")).toBe("/inspector/"); });
test("rate-limits OAuth initiation and callback before invoking the SDK", async () => { const root = await fs.mkdtemp(path.join(os.tmpdir(), "thoughtstream-rate-limit-")); roots.push(root); const fixture = proxyOAuthFixture(root); const proxy = await startAuthenticatedInspectorProxy({ port: 0, oauth: fixture.oauth, oauthRateLimiter: { login: () => ({ allowed: false, retryAfterSeconds: 17 }), callback: () => ({ allowed: false, retryAfterSeconds: 23 }), }, }); servers.push(proxy); const base = baseUrl(proxy);
const login = await fetch(`${base}/oauth/login`, { method: "POST", redirect: "manual" }); expect(login.status).toBe(429); expect(login.headers.get("retry-after")).toBe("17"); const callback = await fetch(`${base}/oauth/callback?state=sdk-protocol-state-1234567890abcdef&code=private-code`); expect(callback.status).toBe(429); expect(callback.headers.get("retry-after")).toBe("23"); expect(fixture.protocol.authorizeCalls).toBe(0); expect(fixture.protocol.callbackCalls).toBe(0); });
test("aborts OAuth discovery when the initiating browser disconnects", async () => { const root = await fs.mkdtemp(path.join(os.tmpdir(), "thoughtstream-disconnect-")); roots.push(root); let authorizeStarted!: () => void; const started = new Promise<void>((resolve) => { authorizeStarted = resolve; }); let aborted = false; const fixture = proxyOAuthFixture(root, { authorize: async (_handle, options) => { authorizeStarted(); await new Promise<void>((_resolve, reject) => { options.signal?.addEventListener("abort", () => { aborted = true; reject(new Error("aborted")); }, { once: true }); }); throw new Error("unreachable"); }, }); const proxy = await startAuthenticatedInspectorProxy({ port: 0, oauth: fixture.oauth, }); servers.push(proxy); const address = proxy.address(); if (!address || typeof address === "string") throw new Error("Missing proxy address"); const request = http.request({ hostname: "127.0.0.1", port: address.port, path: "/oauth/login", method: "POST", }); request.on("error", () => undefined); request.end(); await started; request.destroy(); await waitFor(() => aborted); expect(aborted).toBe(true); });
test("returns a generic content-dark response when the inspector is unavailable", async () => { const unavailablePort = await unusedPort(); const proxy = await startAuthenticatedInspectorProxy({ port: 0, upstreamPort: unavailablePort, username: "cameron", password: "another-correctly-long-private-password", basicFallbackEnabled: true, }); servers.push(proxy); const response = await fetch(`${baseUrl(proxy)}/inspector/`, { headers: { authorization: basic("cameron", "another-correctly-long-private-password") }, }); expect(response.status).toBe(502); expect(await response.text()).toBe("Inspector unavailable.\n"); });
test("refuses public binds, non-loopback upstreams, short passwords, and malformed environment", async () => { await expect(startAuthenticatedInspectorProxy({ host: "0.0.0.0", username: "cameron", password: "correctly-long-private-password", })).rejects.toThrow("loopback"); await expect(startAuthenticatedInspectorProxy({ upstreamHost: "example.com", username: "cameron", password: "correctly-long-private-password", })).rejects.toThrow("upstream must be loopback"); await expect(startAuthenticatedInspectorProxy({ username: "cameron", password: "too-short", basicFallbackEnabled: true, })).rejects.toThrow("at least 20"); await expect(startAuthenticatedInspectorProxy({ username: "cameron:admin", password: "correctly-long-private-password", basicFallbackEnabled: true, })).rejects.toThrow("username"); await expect(startAuthenticatedInspectorProxy({})).rejects.toThrow("requires OAuth or enabled Basic fallback"); expect(() => authenticatedProxyOptionsFromEnv({ PROXY_BASIC_FALLBACK_ENABLED: "1", PROXY_USER: "cameron", PROXY_PASSWORD_B64: "not canonical base64 !!!", })).toThrow("canonical base64"); expect(() => authenticatedProxyOptionsFromEnv({ PROXY_BASIC_FALLBACK_ENABLED: "true", PROXY_USER: "cameron", })).toThrow("PROXY_PASSWORD_B64"); expect(() => authenticatedProxyOptionsFromEnv({ PROXY_BASIC_FALLBACK_ENABLED: "sometimes", })).toThrow("must be one of"); expect(() => authenticatedProxyOptionsFromEnv({ PROXY_BASIC_FALLBACK_ENABLED: " TRUE ", })).toThrow("must be one of"); expect(() => authenticatedProxyOptionsFromEnv({ PROXY_BASIC_FALLBACK_ENABLED: "False", })).toThrow("must be one of"); });
test("defaults Basic off, ignores dormant credentials, and loads credentials only when explicitly enabled", () => { expect(authenticatedProxyOptionsFromEnv({ PROXY_USER: "ignored", PROXY_PASSWORD_B64: "not canonical base64 !!!", })).toEqual({ host: "127.0.0.1", port: 4319, upstreamHost: "127.0.0.1", upstreamPort: 4317, basicFallbackEnabled: false, projectRoot: process.cwd(), }); for (const disabledValue of ["0", "false"]) { expect(authenticatedProxyOptionsFromEnv({ PROXY_BASIC_FALLBACK_ENABLED: disabledValue, PROXY_USER: "ignored", PROXY_PASSWORD_B64: "not canonical base64 !!!", })).toMatchObject({ basicFallbackEnabled: false }); } const options = authenticatedProxyOptionsFromEnv({ PROXY_BASIC_FALLBACK_ENABLED: "1", PROXY_USER: "cameron", PROXY_PASSWORD_B64: Buffer.from("correct-horse-battery-staple-private").toString("base64"), PROXY_PORT: "4319", PROXY_UPSTREAM_PORT: "4317", }); expect(options).toEqual({ host: "127.0.0.1", port: 4319, upstreamHost: "127.0.0.1", upstreamPort: 4317, username: "cameron", password: "correct-horse-battery-staple-private", basicFallbackEnabled: true, projectRoot: process.cwd(), }); });});
function proxyOAuthFixture( root: string, options: { authorize?: OAuthProtocolClient["authorize"] } = {},): { oauth: InspectorOAuthAuth; protocol: ProxyOAuthProtocol } { const key = Buffer.alloc(32, 12); const browserSessions = new SecureJsonStore<BrowserSession>({ directory: root, name: "proxy-browser", key, maxEntries: 8, maxSerializedBytes: 32 * 1024, }); const flowStates = new SecureJsonStore<{ createdAt: number }>({ directory: root, name: "proxy-flow", key, maxEntries: 64, maxSerializedBytes: 32 * 1024, ttlMs: 15 * 60_000, }); const protocol = new ProxyOAuthProtocol(options.authorize); return { oauth: new InspectorOAuthAuth({ protocol, allowedDid: "did:plc:allowed", expectedHandle: "cameron.stream", browserSessions, flowStates, }), protocol, };}
class ProxyOAuthProtocol implements OAuthProtocolClient { readonly clientMetadata = { client_id: "https://thought.stream/oauth/client-metadata.json" }; readonly jwks = { keys: [] }; authorizeCalls = 0; callbackCalls = 0; readonly restoreCalls: string[] = []; callbackError: Error | undefined; private applicationState?: string; private readonly stagedAttempts = new Map<string, string>(); private currentGeneration: number | undefined; private nextGeneration = 0;
constructor(private readonly customAuthorize?: OAuthProtocolClient["authorize"]) {}
async authorize(handle: string, options: { state: string; scope: string; signal?: AbortSignal }): Promise<URL> { this.authorizeCalls += 1; this.applicationState = options.state; if (this.customAuthorize) return this.customAuthorize(handle, options); return new URL("https://pds.example/authorize?request_uri=urn:ietf:params:oauth:request_uri:fixture"); }
async callback(_params: URLSearchParams, attemptId: string): Promise<{ session: { did: string }; state: string | null }> { this.callbackCalls += 1; if (this.callbackError) throw this.callbackError; this.stagedAttempts.set(attemptId, "did:plc:allowed"); return { session: { did: "did:plc:allowed" }, state: this.applicationState ?? null }; }
expireCallbackAttempt(): void {}
async dropCallbackAttempt(attemptId: string): Promise<void> { this.stagedAttempts.delete(attemptId); }
async promoteCallbackSession(attemptId: string, did: string, guard: () => boolean): Promise<number> { if (!guard() || this.stagedAttempts.get(attemptId) !== did) throw new Error("missing staged session"); this.stagedAttempts.delete(attemptId); this.currentGeneration = ++this.nextGeneration; return this.currentGeneration; }
isCurrentGeneration(_did: string, generation: number): boolean { return this.currentGeneration === generation; }
async discardPromotedSession(_did: string, generation: number): Promise<void> { if (this.currentGeneration === generation) this.currentGeneration = undefined; }
async restore(did: string, generation: number): Promise<{ did: string }> { this.restoreCalls.push(did); if (this.currentGeneration !== generation) throw new Error("stale generation"); return { did }; }
async revoke(): Promise<void> {} async deleteSession(): Promise<void> {}}
async function waitFor(predicate: () => boolean): Promise<void> { for (let attempt = 0; attempt < 100; attempt += 1) { if (predicate()) return; await new Promise((resolve) => setTimeout(resolve, 5)); } throw new Error("Timed out waiting for condition");}
function cookieFromHeader(header: string, name: string): string { const match = header.match(new RegExp(`(?:^|,\\s*)${name}=([^;]+)`)); if (!match?.[1]) throw new Error(`Missing ${name} cookie`); return match[1];}
function basic(username: string, password: string): string { return `Basic ${Buffer.from(`${username}:${password}`).toString("base64")}`;}
async function listen(server: http.Server): Promise<number> { await new Promise<void>((resolve, reject) => { server.once("error", reject); server.listen(0, "127.0.0.1", () => resolve()); }); const address = server.address(); if (!address || typeof address === "string") throw new Error("Missing server address"); return address.port;}
function baseUrl(server: http.Server): string { const address = server.address(); if (!address || typeof address === "string") throw new Error("Missing server address"); return `http://127.0.0.1:${address.port}`;}
async function unusedPort(): Promise<number> { const server = http.createServer(); const port = await listen(server); await new Promise<void>((resolve) => server.close(() => resolve())); return port;}