import * as Effect from "effect/Effect"; import { Secret } from "../configuration/secrets"; import { mcpHeaders, type McpHeader } from "../shared/mcp-headers"; import { agentCall, agentValidation } from "./agent-io"; import { decryptCredential, encryptCredential } from "./credential-crypto"; import { OperationFailure } from "./operation-result"; const purpose = "flarebot/mcp-headers/v1"; // Immutable encrypted sets make the SQL reference swap the commit point. // Preparing or replacing a set never needs to decrypt existing credentials. export class McpHeaders { constructor( private readonly storage: DurableObjectStorage, private readonly installationId: string, private readonly secret: Secret, ) {} prepare(id: string, value: unknown) { return Effect.gen({ self: this }, function* () { const headers = yield* agentValidation(() => { const parsed = mcpHeaders.safeParse(value); if (!parsed.success) throw new OperationFailure("mcp_headers_invalid"); return parsed.data; }); if (!headers.length) return { ref: null, names: [] }; const ref = `${id}/${crypto.randomUUID()}`; const ciphertext = yield* encryptCredential( new Secret(JSON.stringify(headers)), this.secret, this.installationId, purpose, ref, ); yield* agentCall(() => this.storage.put(`${purpose}/${ref}`, ciphertext)); return { ref, names: headers.map((header) => header.name) }; }); } resolve(ref: string | null) { return Effect.gen({ self: this }, function* () { if (!ref) return yield* Effect.fail( new OperationFailure("mcp_credentials_unreadable"), ); const ciphertext = yield* agentCall(() => this.storage.get(`${purpose}/${ref}`), ); if (ciphertext === undefined) return yield* Effect.fail( new OperationFailure("mcp_credentials_unreadable"), ); const credential = yield* decryptCredential( ciphertext, this.secret, this.installationId, purpose, ref, ); yield* agentValidation(() => mcpHeaders.parse(JSON.parse(credential.reveal())), ); return credential; }).pipe( Effect.catch(() => Effect.fail(new OperationFailure("mcp_credentials_unreadable")), ), ); } discard(ref: string | null) { return ref ? agentCall(() => this.storage.delete(`${purpose}/${ref}`)).pipe( Effect.asVoid, ) : Effect.void; } // Awaited inside parent startup before accepting writes: no in-flight prepared // set can be mistaken for an orphan during cleanup. prune(references: string[]) { return agentCall(async () => { const retained = new Set(references.map((ref) => `${purpose}/${ref}`)); const records = await this.storage.list({ prefix: `${purpose}/` }); const orphaned = [...records.keys()].filter((key) => !retained.has(key)); for (let offset = 0; offset < orphaned.length; offset += 128) await this.storage.delete(orphaned.slice(offset, offset + 128)); }); } } export function mcpHeaderFetch( endpoint: string, credentials: Secret, authorized: () => boolean, rejected: () => void, ): typeof fetch { const origin = new URL(endpoint).origin; const parsed: McpHeader[] = mcpHeaders.parse( JSON.parse(credentials.reveal()), ); return async (input, init) => { const request = new Request(input, init); if (!authorized() || new URL(request.url).origin !== origin) throw new Error("MCP credentials are no longer authorized"); const headers = new Headers(request.headers); for (const header of parsed) headers.set(header.name, header.value); // Never forward custom secrets through a redirect or an SSE endpoint on // another origin. The native client still owns protocol and cancellation. const response = await fetch( new Request(request, { headers, redirect: "manual" }), ); if (response.status >= 300 && response.status < 400) { await response.body?.cancel(); throw new Error("MCP credential requests cannot follow redirects"); } if (response.status === 401 || response.status === 403) rejected(); if (!response.ok) { await response.body?.cancel(); return new Response("MCP request failed", { status: response.status }); } return response; }; }