Something went wrong. Try again.
A local-first event pipeline for independent agents, built on Jazz.
Something went wrong. Try again.
31 kB · 793 lines
TypeScript
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794import { AsyncLocalStorage } from "node:async_hooks";import { randomBytes, timingSafeEqual } from "node:crypto";import os from "node:os";import path from "node:path";import { JoseKey, NodeOAuthClient, type NodeSavedSession, type NodeSavedState,} from "@atproto/oauth-client-node";import { SecureJsonStore, decodeKey32, opaqueKey } from "./secure-store.js";import { SingleProcessDirectoryLock } from "./single-process-lock.js";
const SESSION_COOKIE = "__Host-thoughtstream_session";const FLOW_COOKIE = "__Host-thoughtstream_oauth";const OAUTH_STATE_TTL_MS = 15 * 60_000;const DEFAULT_BROWSER_SESSION_TTL_MS = 12 * 60 * 60_000;const SDK_STATE_MAX_ENTRIES = 128;const SDK_STATE_MAX_BYTES = 512 * 1024;const SDK_SESSION_MAX_ENTRIES = 8;const SDK_SESSION_MAX_BYTES = 512 * 1024;const APP_FLOW_MAX_ENTRIES = 64;const APP_FLOW_MAX_BYTES = 32 * 1024;const BROWSER_SESSION_MAX_ENTRIES = 8;const BROWSER_SESSION_MAX_BYTES = 32 * 1024;const DEFAULT_CALLBACK_TIMEOUT_MS = 20_000;
export interface BrowserSession { did: string; generation: number; csrfToken: string; expiresAt: number;}
export interface OAuthProtocolClient { readonly clientMetadata: unknown; readonly jwks: unknown; authorize(handle: string, options: { state: string; scope: string; signal?: AbortSignal }): Promise<URL>; callback(params: URLSearchParams, attemptId: string): Promise<{ session: { did: string }; state: string | null }>; expireCallbackAttempt(attemptId: string): void; dropCallbackAttempt(attemptId: string): Promise<void>; promoteCallbackSession(attemptId: string, did: string, guard: () => boolean): Promise<number>; isCurrentGeneration(did: string, generation: number): boolean; discardPromotedSession(did: string, generation: number): Promise<void>; restore(did: string, generation: number): Promise<{ did: string }>;}
export interface InspectorOAuthAuthOptions { protocol: OAuthProtocolClient; allowedDid: string; expectedHandle: string; browserSessions: SecureJsonStore<BrowserSession>; flowStates: SecureJsonStore<{ createdAt: number }>; processLock?: SingleProcessDirectoryLock; sessionTtlMs?: number; callbackTimeoutMs?: number; now?: () => number;}
export interface OAuthStartResult { redirect: URL; setCookie: string;}
export interface OAuthCallbackResult { redirect: string; setCookies: string[];}
export interface OAuthEnvironmentConfiguration { enabled: boolean; publicOrigin?: string; allowedDid?: string; expectedHandle?: string; storeDirectory?: string; storeKey?: Buffer; privateJwk?: Record<string, unknown>; sessionTtlMs?: number;}
class OAuthCallbackWatchdogError extends Error { constructor() { super("OAuth callback timed out"); }}
export interface PersistedOAuthSession { generation: number; session?: NodeSavedSession;}
interface CallbackAttempt { status: "active" | "expired" | "promoting"; sessions: Map<string, NodeSavedSession>;}
type SessionOperationContext = | { kind: "callback"; attemptId: string } | { kind: "restore"; did: string; generation: number };
export class OAuthCallbackQuarantineCapacityError extends Error { constructor(readonly capacity: number) { super(`OAuth callback quarantine capacity ${capacity} reached; recycle the proxy process`); }}
export class GenerationSessionStore { private readonly context = new AsyncLocalStorage<SessionOperationContext>(); private readonly attempts = new Map<string, CallbackAttempt>(); private readonly currentGenerations = new Map<string, number>(); private readonly lastGenerations = new Map<string, number>();
constructor( private readonly persistent: SecureJsonStore<PersistedOAuthSession>, private readonly maxRetainedAttempts = 8, ) { if (!Number.isSafeInteger(maxRetainedAttempts) || maxRetainedAttempts < 1 || maxRetainedAttempts > 128) { throw new Error("OAuth callback quarantine capacity is invalid"); } }
async get(did: string): Promise<NodeSavedSession | undefined> { const context = this.context.getStore(); if (context?.kind === "callback") { return structuredClone(this.attempts.get(context.attemptId)?.sessions.get(did)); } const record = await this.persistent.get(did); if (!record) return undefined; this.lastGenerations.set(did, Math.max(record.generation, this.lastGenerations.get(did) ?? 0)); if (context?.kind === "restore" && (context.did !== did || context.generation !== record.generation)) { return undefined; } if (!record.session) return undefined; this.currentGenerations.set(did, record.generation); return structuredClone(record.session); }
async set(did: string, value: NodeSavedSession): Promise<void> { const context = this.context.getStore(); if (context?.kind === "callback") { const attempt = this.attempts.get(context.attemptId); if (!attempt) throw new Error("Unknown OAuth callback attempt"); // Expired attempts remain as local quarantine. This avoids one known // session-store failure path, but the SDK may still revoke on issuer errors. attempt.sessions.set(did, structuredClone(value)); return; } if (context?.kind === "restore") { if (context.did !== did) return; const replaced = await this.persistent.replaceIf( did, (record) => record.generation === context.generation && record.session !== undefined, { generation: context.generation, session: value }, ); if (replaced) { this.currentGenerations.set(did, context.generation); this.lastGenerations.set(did, Math.max(context.generation, this.lastGenerations.get(did) ?? 0)); } return; } const current = await this.persistent.get(did); if (!current?.session) throw new Error("Cannot update an unknown persistent OAuth session"); this.currentGenerations.set(did, current.generation); await this.persistent.set( did, { generation: current.generation, session: value }, () => this.currentGenerations.get(did) === current.generation, ); this.currentGenerations.set(did, current.generation); this.lastGenerations.set(did, Math.max(current.generation, this.lastGenerations.get(did) ?? 0)); }
async del(did: string): Promise<void> { const context = this.context.getStore(); if (context?.kind === "callback") { this.attempts.get(context.attemptId)?.sessions.delete(did); return; } if (context?.kind === "restore") { if (context.did !== did) return; const replaced = await this.persistent.replaceIf( did, (record) => record.generation === context.generation, { generation: context.generation }, ); if (replaced && this.currentGenerations.get(did) === context.generation) { this.currentGenerations.delete(did); } return; } const current = await this.persistent.get(did); if (current) { await this.persistent.replaceIf( did, (record) => record.generation === current.generation, { generation: current.generation }, ); if (this.currentGenerations.get(did) === current.generation) this.currentGenerations.delete(did); } }
run<T>(attemptId: string, operation: () => Promise<T>): Promise<T> { if (this.attempts.has(attemptId)) throw new Error("Duplicate OAuth callback attempt"); if (this.attempts.size >= this.maxRetainedAttempts) { throw new OAuthCallbackQuarantineCapacityError(this.maxRetainedAttempts); } this.attempts.set(attemptId, { status: "active", sessions: new Map() }); return this.context.run({ kind: "callback", attemptId }, operation); }
runRestore<T>(did: string, generation: number, operation: () => Promise<T>): Promise<T> { return this.context.run({ kind: "restore", did, generation }, operation); }
async restoreExactGeneration<T>(did: string, generation: number, operation: () => Promise<T>): Promise<T> { if (!await this.matchesPersistedGeneration(did, generation)) throw new Error("OAuth session generation is stale"); const result = await this.runRestore(did, generation, operation); if (!await this.matchesPersistedGeneration(did, generation)) { throw new Error("OAuth session generation changed during restore"); } return result; }
quarantineStatus(): { retained: number; capacity: number; recycleRequired: boolean } { return { retained: this.attempts.size, capacity: this.maxRetainedAttempts, recycleRequired: this.attempts.size >= this.maxRetainedAttempts, }; }
expire(attemptId: string): void { const attempt = this.attempts.get(attemptId); if (attempt) attempt.status = "expired"; }
async drop(attemptId: string): Promise<void> { const attempt = this.attempts.get(attemptId); if (!attempt) return; attempt.status = "expired"; attempt.sessions.clear(); this.attempts.delete(attemptId); }
async promote(attemptId: string, did: string, guard: () => boolean): Promise<number> { const attempt = this.attempts.get(attemptId); if (!attempt || attempt.status !== "active" || !guard()) { throw new Error("OAuth callback attempt cannot be promoted"); } const session = attempt.sessions.get(did); if (!session) throw new Error("OAuth callback produced no staged session"); const current = await this.persistent.get(did); const generation = Math.max( current?.generation ?? 0, this.currentGenerations.get(did) ?? 0, this.lastGenerations.get(did) ?? 0, ) + 1; if (attempt.status !== "active" || !guard()) throw new Error("OAuth callback attempt cannot be promoted"); attempt.status = "promoting"; try { await this.persistent.set( did, { generation, session }, () => attempt.status === "promoting" && guard() && generation > (this.currentGenerations.get(did) ?? 0), ); if (attempt.status !== "promoting" || !guard()) { await this.persistent.replaceIf( did, (record) => record.generation === generation, { generation }, ); throw new Error("OAuth callback attempt expired during promotion"); } this.currentGenerations.set(did, generation); this.lastGenerations.set(did, generation); this.attempts.delete(attemptId); return generation; } catch (error) { attempt.status = "expired"; throw error; } }
isCurrentGeneration(did: string, generation: number): boolean { return this.currentGenerations.get(did) === generation; }
async matchesPersistedGeneration(did: string, generation: number): Promise<boolean> { const record = await this.persistent.get(did); if (!record?.session || record.generation !== generation) return false; this.currentGenerations.set(did, generation); this.lastGenerations.set(did, Math.max(generation, this.lastGenerations.get(did) ?? 0)); return true; }
async discardPromoted(did: string, generation: number): Promise<void> { await this.persistent.replaceIf( did, (record) => record.generation === generation, { generation }, ); if (this.currentGenerations.get(did) === generation) this.currentGenerations.delete(did); this.lastGenerations.set(did, Math.max(generation, this.lastGenerations.get(did) ?? 0)); }}
export class InspectorOAuthAuth { private readonly now: () => number; private readonly sessionTtlMs: number; private readonly callbackTimeoutMs: number; private callbackOperation = Promise.resolve();
constructor(private readonly options: InspectorOAuthAuthOptions) { if (!/^did:(plc|web):/.test(options.allowedDid)) throw new Error("OAuth allowlisted DID is invalid"); if (!/^[a-z0-9][a-z0-9.-]+$/i.test(options.expectedHandle)) throw new Error("OAuth expected handle is invalid"); this.now = options.now ?? Date.now; this.sessionTtlMs = options.sessionTtlMs ?? DEFAULT_BROWSER_SESSION_TTL_MS; if (!Number.isSafeInteger(this.sessionTtlMs) || this.sessionTtlMs < 60_000 || this.sessionTtlMs > 7 * 86_400_000) { throw new Error("OAuth browser session TTL must be between one minute and seven days"); } this.callbackTimeoutMs = options.callbackTimeoutMs ?? DEFAULT_CALLBACK_TIMEOUT_MS; if (!Number.isSafeInteger(this.callbackTimeoutMs) || this.callbackTimeoutMs < 10 || this.callbackTimeoutMs > 120_000) { throw new Error("OAuth callback timeout must be between 10ms and two minutes"); } }
get clientMetadata(): unknown { return this.options.protocol.clientMetadata; }
get jwks(): unknown { return this.options.protocol.jwks; }
async begin(signal?: AbortSignal): Promise<OAuthStartResult> { const state = randomBytes(32).toString("base64url"); await this.options.flowStates.set(state, { createdAt: this.now() }); let redirect: URL; try { redirect = await this.options.protocol.authorize(this.options.expectedHandle, { state, scope: "atproto", ...(signal ? { signal } : {}), }); } catch (error) { await this.options.flowStates.del(state); throw error; } return { redirect, setCookie: serializeCookie(FLOW_COOKIE, state, { maxAgeSeconds: Math.floor(OAUTH_STATE_TTL_MS / 1_000), }), }; }
async finish(params: URLSearchParams, cookieHeader: string | undefined): Promise<OAuthCallbackResult> { const previous = this.callbackOperation; let release!: () => void; const barrier = new Promise<void>((resolve) => { release = resolve; }); this.callbackOperation = previous.catch(() => undefined).then(() => barrier); await previous.catch(() => undefined); try { return await this.finishSerialized(params, cookieHeader); } finally { release(); } }
private async finishSerialized(params: URLSearchParams, cookieHeader: string | undefined): Promise<OAuthCallbackResult> { const protocolStates = params.getAll("state"); const applicationState = readUniqueCookie(cookieHeader, FLOW_COOKIE); if (protocolStates.length !== 1 || !validProtocolState(protocolStates[0]!) || !validCallbackMultiplicity(params) || !applicationState || !validOpaqueState(applicationState)) { throw new Error("OAuth callback could not be verified"); } const attemptId = randomBytes(24).toString("base64url"); let authoritative = true; const guard = () => authoritative; const settlement = this.settleCallbackAttempt(params, applicationState, attemptId, guard); const watchdog = callbackWatchdog(this.callbackTimeoutMs); const settled = await Promise.race([ settlement.then( (value) => ({ kind: "result" as const, value }), (error: unknown) => ({ kind: "failure" as const, error }), ), watchdog.promise, ]); watchdog.cancel(); if (settled.kind === "timeout") { authoritative = false; this.options.protocol.expireCallbackAttempt(attemptId); this.observeLateSettlement(settlement); this.deleteFlowStateDetached(applicationState); throw new OAuthCallbackWatchdogError(); } authoritative = false; if (settled.kind === "failure") throw settled.error; return settled.value; }
private async settleCallbackAttempt( params: URLSearchParams, applicationState: string, attemptId: string, guard: () => boolean, ): Promise<OAuthCallbackResult> { let did: string | undefined; let generation: number | undefined; let browserSessionKey: string | undefined; try { const result = await this.options.protocol.callback(params, attemptId); this.assertCallbackAuthority(guard); did = result.session.did; const returnedApplicationState = result.state; const pending = returnedApplicationState && constantTimeTextEqual(returnedApplicationState, applicationState) ? await this.options.flowStates.take(applicationState, guard) : undefined; this.assertCallbackAuthority(guard); if (!pending) throw new Error("OAuth callback could not be verified"); if (!constantTimeTextEqual(did, this.options.allowedDid)) throw new Error("OAuth identity is not authorized"); this.assertCallbackAuthority(guard); generation = await this.options.protocol.promoteCallbackSession(attemptId, did, guard); this.assertCallbackAuthority(guard); if (!this.options.protocol.isCurrentGeneration(did, generation)) { throw new Error("OAuth callback generation is no longer current"); } const sessionId = randomBytes(32).toString("base64url"); browserSessionKey = opaqueKey(sessionId); const browserSession: BrowserSession = { did, generation, csrfToken: randomBytes(32).toString("base64url"), expiresAt: this.now() + this.sessionTtlMs, }; await this.options.browserSessions.replaceAll( browserSessionKey, browserSession, () => guard() && this.options.protocol.isCurrentGeneration(did!, generation!), ); this.assertCallbackAuthority(guard); if (!this.options.protocol.isCurrentGeneration(did, generation)) { throw new Error("OAuth callback generation is no longer current"); } return { redirect: "/inspector/", setCookies: [ serializeCookie(SESSION_COOKIE, sessionId, { maxAgeSeconds: Math.floor(this.sessionTtlMs / 1_000) }), clearCookie(FLOW_COOKIE), ], }; } catch (error) { if (error instanceof OAuthCallbackQuarantineCapacityError) throw error; this.options.protocol.expireCallbackAttempt(attemptId); if (browserSessionKey && generation !== undefined) { this.detachCleanup(this.options.browserSessions.deleteIf( browserSessionKey, (session) => session.generation === generation, )); } if (did && generation !== undefined) { this.detachCleanup(this.options.protocol.discardPromotedSession(did, generation)); } this.detachCleanup(this.options.protocol.dropCallbackAttempt(attemptId)); this.deleteFlowStateDetached(applicationState); throw error; } }
private assertCallbackAuthority(guard: () => boolean): void { if (!guard()) throw new OAuthCallbackWatchdogError(); }
async authenticate(cookieHeader: string | undefined): Promise<BrowserSession | undefined> { const sessionId = readUniqueCookie(cookieHeader, SESSION_COOKIE); if (!sessionId || !validOpaqueState(sessionId)) return undefined; const key = opaqueKey(sessionId); const stored = await this.options.browserSessions.getWithExpiration(key); if (!stored) return undefined; const browserSession = stored.value; if (stored.expired || browserSession.expiresAt <= this.now() || !constantTimeTextEqual(browserSession.did, this.options.allowedDid)) { await this.options.browserSessions.del(key); await this.options.protocol.discardPromotedSession(browserSession.did, browserSession.generation); return undefined; } try { const restored = await this.options.protocol.restore(browserSession.did, browserSession.generation); if (!constantTimeTextEqual(restored.did, this.options.allowedDid)) throw new Error("restored DID mismatch"); return browserSession; } catch { await this.options.browserSessions.del(key); await this.options.protocol.discardPromotedSession(browserSession.did, browserSession.generation); return undefined; } }
async logout(cookieHeader: string | undefined, csrfToken: string | undefined): Promise<string> { const sessionId = readUniqueCookie(cookieHeader, SESSION_COOKIE); if (!sessionId || !csrfToken) throw new Error("Logout request could not be verified"); const key = opaqueKey(sessionId); const browserSession = await this.options.browserSessions.get(key); if (!browserSession || !constantTimeTextEqual(browserSession.csrfToken, csrfToken)) { throw new Error("Logout request could not be verified"); } await this.options.browserSessions.del(key); await this.options.protocol.discardPromotedSession(browserSession.did, browserSession.generation); return clearCookie(SESSION_COOKIE); }
private observeLateSettlement(settlement: Promise<OAuthCallbackResult>): void { void settlement.catch(() => undefined); }
private detachCleanup(cleanup: Promise<unknown>): void { void cleanup.catch(() => undefined); }
private deleteFlowStateDetached(applicationState: string): void { this.detachCleanup(this.options.flowStates.del(applicationState)); }
async close(): Promise<void> { await this.options.processLock?.release(); }
}
export async function createInspectorOAuthAuth( configuration: Required<Omit<OAuthEnvironmentConfiguration, "enabled" | "sessionTtlMs">> & { sessionTtlMs?: number },): Promise<InspectorOAuthAuth> { const origin = new URL(configuration.publicOrigin); if (origin.protocol !== "https:" || origin.username || origin.password || origin.pathname !== "/" || origin.search || origin.hash) { throw new Error("OAUTH_PUBLIC_ORIGIN must be one bare HTTPS origin"); } const key = await JoseKey.fromJWK(configuration.privateJwk); const processLock = await SingleProcessDirectoryLock.acquire(configuration.storeDirectory); try { const stateStore = new SecureJsonStore<NodeSavedState>({ directory: configuration.storeDirectory, name: "oauth-state", key: configuration.storeKey, maxEntries: SDK_STATE_MAX_ENTRIES, maxSerializedBytes: SDK_STATE_MAX_BYTES, ttlMs: OAUTH_STATE_TTL_MS, validate: objectValue as (value: unknown) => NodeSavedState, }); const persistedSessions = new SecureJsonStore<PersistedOAuthSession>({ directory: configuration.storeDirectory, name: "oauth-session", key: configuration.storeKey, maxEntries: SDK_SESSION_MAX_ENTRIES, maxSerializedBytes: SDK_SESSION_MAX_BYTES, validate: (value) => { const object = objectValue(value); if (!Number.isSafeInteger(object.generation) || Number(object.generation) < 1) { throw new Error("Invalid OAuth session generation"); } return { generation: Number(object.generation), ...(object.session === undefined ? {} : { session: objectValue(object.session) as NodeSavedSession }), }; }, }); const sessionStore = new GenerationSessionStore(persistedSessions); const flowStates = new SecureJsonStore<{ createdAt: number }>({ directory: configuration.storeDirectory, name: "browser-oauth-flow", key: configuration.storeKey, maxEntries: APP_FLOW_MAX_ENTRIES, maxSerializedBytes: APP_FLOW_MAX_BYTES, ttlMs: OAUTH_STATE_TTL_MS, validate: (value) => { const object = objectValue(value); if (!Number.isSafeInteger(object.createdAt)) throw new Error("Invalid OAuth flow state"); return { createdAt: Number(object.createdAt) }; }, }); const browserSessions = new SecureJsonStore<BrowserSession>({ directory: configuration.storeDirectory, name: "browser-session", key: configuration.storeKey, maxEntries: BROWSER_SESSION_MAX_ENTRIES, maxSerializedBytes: BROWSER_SESSION_MAX_BYTES, validate: browserSessionValue, }); await Promise.all([stateStore.initialize(), persistedSessions.initialize(), flowStates.initialize(), browserSessions.initialize()]); const clientId = new URL("/oauth/client-metadata.json", origin).href; const callback = new URL("/oauth/callback", origin).href; const jwksUri = new URL("/oauth/jwks.json", origin).href; const client = new NodeOAuthClient({ clientMetadata: { client_id: clientId, client_name: "thought stream inspector", client_uri: origin.href, redirect_uris: [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: jwksUri, }, keyset: [key], stateStore, sessionStore, requestLock: keyedRuntimeLock(), }); const protocol: OAuthProtocolClient = { clientMetadata: client.clientMetadata, jwks: client.jwks, authorize: (handle, options) => client.authorize(handle, options), callback: (params, attemptId) => sessionStore.run(attemptId, () => client.callback(params)), expireCallbackAttempt: (attemptId) => sessionStore.expire(attemptId), dropCallbackAttempt: (attemptId) => sessionStore.drop(attemptId), promoteCallbackSession: (attemptId, did, guard) => sessionStore.promote(attemptId, did, guard), isCurrentGeneration: (did, generation) => sessionStore.isCurrentGeneration(did, generation), discardPromotedSession: (did, generation) => sessionStore.discardPromoted(did, generation), restore: (did, generation) => sessionStore.restoreExactGeneration( did, generation, () => client.restore(did), ), }; return new InspectorOAuthAuth({ protocol, allowedDid: configuration.allowedDid, expectedHandle: configuration.expectedHandle, browserSessions, flowStates, processLock, ...(configuration.sessionTtlMs ? { sessionTtlMs: configuration.sessionTtlMs } : {}), }); } catch (error) { await processLock.release(); throw error; }}
export function oauthConfigurationFromEnv(env: NodeJS.ProcessEnv = process.env): OAuthEnvironmentConfiguration { const enabled = env.OAUTH_ENABLED === "1" || env.OAUTH_ENABLED === "true"; if (!enabled) return { enabled: false }; const privateJwk = parsePrivateJwk(env.OAUTH_PRIVATE_JWK_B64); const sessionTtlMs = env.OAUTH_SESSION_TTL_MS ? Number(env.OAUTH_SESSION_TTL_MS) : undefined; const home = env.HOME?.trim() || os.homedir(); const expectedStoreDirectory = path.resolve(home, ".local", "share", "thoughtstream-inspector-auth"); const storeDirectory = path.resolve(required(env.OAUTH_STORE_DIR, "OAUTH_STORE_DIR")); if (storeDirectory !== expectedStoreDirectory) { throw new Error(`OAUTH_STORE_DIR must be ${expectedStoreDirectory}; the systemd sandbox permits only that owner-only path`); } return { enabled: true, publicOrigin: required(env.OAUTH_PUBLIC_ORIGIN, "OAUTH_PUBLIC_ORIGIN"), allowedDid: required(env.OAUTH_ALLOWED_DID, "OAUTH_ALLOWED_DID"), expectedHandle: required(env.OAUTH_EXPECTED_HANDLE, "OAUTH_EXPECTED_HANDLE"), storeDirectory, storeKey: decodeKey32(env.OAUTH_STORE_KEY_B64, "OAUTH_STORE_KEY_B64"), privateJwk, ...(sessionTtlMs !== undefined ? { sessionTtlMs } : {}), };}
export function sessionCookieName(): string { return SESSION_COOKIE;}
export function flowCookieName(): string { return FLOW_COOKIE;}
function callbackWatchdog(timeoutMs: number): { promise: Promise<{ kind: "timeout" }>; cancel: () => void;} { let timer: NodeJS.Timeout | undefined; const promise = new Promise<{ kind: "timeout" }>((resolve) => { timer = setTimeout(() => resolve({ kind: "timeout" }), timeoutMs); timer.unref(); }); return { promise, cancel: () => { if (timer) clearTimeout(timer); timer = undefined; }, };}
function keyedRuntimeLock(): <T>(key: string, fn: () => T | PromiseLike<T>) => Promise<T> { const locks = new Map<string, Promise<unknown>>(); return async <T>(key: string, fn: () => T | PromiseLike<T>): Promise<T> => { const previous = locks.get(key) ?? Promise.resolve(); let release!: () => void; const barrier = new Promise<void>((resolve) => { release = resolve; }); const chain = previous.catch(() => undefined).then(() => barrier); locks.set(key, chain); await previous.catch(() => undefined); try { return await fn(); } finally { release(); if (locks.get(key) === chain) locks.delete(key); } };}
function serializeCookie(name: string, value: string, options: { maxAgeSeconds: number }): string { return `${name}=${value}; Path=/; Max-Age=${options.maxAgeSeconds}; HttpOnly; Secure; SameSite=Lax`;}
function clearCookie(name: string): string { return `${name}=; Path=/; Max-Age=0; HttpOnly; Secure; SameSite=Lax`;}
function readUniqueCookie(header: string | undefined, name: string): string | undefined { if (!header || header.length > 8_192) return undefined; const values = header.split(";").map((part) => part.trim()).flatMap((part) => { const index = part.indexOf("="); return index > 0 && part.slice(0, index) === name ? [part.slice(index + 1)] : []; }); return values.length === 1 ? values[0] : undefined;}
function validCallbackMultiplicity(params: URLSearchParams): boolean { const codes = params.getAll("code"); const errors = params.getAll("error"); return codes.length <= 1 && errors.length <= 1 && params.getAll("iss").length <= 1 && params.getAll("response").length <= 1 && !(codes.length === 1 && errors.length === 1);}
function validProtocolState(value: string): boolean { return value.length >= 16 && value.length <= 512 && /^[A-Za-z0-9._~-]+$/.test(value);}
function validOpaqueState(value: string): boolean { return /^[A-Za-z0-9_-]{43}$/.test(value);}
function constantTimeTextEqual(left: string, right: string): boolean { const leftBytes = Buffer.from(left, "utf8"); const rightBytes = Buffer.from(right, "utf8"); if (leftBytes.length !== rightBytes.length) return false; return timingSafeEqual(leftBytes, rightBytes);}
function parsePrivateJwk(encoded: string | undefined): Record<string, unknown> { if (!encoded) throw new Error("OAUTH_PRIVATE_JWK_B64 is required"); const compact = encoded.trim(); const bytes = Buffer.from(compact, "base64"); if (bytes.length === 0 || bytes.toString("base64") !== compact) { throw new Error("OAUTH_PRIVATE_JWK_B64 must be canonical base64"); } try { const value = JSON.parse(bytes.toString("utf8")); if (!value || typeof value !== "object" || Array.isArray(value) || typeof value.d !== "string") throw new Error("shape"); return value as Record<string, unknown>; } catch { throw new Error("OAUTH_PRIVATE_JWK_B64 must encode one private JWK object"); }}
function browserSessionValue(value: unknown): BrowserSession { const object = objectValue(value); if ( typeof object.did !== "string" || !Number.isSafeInteger(object.generation) || Number(object.generation) < 1 || typeof object.csrfToken !== "string" || !Number.isSafeInteger(object.expiresAt) ) { throw new Error("Invalid browser session value"); } return { did: object.did, generation: Number(object.generation), csrfToken: object.csrfToken, expiresAt: Number(object.expiresAt), };}
function objectValue(value: unknown): Record<string, unknown> { if (!value || typeof value !== "object" || Array.isArray(value)) throw new Error("Expected object value"); return value as Record<string, unknown>;}
function required(value: string | undefined, label: string): string { const trimmed = value?.trim(); if (!trimmed) throw new Error(`${label} is required`); return trimmed;}