Something went wrong. Try again.
A local-first event pipeline for independent agents, built on Jazz.
Something went wrong. Try again.
36 kB · 830 lines
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831import fs from "node:fs/promises";import os from "node:os";import path from "node:path";import { afterEach, describe, expect, test } from "vitest";import { JoseKey, NodeOAuthClient, type NodeSavedSession, type NodeSavedState } from "@atproto/oauth-client-node";import { createInspectorOAuthAuth, flowCookieName, GenerationSessionStore, InspectorOAuthAuth, oauthConfigurationFromEnv, type BrowserSession, type OAuthProtocolClient, type PersistedOAuthSession, sessionCookieName,} from "../src/web/oauth-auth.js";import { SecureJsonStore } from "../src/web/secure-store.js";
const roots: string[] = [];
afterEach(async () => { await Promise.all(roots.splice(0).map((root) => fs.rm(root, { recursive: true, force: true })));});
describe("ATProto OAuth inspector authentication", () => { test("keeps SDK protocol state distinct from browser-bound application state", async () => { const fixture = await authFixture(); const started = await fixture.auth.begin(); const applicationState = cookieValue(started.setCookie, flowCookieName()); const protocolState = fixture.protocol.protocolState;
expect(applicationState).not.toBe(protocolState); expect(started.redirect.searchParams.get("request_uri")).toBe("urn:ietf:params:oauth:request_uri:fixture"); expect(fixture.protocol.authorizeCalls).toEqual([{ handle: "cameron.stream", applicationState, scope: "atproto", signal: undefined, }]);
const finished = await fixture.auth.finish( new URLSearchParams({ state: protocolState, code: "fixture-code" }), `${flowCookieName()}=${applicationState}`, ); expect(finished.redirect).toBe("/inspector/"); expect(fixture.protocol.callbackStates).toEqual([protocolState]); const sessionCookie = finished.setCookies.find((value) => value.startsWith(`${sessionCookieName()}=`)); expect(sessionCookie).toBeDefined(); expect(sessionCookie).not.toContain("did:plc:allowed"); expect(sessionCookie).not.toContain("fixture-code");
const sessionId = cookieValue(sessionCookie!, sessionCookieName()); const browser = await fixture.auth.authenticate(`${sessionCookieName()}=${sessionId}`); expect(browser).toMatchObject({ did: "did:plc:allowed", csrfToken: expect.any(String) }); expect(fixture.protocol.restoreCalls).toEqual(["did:plc:allowed"]);
await expect(fixture.auth.finish( new URLSearchParams({ state: protocolState, code: "replayed-code" }), `${flowCookieName()}=${applicationState}`, )).rejects.toThrow("protocol state replayed"); });
test("serializes concurrent callbacks so SDK protocol-state get/delete cannot race", async () => { const fixture = await authFixture(); const started = await fixture.auth.begin(); const applicationState = cookieValue(started.setCookie, flowCookieName()); let release!: () => void; fixture.protocol.callbackGate = new Promise<void>((resolve) => { release = resolve; }); const first = fixture.auth.finish(callbackParams(fixture.protocol), `${flowCookieName()}=${applicationState}`); const second = fixture.auth.finish(callbackParams(fixture.protocol), `${flowCookieName()}=${applicationState}`); await new Promise((resolve) => setTimeout(resolve, 5)); expect(fixture.protocol.callbackCalls).toBe(1); release(); const results = await Promise.allSettled([first, second]); expect(results.filter((result) => result.status === "fulfilled")).toHaveLength(1); expect(results.filter((result) => result.status === "rejected")).toHaveLength(1); expect(fixture.protocol.callbackCalls).toBe(2); });
test("watchdog advances past a never-resolving callback and permits a later callback", async () => { const protocol = new WatchdogProtocol(["never", "success"], true); const auth = await authWithProtocol(protocol, 20); const firstStart = await auth.begin(); const firstAppState = cookieValue(firstStart.setCookie, flowCookieName()); await expect(auth.finish( new URLSearchParams({ state: protocol.protocolState(0), code: "first" }), `${flowCookieName()}=${firstAppState}`, )).rejects.toThrow("timed out");
const secondStart = await auth.begin(); const secondAppState = cookieValue(secondStart.setCookie, flowCookieName()); const second = await auth.finish( new URLSearchParams({ state: protocol.protocolState(1), code: "second" }), `${flowCookieName()}=${secondAppState}`, ); const sessionCookie = second.setCookies.find((value) => value.startsWith(`${sessionCookieName()}=`)); expect(sessionCookie).toBeDefined(); expect(protocol.callbackCalls).toBe(2); expect(protocol.promotedAttempts).toHaveLength(1); });
test("watchdog advances past hung promotion and permits a later generation", async () => { const protocol = new WatchdogProtocol(["hung-promotion", "success"]); const auth = await authWithProtocol(protocol, 20); const firstStart = await auth.begin(); const firstAppState = cookieValue(firstStart.setCookie, flowCookieName()); await expect(auth.finish( new URLSearchParams({ state: protocol.protocolState(0), code: "first" }), `${flowCookieName()}=${firstAppState}`, )).rejects.toThrow("timed out");
const secondStart = await auth.begin(); const secondAppState = cookieValue(secondStart.setCookie, flowCookieName()); const second = await auth.finish( new URLSearchParams({ state: protocol.protocolState(1), code: "second" }), `${flowCookieName()}=${secondAppState}`, ); expect(second.setCookies.some((value) => value.startsWith(`${sessionCookieName()}=`))).toBe(true); expect(protocol.currentGeneration).toBe(1); });
test("hung non-timeout cleanup cannot block a subsequent successful callback", async () => { const protocol = new WatchdogProtocol(["success", "success"], true); const auth = await authWithProtocol(protocol, 50); await auth.begin(); await expect(auth.finish( new URLSearchParams({ state: protocol.protocolState(0), code: "first" }), `${flowCookieName()}=${"u".repeat(43)}`, )).rejects.toThrow("could not be verified");
const secondStart = await auth.begin(); const secondAppState = cookieValue(secondStart.setCookie, flowCookieName()); const second = await auth.finish( new URLSearchParams({ state: protocol.protocolState(1), code: "second" }), `${flowCookieName()}=${secondAppState}`, ); expect(second.setCookies.some((value) => value.startsWith(`${sessionCookieName()}=`))).toBe(true); expect(protocol.currentGeneration).toBe(1); });
test("late callback completion is discarded without revoking newer browser authority", async () => { const protocol = new WatchdogProtocol(["late", "success"]); const auth = await authWithProtocol(protocol, 20); const firstStart = await auth.begin(); const firstAppState = cookieValue(firstStart.setCookie, flowCookieName()); await expect(auth.finish( new URLSearchParams({ state: protocol.protocolState(0), code: "first" }), `${flowCookieName()}=${firstAppState}`, )).rejects.toThrow("timed out");
const secondStart = await auth.begin(); const secondAppState = cookieValue(secondStart.setCookie, flowCookieName()); const second = await auth.finish( new URLSearchParams({ state: protocol.protocolState(1), code: "second" }), `${flowCookieName()}=${secondAppState}`, ); const sessionCookie = second.setCookies.find((value) => value.startsWith(`${sessionCookieName()}=`))!; const sessionId = cookieValue(sessionCookie, sessionCookieName()); expect(await auth.authenticate(`${sessionCookieName()}=${sessionId}`)).toMatchObject({ did: "did:plc:allowed" });
protocol.resolveLate(); await waitFor(() => protocol.lateDiscarded); expect(await auth.authenticate(`${sessionCookieName()}=${sessionId}`)).toMatchObject({ did: "did:plc:allowed" }); expect(protocol.persistentDid).toBe("did:plc:allowed"); expect(protocol.currentGeneration).toBe(1); expect(protocol.remoteRevocations).toBe(0); });
test("rejects missing, duplicate, expired, or mismatched browser application state", async () => { const missing = await authFixture(); const missingStart = await missing.auth.begin(); await expect(missing.auth.finish( callbackParams(missing.protocol), undefined, )).rejects.toThrow("could not be verified"); expect(missing.protocol.callbackCalls).toBe(0);
const duplicate = await authFixture(); const duplicateStart = await duplicate.auth.begin(); const duplicateAppState = cookieValue(duplicateStart.setCookie, flowCookieName()); await expect(duplicate.auth.finish( callbackParams(duplicate.protocol), `${flowCookieName()}=${duplicateAppState}; ${flowCookieName()}=${duplicateAppState}`, )).rejects.toThrow("could not be verified"); expect(duplicate.protocol.callbackCalls).toBe(0);
const duplicateProtocol = await authFixture(); const duplicateProtocolStart = await duplicateProtocol.auth.begin(); const duplicateProtocolAppState = cookieValue(duplicateProtocolStart.setCookie, flowCookieName()); const duplicatedParams = callbackParams(duplicateProtocol.protocol); duplicatedParams.append("state", duplicateProtocol.protocol.protocolState); await expect(duplicateProtocol.auth.finish( duplicatedParams, `${flowCookieName()}=${duplicateProtocolAppState}`, )).rejects.toThrow("could not be verified"); expect(duplicateProtocol.protocol.callbackCalls).toBe(0);
const mismatch = await authFixture(); await mismatch.auth.begin(); const unrelatedState = "u".repeat(43); await expect(mismatch.auth.finish( callbackParams(mismatch.protocol), `${flowCookieName()}=${unrelatedState}`, )).rejects.toThrow("could not be verified"); expect(mismatch.protocol.callbackCalls).toBe(1); expect(mismatch.protocol.revoked).toEqual([]); expect(mismatch.protocol.deletedSessions).toEqual([]); expect(mismatch.protocol.discardedAttempts).toHaveLength(1);
let now = 1_000_000; const expired = await authFixture({ now: () => now }); const expiredStart = await expired.auth.begin(); const expiredAppState = cookieValue(expiredStart.setCookie, flowCookieName()); now += 15 * 60_000 + 1; await expect(expired.auth.finish( callbackParams(expired.protocol), `${flowCookieName()}=${expiredAppState}`, )).rejects.toThrow("could not be verified"); expect(expired.protocol.revoked).toEqual([]); expect(expired.protocol.deletedSessions).toEqual([]); expect(expired.protocol.discardedAttempts).toHaveLength(1); });
test("deletes staged SDK authority locally when flow-store settlement fails", async () => { const fixture = await authFixture(); const started = await fixture.auth.begin(); const applicationState = cookieValue(started.setCookie, flowCookieName()); fixture.flowStates.take = async () => { throw new Error("synthetic flow-store disk failure"); };
await expect(fixture.auth.finish( callbackParams(fixture.protocol), `${flowCookieName()}=${applicationState}`, )).rejects.toThrow("synthetic flow-store disk failure"); expect(fixture.protocol.revoked).toEqual([]); expect(fixture.protocol.deletedSessions).toEqual([]); expect(fixture.protocol.discardedAttempts).toHaveLength(1); expect(fixture.protocol.promotedAttempts).toHaveLength(0); });
test("rejects duplicate session cookies before SDK restore", async () => { const fixture = await authFixture(); const { sessionId } = await completeLogin(fixture); const cookie = `${sessionCookieName()}=${sessionId}`; expect(await fixture.auth.authenticate(`${cookie}; ${cookie}`)).toBeUndefined(); expect(fixture.protocol.restoreCalls).toEqual([]); });
test("expires browser sessions and deletes only their matching local generation", async () => { let now = 2_000_000; const fixture = await authFixture({ now: () => now, sessionTtlMs: 60_000 }); const { sessionId } = await completeLogin(fixture); now += 60_001;
expect(await fixture.auth.authenticate(`${sessionCookieName()}=${sessionId}`)).toBeUndefined(); expect(fixture.protocol.revoked).toEqual([]); expect(fixture.protocol.deletedSessions).toEqual([]); expect(fixture.protocol.discardedPromotions).toEqual([1]); expect(fixture.protocol.restoreCalls).toEqual([]); });
test("deletes the matching local generation after restore failure", async () => { const fixture = await authFixture(); const { sessionId } = await completeLogin(fixture); fixture.protocol.restoreFailure = true;
expect(await fixture.auth.authenticate(`${sessionCookieName()}=${sessionId}`)).toBeUndefined(); expect(fixture.protocol.restoreCalls).toEqual(["did:plc:allowed"]); expect(fixture.protocol.revoked).toEqual([]); expect(fixture.protocol.deletedSessions).toEqual([]); expect(fixture.protocol.discardedPromotions).toEqual([1]); expect(await fixture.auth.authenticate(`${sessionCookieName()}=${sessionId}`)).toBeUndefined(); });
test("logout clears browser authority and only its exact local generation", async () => { const fixture = await authFixture(); const { sessionId } = await completeLogin(fixture); const browser = await fixture.auth.authenticate(`${sessionCookieName()}=${sessionId}`);
const cleared = await fixture.auth.logout(`${sessionCookieName()}=${sessionId}`, browser!.csrfToken); expect(cleared).toContain("Max-Age=0"); expect(fixture.protocol.revoked).toEqual([]); expect(fixture.protocol.deletedSessions).toEqual([]); expect(fixture.protocol.discardedPromotions).toEqual([1]); expect(await fixture.auth.authenticate(`${sessionCookieName()}=${sessionId}`)).toBeUndefined(); });
test("drops a non-allowlisted staged DID without remote revocation", async () => { const fixture = await authFixture({ callbackDid: "did:plc:other" }); const started = await fixture.auth.begin(); const applicationState = cookieValue(started.setCookie, flowCookieName());
await expect(fixture.auth.finish( callbackParams(fixture.protocol), `${flowCookieName()}=${applicationState}`, )).rejects.toThrow("not authorized"); expect(fixture.protocol.revoked).toEqual([]); expect(fixture.protocol.deletedSessions).toEqual([]); expect(fixture.protocol.discardedAttempts).toHaveLength(1); });
test("scopes deferred restore refresh writes to their initiating generation", async () => { const { store } = await generationStoreFixture(); const did = "did:plc:allowed"; const generationOne = await stageAndPromote(store, "attempt-one", did, savedSession("one")); let release!: () => void; let entered!: () => void; const gate = new Promise<void>((resolve) => { release = resolve; }); const started = new Promise<void>((resolve) => { entered = resolve; }); const oldRestore = store.restoreExactGeneration(did, generationOne, async () => { expect(await store.get(did)).toEqual(savedSession("one")); entered(); await gate; await store.set(did, savedSession("one-refreshed")); }); await started; const generationTwo = await stageAndPromote(store, "attempt-two", did, savedSession("two")); release(); await expect(oldRestore).rejects.toThrow("changed during restore");
expect(await store.matchesPersistedGeneration(did, generationOne)).toBe(false); expect(await store.matchesPersistedGeneration(did, generationTwo)).toBe(true); expect(await store.get(did)).toEqual(savedSession("two")); });
test("scopes deferred restore failure deletion to its initiating generation", async () => { const { store } = await generationStoreFixture(); const did = "did:plc:allowed"; const generationOne = await stageAndPromote(store, "attempt-one", did, savedSession("one")); let release!: () => void; let entered!: () => void; const gate = new Promise<void>((resolve) => { release = resolve; }); const started = new Promise<void>((resolve) => { entered = resolve; }); const oldRestore = store.runRestore(did, generationOne, async () => { expect(await store.get(did)).toEqual(savedSession("one")); entered(); await gate; await store.del(did); throw new Error("synthetic generation-one refresh failure"); }); await started; const generationTwo = await stageAndPromote(store, "attempt-two", did, savedSession("two")); release(); await expect(oldRestore).rejects.toThrow("generation-one refresh failure");
expect(await store.matchesPersistedGeneration(did, generationOne)).toBe(false); expect(await store.matchesPersistedGeneration(did, generationTwo)).toBe(true); expect(await store.get(did)).toEqual(savedSession("two")); });
test("bounds retained callback quarantines and recovers after settlement or process recycle", async () => { const { persistent, store } = await generationStoreFixture(2); let releaseOne!: () => void; let releaseTwo!: () => void; const first = store.run("retained-one", () => new Promise<void>((resolve) => { releaseOne = resolve; })); const second = store.run("retained-two", () => new Promise<void>((resolve) => { releaseTwo = resolve; })); store.expire("retained-one"); store.expire("retained-two"); expect(store.quarantineStatus()).toEqual({ retained: 2, capacity: 2, recycleRequired: true }); expect(() => store.run("refused", async () => undefined)).toThrow("recycle the proxy process");
releaseOne(); await first; await store.drop("retained-one"); expect(store.quarantineStatus()).toEqual({ retained: 1, capacity: 2, recycleRequired: false }); await store.run("accepted-after-settlement", async () => undefined); await store.drop("accepted-after-settlement"); releaseTwo(); await second; await store.drop("retained-two");
const restarted = new GenerationSessionStore(persistent, 2); expect(restarted.quarantineStatus()).toEqual({ retained: 0, capacity: 2, recycleRequired: false }); await restarted.run("accepted-after-recycle", async () => undefined); await restarted.drop("accepted-after-recycle"); });
test("fails closed when OAUTH_STORE_DIR escapes the systemd writable path", () => { const home = "/tmp/thoughtstream-fixture-home"; const expected = `${home}/.local/share/thoughtstream-inspector-auth`; const base = { HOME: home, OAUTH_ENABLED: "1", OAUTH_PUBLIC_ORIGIN: "https://thought.stream", OAUTH_ALLOWED_DID: "did:plc:allowed", OAUTH_EXPECTED_HANDLE: "cameron.stream", OAUTH_STORE_KEY_B64: Buffer.alloc(32, 16).toString("base64"), OAUTH_PRIVATE_JWK_B64: Buffer.from(JSON.stringify({ kty: "EC", d: "private-fixture" })).toString("base64"), }; expect(oauthConfigurationFromEnv({ ...base, OAUTH_STORE_DIR: expected }).storeDirectory).toBe(expected); expect(() => oauthConfigurationFromEnv({ ...base, OAUTH_STORE_DIR: `${home}/custom` })) .toThrow("systemd sandbox permits only"); });
test("the installed official SDK generates protocol state separately from appState", async () => { const key = await JoseKey.generate(["ES256"], "fixture-key"); const states = new Map<string, NodeSavedState>(); const sessions = new Map<string, NodeSavedSession>(); const stateStore = mapStore(states); const sessionStore = mapStore(sessions); const client = new NodeOAuthClient({ clientMetadata: { client_id: "https://thought.stream/oauth/client-metadata.json", client_name: "thought stream inspector", client_uri: "https://thought.stream/", redirect_uris: ["https://thought.stream/oauth/callback"], grant_types: ["authorization_code", "refresh_token"], scope: "atproto", response_types: ["code"], application_type: "web", token_endpoint_auth_method: "private_key_jwt", token_endpoint_auth_signing_alg: "ES256", dpop_bound_access_tokens: true, jwks_uri: "https://thought.stream/oauth/jwks.json", }, keyset: [key], stateStore, sessionStore, requestLock: async (_name, operation) => operation(), }); const authorizationMetadata = { issuer: "https://auth.example", authorization_endpoint: "https://auth.example/authorize", token_endpoint: "https://auth.example/token", pushed_authorization_request_endpoint: "https://auth.example/par", require_pushed_authorization_requests: true, token_endpoint_auth_methods_supported: ["private_key_jwt"], token_endpoint_auth_signing_alg_values_supported: ["ES256"], dpop_signing_alg_values_supported: ["ES256"], scopes_supported: ["atproto"], response_types_supported: ["code"], grant_types_supported: ["authorization_code", "refresh_token"], code_challenge_methods_supported: ["S256"], authorization_response_iss_parameter_supported: true, }; let pushedState: string | undefined; (client.oauthResolver as unknown as { resolve: () => Promise<unknown> }).resolve = async () => ({ identityInfo: { did: "did:plc:allowed", handle: "cameron.stream" }, metadata: authorizationMetadata, }); (client.serverFactory as unknown as { fromMetadata: () => Promise<unknown> }).fromMetadata = async () => ({ request: async (_endpoint: string, payload: { state?: string }) => { pushedState = payload.state; return { request_uri: "urn:ietf:params:oauth:request_uri:fixture", expires_in: 90 }; }, }); const applicationState = "a".repeat(43); const redirect = await client.authorize("did:plc:allowed", { state: applicationState, scope: "atproto" }); const [protocolState, stored] = [...states.entries()][0]!;
expect(protocolState).not.toBe(applicationState); expect(stored.appState).toBe(applicationState); expect(pushedState).toBe(protocolState); expect(redirect.searchParams.get("request_uri")).toBe("urn:ietf:params:oauth:request_uri:fixture"); });
test("official client metadata is exact, JWKS is public-only, and the store is singleton", async () => { const root = await temporaryRoot("thoughtstream-oauth-client-"); const key = await JoseKey.generate(["ES256"], "fixture-key"); const configuration = { publicOrigin: "https://thought.stream", allowedDid: "did:plc:allowed", expectedHandle: "cameron.stream", storeDirectory: root, storeKey: Buffer.alloc(32, 4), privateJwk: { ...key.privateJwk! }, }; const auth = await createInspectorOAuthAuth(configuration);
expect(auth.clientMetadata).toMatchObject({ client_id: "https://thought.stream/oauth/client-metadata.json", client_uri: "https://thought.stream/", redirect_uris: ["https://thought.stream/oauth/callback"], grant_types: ["authorization_code", "refresh_token"], scope: "atproto", response_types: ["code"], token_endpoint_auth_method: "private_key_jwt", token_endpoint_auth_signing_alg: "ES256", dpop_bound_access_tokens: true, jwks_uri: "https://thought.stream/oauth/jwks.json", }); expect(JSON.stringify(auth.jwks)).not.toContain('"d"'); expect(JSON.stringify(auth.clientMetadata)).not.toContain('"d"'); await expect(createInspectorOAuthAuth(configuration)).rejects.toThrow("already owned by process"); await auth.close(); const reacquired = await createInspectorOAuthAuth(configuration); await reacquired.close(); });});
interface FixtureOptions { callbackDid?: string; now?: () => number; sessionTtlMs?: number; callbackTimeoutMs?: number;}
async function authFixture(options: FixtureOptions = {}): Promise<{ auth: InspectorOAuthAuth; protocol: FakeProtocol; flowStates: SecureJsonStore<{ createdAt: number }>;}> { const root = await temporaryRoot("thoughtstream-oauth-auth-"); const key = Buffer.alloc(32, 3); const now = options.now ?? Date.now; 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, now, }); const protocol = new FakeProtocol(options.callbackDid ?? "did:plc:allowed"); return { auth: new InspectorOAuthAuth({ protocol, allowedDid: "did:plc:allowed", expectedHandle: "cameron.stream", browserSessions, flowStates, now, ...(options.sessionTtlMs ? { sessionTtlMs: options.sessionTtlMs } : {}), ...(options.callbackTimeoutMs ? { callbackTimeoutMs: options.callbackTimeoutMs } : {}), }), protocol, flowStates, };}
class FakeProtocol implements OAuthProtocolClient { readonly clientMetadata = { client_id: "https://thought.stream/oauth/client-metadata.json" }; readonly jwks = { keys: [] }; readonly protocolState = "sdk-protocol-state-1234567890abcdef"; readonly authorizeCalls: Array<{ handle: string; applicationState: string; scope: string; signal: AbortSignal | undefined }> = []; readonly callbackStates: string[] = []; readonly restoreCalls: string[] = []; readonly revoked: string[] = []; readonly deletedSessions: string[] = []; readonly discardedAttempts: string[] = []; readonly expiredAttempts: string[] = []; readonly discardedPromotions: number[] = []; readonly promotedAttempts: string[] = []; callbackCalls = 0; restoreFailure = false; revokeFailure = false; callbackGate?: Promise<void>; private applicationState?: string; private callbackUsed = false; private readonly attemptDids = new Map<string, string>(); private currentGeneration: number | undefined; private nextGeneration = 0;
constructor(private readonly callbackDid: string) {}
async authorize(handle: string, options: { state: string; scope: string; signal?: AbortSignal }): Promise<URL> { this.applicationState = options.state; this.authorizeCalls.push({ handle, applicationState: options.state, scope: options.scope, signal: options.signal }); 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; const protocolState = params.get("state"); this.callbackStates.push(protocolState ?? ""); if (this.callbackUsed) throw new Error("protocol state replayed"); if (protocolState !== this.protocolState) throw new Error("unknown protocol state"); await this.callbackGate; this.callbackUsed = true; this.attemptDids.set(attemptId, this.callbackDid); return { session: { did: this.callbackDid }, state: this.applicationState ?? null }; }
expireCallbackAttempt(attemptId: string): void { this.expiredAttempts.push(attemptId); }
async dropCallbackAttempt(attemptId: string): Promise<void> { this.discardedAttempts.push(attemptId); this.attemptDids.delete(attemptId); }
async promoteCallbackSession(attemptId: string, did: string, guard: () => boolean): Promise<number> { if (!guard() || this.attemptDids.get(attemptId) !== did) throw new Error("missing staged session"); this.attemptDids.delete(attemptId); this.promotedAttempts.push(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> { this.discardedPromotions.push(generation); if (this.currentGeneration === generation) this.currentGeneration = undefined; }
async restore(did: string, generation: number): Promise<{ did: string }> { this.restoreCalls.push(did); if (this.restoreFailure || this.currentGeneration !== generation) throw new Error("synthetic restore failure"); return { did }; }
async revoke(did: string): Promise<void> { this.revoked.push(did); if (this.revokeFailure) throw new Error("synthetic revoke failure"); }
async deleteSession(did: string): Promise<void> { this.deletedSessions.push(did); this.currentGeneration = undefined; }}
async function generationStoreFixture(maxRetainedAttempts = 8): Promise<{ persistent: SecureJsonStore<PersistedOAuthSession>; store: GenerationSessionStore;}> { const root = await temporaryRoot("thoughtstream-generation-store-"); const persistent = new SecureJsonStore<PersistedOAuthSession>({ directory: root, name: "generation-sessions", key: Buffer.alloc(32, 20), maxEntries: 8, maxSerializedBytes: 64 * 1024, }); return { persistent, store: new GenerationSessionStore(persistent, maxRetainedAttempts) };}
async function stageAndPromote( store: GenerationSessionStore, attemptId: string, did: string, session: NodeSavedSession,): Promise<number> { await store.run(attemptId, () => store.set(did, session)); return store.promote(attemptId, did, () => true);}
function savedSession(label: string): NodeSavedSession { return { fixture: label } as unknown as NodeSavedSession;}
async function authWithProtocol(protocol: OAuthProtocolClient, callbackTimeoutMs: number): Promise<InspectorOAuthAuth> { const root = await temporaryRoot("thoughtstream-oauth-watchdog-"); const key = Buffer.alloc(32, 14); return new InspectorOAuthAuth({ protocol, allowedDid: "did:plc:allowed", expectedHandle: "cameron.stream", browserSessions: new SecureJsonStore<BrowserSession>({ directory: root, name: "watchdog-browser", key, maxEntries: 8, maxSerializedBytes: 32 * 1024, }), flowStates: new SecureJsonStore<{ createdAt: number }>({ directory: root, name: "watchdog-flow", key, maxEntries: 64, maxSerializedBytes: 32 * 1024, ttlMs: 15 * 60_000, }), callbackTimeoutMs, });}
class WatchdogProtocol implements OAuthProtocolClient { readonly clientMetadata = { client_id: "https://thought.stream/oauth/client-metadata.json" }; readonly jwks = { keys: [] }; readonly promotedAttempts: string[] = []; readonly discardedAttempts: string[] = []; callbackCalls = 0; lateDiscarded = false; persistentDid: string | undefined; currentGeneration: number | undefined; remoteRevocations = 0; private nextGeneration = 0; private readonly expiredAttempts = new Set<string>(); private readonly applicationStates: string[] = []; private readonly stagedAttempts = new Map<string, { did: string; index: number }>(); private readonly attemptIndexes = new Map<string, number>(); private callbackIndex = 0; private lateResolve!: () => void; private readonly latePromise = new Promise<void>((resolve) => { this.lateResolve = resolve; });
constructor( private readonly behaviors: Array<"never" | "late" | "success" | "hung-promotion">, private readonly hangFirstDrop = false, ) {}
protocolState(index: number): string { return `watchdog-protocol-state-0000000000000000-${index}`; }
resolveLate(): void { this.lateResolve(); }
async authorize(_handle: string, options: { state: string; scope: string; signal?: AbortSignal }): Promise<URL> { this.applicationStates.push(options.state); 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 }> { const index = this.callbackIndex++; this.attemptIndexes.set(attemptId, index); this.callbackCalls += 1; if (params.get("state") !== this.protocolState(index)) throw new Error("unexpected protocol state"); const behavior = this.behaviors[index]; if (behavior === "never") return new Promise(() => undefined); if (behavior === "late") await this.latePromise; this.stagedAttempts.set(attemptId, { did: "did:plc:allowed", index }); return { session: { did: "did:plc:allowed" }, state: this.applicationStates[index] ?? null }; }
expireCallbackAttempt(attemptId: string): void { this.expiredAttempts.add(attemptId); }
async dropCallbackAttempt(attemptId: string): Promise<void> { this.discardedAttempts.push(attemptId); const index = this.attemptIndexes.get(attemptId); if (index === 0 && this.hangFirstDrop) return new Promise(() => undefined); const staged = this.stagedAttempts.get(attemptId); this.stagedAttempts.delete(attemptId); if (staged?.index === 0) this.lateDiscarded = true; }
async promoteCallbackSession(attemptId: string, did: string, guard: () => boolean): Promise<number> { const index = this.attemptIndexes.get(attemptId); if (index !== undefined && this.behaviors[index] === "hung-promotion") return new Promise(() => undefined); if (!guard() || this.expiredAttempts.has(attemptId) || this.stagedAttempts.get(attemptId)?.did !== did) { throw new Error("missing staged session"); } this.stagedAttempts.delete(attemptId); this.promotedAttempts.push(attemptId); this.persistentDid = did; this.currentGeneration = ++this.nextGeneration; return this.currentGeneration; }
isCurrentGeneration(did: string, generation: number): boolean { return this.persistentDid === did && this.currentGeneration === generation; }
async discardPromotedSession(did: string, generation: number): Promise<void> { if (this.persistentDid === did && this.currentGeneration === generation) { this.persistentDid = undefined; this.currentGeneration = undefined; } }
async restore(did: string, generation: number): Promise<{ did: string }> { if (this.persistentDid !== did || this.currentGeneration !== generation) throw new Error("missing persistent session"); return { did }; }
async revoke(did: string): Promise<void> { this.remoteRevocations += 1; if (this.persistentDid === did) this.persistentDid = undefined; }
async deleteSession(did: string): Promise<void> { if (this.persistentDid === did) { this.persistentDid = undefined; this.currentGeneration = undefined; } }}
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 mapStore<T>(map: Map<string, T>): { get(key: string): Promise<T | undefined>; set(key: string, value: T): Promise<void>; del(key: string): Promise<void>;} { return { get: async (key) => map.get(key), set: async (key, value) => { map.set(key, value); }, del: async (key) => { map.delete(key); }, };}
async function completeLogin(fixture: { auth: InspectorOAuthAuth; protocol: FakeProtocol }): Promise<{ sessionId: string }> { const started = await fixture.auth.begin(); const applicationState = cookieValue(started.setCookie, flowCookieName()); const finished = await fixture.auth.finish( callbackParams(fixture.protocol), `${flowCookieName()}=${applicationState}`, ); const sessionCookie = finished.setCookies.find((value) => value.startsWith(`${sessionCookieName()}=`)); if (!sessionCookie) throw new Error("Missing session cookie"); return { sessionId: cookieValue(sessionCookie, sessionCookieName()) };}
function callbackParams(protocol: FakeProtocol): URLSearchParams { return new URLSearchParams({ state: protocol.protocolState, code: "fixture-code" });}
function cookieValue(serialized: string, name: string): string { const first = serialized.split(";", 1)[0]!; const [cookieName, value] = first.split("=", 2); if (cookieName !== name || !value) throw new Error(`Missing ${name} cookie`); return value;}
async function temporaryRoot(prefix: string): Promise<string> { const root = await fs.mkdtemp(path.join(os.tmpdir(), prefix)); roots.push(root); return root;}