Something went wrong. Try again.
Monorepo for Tangled tangled.org
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195import { Container } from "@cloudflare/containers";
export { ContainerProxy } from "@cloudflare/containers";
export interface Env { BOBBIN: DurableObjectNamespace<BobbinContainer>; HYDRANT: Fetcher; RATE_LIMITER: RateLimit; BOBBIN_HYDRANT_URL: string; BOBBIN_SLINGSHOT_URL: string; BOBBIN_MIRROR_URL?: string; BOBBIN_MIRROR_V2_URL?: string; BOBBIN_SERVICE_DID?: string; BOBBIN_LOG: 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<Response> { 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<Env> { 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, ...(env.BOBBIN_MIRROR_URL ? { BOBBIN_MIRROR_URL: env.BOBBIN_MIRROR_URL } : {}), ...(env.BOBBIN_MIRROR_V2_URL ? { BOBBIN_MIRROR_V2_URL: env.BOBBIN_MIRROR_V2_URL } : {}), ...(env.BOBBIN_SERVICE_DID ? { BOBBIN_SERVICE_DID: env.BOBBIN_SERVICE_DID } : {}), BOBBIN_LOG: env.BOBBIN_LOG, BOBBIN_LOG_FORMAT: "json", }, }); }
async onActivityExpired(): Promise<void> { // 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<Response> { 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"); 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<Env>;