diff --git a/bobbin/worker/src/index.ts b/bobbin/worker/src/index.ts index a46f47fd2..9000c8c61 100644 --- a/bobbin/worker/src/index.ts +++ b/bobbin/worker/src/index.ts @@ -1,7 +1,10 @@ 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; @@ -10,6 +13,69 @@ export interface Env { 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 { + 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; @@ -83,6 +149,16 @@ export default { ); } + 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" + ) { + return proxyHydrant(request, env); + } + const ip = request.headers.get("cf-connecting-ip") ?? "unknown"; const { success } = await env.RATE_LIMITER.limit({ key: ip }); if (!success) { diff --git a/bobbin/worker/wrangler.dev.jsonc b/bobbin/worker/wrangler.dev.jsonc index 2386fda1e..5a74dc9d5 100644 --- a/bobbin/worker/wrangler.dev.jsonc +++ b/bobbin/worker/wrangler.dev.jsonc @@ -2,7 +2,8 @@ "$schema": "node_modules/wrangler/config-schema.json", "name": "bobbin-svfe-dev", "main": "src/index.ts", - "workers_dev": false, + "workers_dev": true, + "preview_urls": false, "compatibility_date": "2026-03-07", "observability": { "enabled": true, @@ -28,6 +29,12 @@ }, ], }, + "vpc_services": [ + { + "binding": "HYDRANT", + "service_id": "01a05715-1857-7561-abab-104563fe3970", + }, + ], "migrations": [ { "tag": "v1", @@ -45,8 +52,8 @@ }, ], "vars": { - "BOBBIN_HYDRANT_URL": "https://hyd.bob.oyster.cafe", - "BOBBIN_SLINGSHOT_URL": "https://slng.bob.oyster.cafe", + "BOBBIN_HYDRANT_URL": "https://bobbin-svfe-dev.anirudh-s-account.workers.dev", + "BOBBIN_SLINGSHOT_URL": "https://bobbin-svfe-dev.anirudh-s-account.workers.dev", "BOBBIN_MIRROR_V2_URL": "https://mirror-fsn.tangled.network", "BOBBIN_SERVICE_DID": "did:web:next.tangled.org", "BOBBIN_LOG": "info", diff --git a/bobbin/worker/wrangler.jsonc b/bobbin/worker/wrangler.jsonc index 803bf81f0..aa9c63467 100644 --- a/bobbin/worker/wrangler.jsonc +++ b/bobbin/worker/wrangler.jsonc @@ -23,6 +23,12 @@ }, ], }, + "vpc_services": [ + { + "binding": "HYDRANT", + "service_id": "01a05715-1857-7561-abab-104563fe3970", + }, + ], "migrations": [ { "tag": "v1", @@ -45,8 +51,8 @@ }, ], "vars": { - "BOBBIN_HYDRANT_URL": "https://hyd.bob.oyster.cafe", - "BOBBIN_SLINGSHOT_URL": "https://slng.bob.oyster.cafe", + "BOBBIN_HYDRANT_URL": "https://api.tangled.org", + "BOBBIN_SLINGSHOT_URL": "https://api.tangled.org", "BOBBIN_MIRROR_V2_URL": "https://mirror-fsn.tangled.network", "BOBBIN_SERVICE_DID": "did:web:api.tangled.org", "BOBBIN_LOG": "info",