import { Container } from "@cloudflare/containers"; export { ContainerProxy } from "@cloudflare/containers"; export interface Env { BOBBIN: DurableObjectNamespace; HYDRANT: Fetcher; RATE_LIMITER: RateLimit; BOBBIN_HYDRANT_URL: string; BOBBIN_SLINGSHOT_URL: string; BOBBIN_MIRROR_URL: string; BOBBIN_MIRROR_V2_URL?: string; BOBBIN_POCKET_URL?: string; BOBBIN_SERVICE_DID?: string; BOBBIN_LOG: string; BOBBIN_CODESEARCH_ZOEKT_URL?: string; } function closeWebSocket(webSocket: WebSocket, code: number, reason: string) { const closeCode = [1005, 1006, 1015].includes(code) ? 1011 : code; try { webSocket.close(closeCode, reason); } catch { // The other close handler may already have closed this side. } } async function proxyHydrant(request: Request, env: Env): Promise { const target = new URL(request.url); target.protocol = "http:"; target.host = "hydrant.internal"; const response = await env.HYDRANT.fetch(new Request(target, request)); const upstream = response.webSocket; if (upstream === null) { return response; } const pair = new WebSocketPair(); const client = pair[0]; const server = pair[1]; upstream.accept(); server.accept(); upstream.addEventListener("message", async (event) => { try { const data = event.data instanceof Blob ? await event.data.arrayBuffer() : event.data; server.send(data); } catch { closeWebSocket(upstream, 1011, "Failed to forward message to Bobbin"); } }); server.addEventListener("message", async (event) => { try { const data = event.data instanceof Blob ? await event.data.arrayBuffer() : event.data; upstream.send(data); } catch { closeWebSocket(server, 1011, "Failed to forward message to Hydrant"); } }); upstream.addEventListener("close", (event) => closeWebSocket(server, event.code, event.reason), ); server.addEventListener("close", (event) => closeWebSocket(upstream, event.code, event.reason), ); upstream.addEventListener("error", () => closeWebSocket(server, 1011, "Hydrant WebSocket error"), ); server.addEventListener("error", () => closeWebSocket(upstream, 1011, "Bobbin WebSocket error"), ); return new Response(null, { status: response.status, headers: response.headers, webSocket: client, }); } export class BobbinContainer extends Container { defaultPort = 8090; enableInternet = true; // Bobbin maintains an in-memory index rebuilt from Hydrant replay on every // restart, so we want to avoid sleeping where possible. Override // onActivityExpired() to keep the container alive indefinitely. sleepAfter = "24h"; constructor(ctx: DurableObjectState<{}>, env: Env) { super(ctx, env, { envVars: { BOBBIN_HYDRANT_URL: env.BOBBIN_HYDRANT_URL, BOBBIN_SLINGSHOT_URL: env.BOBBIN_SLINGSHOT_URL, BOBBIN_MIRROR_URL: env.BOBBIN_MIRROR_URL, ...(env.BOBBIN_MIRROR_V2_URL ? { BOBBIN_MIRROR_V2_URL: env.BOBBIN_MIRROR_V2_URL } : {}), ...(env.BOBBIN_POCKET_URL ? { BOBBIN_POCKET_URL: env.BOBBIN_POCKET_URL } : {}), ...(env.BOBBIN_SERVICE_DID ? { BOBBIN_SERVICE_DID: env.BOBBIN_SERVICE_DID } : {}), BOBBIN_LOG: env.BOBBIN_LOG, BOBBIN_LOG_FORMAT: "json", ...(env.BOBBIN_CODESEARCH_ZOEKT_URL ? { BOBBIN_CODESEARCH_ZOEKT_URL: env.BOBBIN_CODESEARCH_ZOEKT_URL } : {}), }, }); } async onActivityExpired(): Promise { // Keep the container running; bobbin's in-memory index is expensive to // rebuild. Renew the timeout instead of stopping so we're pinged again later. this.renewActivityTimeout(); } onError(error: Error) { console.error("bobbin container error:", error); } } const INDEX = `This is bobbin, Tangled's stateless XRPC API service: https://tangled.org/tangled.org/core/tree/master/bobbin`; const CORS_HEADERS = { "Access-Control-Allow-Origin": "*", "Access-Control-Allow-Methods": "GET, HEAD, POST, OPTIONS", "Access-Control-Allow-Headers": "Content-Type, Authorization, atproto-proxy", "Access-Control-Max-Age": "86400", }; function withCors(response: Response): Response { const headers = new Headers(response.headers); for (const [name, value] of Object.entries(CORS_HEADERS)) { headers.set(name, value); } // Responses with these statuses must not carry a body; a non-null body // crashes workerd, so null it out. const body = [101, 204, 205, 304].includes(response.status) ? null : response.body; return new Response(body, { status: response.status, statusText: response.statusText, headers, }); } export default { async fetch(request: Request, env: Env): Promise { if (request.method === "OPTIONS") { return new Response(null, { status: 204, headers: CORS_HEADERS }); } const url = new URL(request.url); if (url.pathname === "/" || url.pathname === "") { return withCors( new Response(INDEX, { headers: { "Content-Type": "text/plain" } }), ); } if ( url.pathname === "/stream" || url.pathname === "/health" || url.pathname.startsWith("/repos/") || url.pathname === "/xrpc/com.atproto.repo.getRecord" || url.pathname === "/xrpc/com.bad-example.identity.resolveMiniDoc" || url.pathname === "/xrpc/blue.microcosm.identity.resolveMiniDoc" ) { return proxyHydrant(request, env); } const ip = request.headers.get("cf-connecting-ip") ?? "unknown"; const { success } = await env.RATE_LIMITER.limit({ key: ip }); if (!success) { return withCors( new Response( JSON.stringify({ error: "RateLimitExceeded", message: "too many requests, slow down", }), { status: 429, headers: { "Content-Type": "application/json", "Retry-After": "60", }, }, ), ); } const container = env.BOBBIN.getByName("primary-eu-v2", { locationHint: "weur", }); const response = await container.fetch(request); // A 101 is a protocol switch (e.g. WebSocket upgrade); return it untouched // so we don't strip the connection off the response. if (response.status === 101) { return response; } return withCors(response); }, } satisfies ExportedHandler;