File
Blob: src/worker/index.ts
| 1 | import { routePartykitRequest } from "partyserver"; |
| 2 | import { eq } from "drizzle-orm"; |
| 3 | import { app } from "@/worker/router"; |
| 4 | import { createSessionDb } from "@/worker/db/d1/client"; |
| 5 | import { pages, users } from "@/worker/db/d1/schema"; |
| 6 | import { handleHttpRequest } from "@/worker/lib/http-entry"; |
| 7 | import { sitesApp } from "@/worker/sites/router"; |
| 8 | import { isWriterRole, resolvePageAccessLevels, resolvePrincipal } from "@/worker/lib/permissions"; |
| 9 | import { verifyAccessToken } from "@/worker/lib/auth"; |
| 10 | import { createLogger, errorContext, setLevel } from "@/worker/lib/logger"; |
| 11 | import { isAllowedOrigin } from "@/worker/lib/origins"; |
| 12 | import { applyBaselineSecurityHeaders } from "@/worker/lib/security-headers"; |
| 13 | import { renderSpaShell } from "@/worker/lib/spa-shell"; |
| 14 | import type { TasksQueueMessage, TasksQueueResult } from "@/worker/queues/messages"; |
| 15 | import { handlePageProjection } from "@/worker/queues/page-projection"; |
| 16 | import { handleSearchIndexMessage } from "@/worker/queues/search-indexer"; |
| 17 | import { handleSiteCover } from "@/worker/queues/site-cover"; |
| 18 | import { handleWorkspaceSitesCleanup } from "@/worker/queues/workspace-sites-cleanup"; |
| 19 | |
| 20 | export { DocSync } from "@/worker/durable-objects/doc-sync"; |
| 21 | export { WorkspaceIndexer } from "@/worker/durable-objects/workspace-indexer"; |
| 22 | |
| 23 | const log = createLogger("websocket"); |
| 24 | |
| 25 | async function handlePartyRequest(request: Request, env: Cloudflare.Env) { |
| 26 | const response = await routePartykitRequest(request, env, { |
| 27 | onBeforeConnect: async (req, lobby) => { |
| 28 | const pageId = lobby.name; |
| 29 | log.debug("connection_attempt", { pageId }); |
| 30 | |
| 31 | // Validate browser-provided origins before upgrading the socket. |
| 32 | const origin = req.headers.get("origin"); |
| 33 | if (origin && !isAllowedOrigin(origin, env)) { |
| 34 | log.warn("origin_rejected", { origin, pageId }); |
| 35 | return new Response("Forbidden origin", { status: 403 }); |
| 36 | } |
| 37 | |
| 38 | const url = new URL(req.url); |
| 39 | const token = url.searchParams.get("token"); |
| 40 | const shareToken = url.searchParams.get("share"); |
| 41 | |
| 42 | if (!token && !shareToken) { |
| 43 | log.warn("auth_missing", { pageId }); |
| 44 | return new Response("Authentication required", { status: 401 }); |
| 45 | } |
| 46 | |
| 47 | const { db } = createSessionDb(env.DB, "first-primary"); |
| 48 | |
| 49 | // Load page first — needed for both auth paths |
| 50 | const page = await db |
| 51 | .select({ workspace_id: pages.workspace_id, archived_at: pages.archived_at }) |
| 52 | .from(pages) |
| 53 | .where(eq(pages.id, pageId)) |
| 54 | .get(); |
| 55 | |
| 56 | if (!page || page.archived_at) { |
| 57 | log.warn("page_not_found", { pageId }); |
| 58 | return new Response("Page not found", { status: 404 }); |
| 59 | } |
| 60 | |
| 61 | // Resolve viewer identity. JWT is optional; share-token is authoritative for |
| 62 | // shared-surface connections so `/s/:token` stays link-scoped even when the |
| 63 | // browser also carries a bearer (mirrors HTTP shared-follow-on routes). |
| 64 | let authedUser: { id: string } | null = null; |
| 65 | if (token) { |
| 66 | try { |
| 67 | const { sub } = await verifyAccessToken(token, env); |
| 68 | const row = await db.select({ id: users.id }).from(users).where(eq(users.id, sub)).get(); |
| 69 | if (!row) { |
| 70 | log.warn("auth_failed", { pageId, reason: "user_not_found" }); |
| 71 | return new Response("Invalid token", { status: 401 }); |
| 72 | } |
| 73 | authedUser = { id: row.id }; |
| 74 | } catch { |
| 75 | log.warn("auth_failed", { pageId }); |
| 76 | return new Response("Invalid token", { status: 401 }); |
| 77 | } |
| 78 | } |
| 79 | |
| 80 | const surface = shareToken ? "shared" : "canonical"; |
| 81 | const resolved = await resolvePrincipal(db, authedUser, page.workspace_id, { |
| 82 | surface, |
| 83 | shareToken: shareToken ?? undefined, |
| 84 | }); |
| 85 | if (!resolved) { |
| 86 | log.warn("auth_missing", { pageId, reason: "principal_unresolved" }); |
| 87 | return new Response("Authentication required", { status: 401 }); |
| 88 | } |
| 89 | |
| 90 | const accessLevels = await resolvePageAccessLevels(db, resolved.principal, [pageId], page.workspace_id); |
| 91 | const accessLevel = accessLevels.get(pageId) ?? "none"; |
| 92 | if (accessLevel === "none") { |
| 93 | log.warn("access_denied", { pageId, surface, principalType: resolved.principal.type }); |
| 94 | return new Response( |
| 95 | surface === "shared" ? "Invalid or expired share link" : "You do not have access to this page", |
| 96 | { status: 403 }, |
| 97 | ); |
| 98 | } |
| 99 | const readOnly = accessLevel !== "edit"; |
| 100 | |
| 101 | log.info("connection_authorized", { pageId, surface, principalType: resolved.principal.type, readOnly }); |
| 102 | |
| 103 | // Always return a sanitized Request to prevent client-injected params. The |
| 104 | // `member_edit` tag is reserved for canonical writers (owner/admin/member) so |
| 105 | // the DO can hand out edit headroom only to those connections. Guests reach |
| 106 | // pages via share grants, not the writer fast path, so they do not qualify. |
| 107 | url.searchParams.delete("readOnly"); |
| 108 | url.searchParams.delete("authType"); |
| 109 | if (readOnly) { |
| 110 | url.searchParams.set("readOnly", "1"); |
| 111 | } |
| 112 | if (surface === "canonical" && !readOnly && isWriterRole(resolved.workspaceRole)) { |
| 113 | url.searchParams.set("authType", "member_edit"); |
| 114 | } |
| 115 | return new Request(url.toString(), req); |
| 116 | }, |
| 117 | }); |
| 118 | |
| 119 | return response && response.status !== 101 ? applyBaselineSecurityHeaders(response) : response; |
| 120 | } |
| 121 | |
| 122 | export default { |
| 123 | async fetch(request: Request, env: Cloudflare.Env, ctx: ExecutionContext) { |
| 124 | setLevel(env.LOG_LEVEL); |
| 125 | return handleHttpRequest(request, env, ctx, { |
| 126 | handlePartyRequest, |
| 127 | handleAppRequest: (nextRequest, nextEnv, nextCtx) => app.fetch(nextRequest, nextEnv, nextCtx), |
| 128 | handleSiteRequest: (nextRequest, nextEnv, nextCtx) => sitesApp.fetch(nextRequest, nextEnv, nextCtx), |
| 129 | handleAssetRequest: async (nextRequest, nextEnv) => |
| 130 | applyBaselineSecurityHeaders(await nextEnv.ASSETS.fetch(nextRequest)), |
| 131 | handleShellRequest: renderSpaShell, |
| 132 | }); |
| 133 | }, |
| 134 | async queue(batch: MessageBatch<TasksQueueMessage>, env: Cloudflare.Env, _ctx: ExecutionContext) { |
| 135 | setLevel(env.LOG_LEVEL); |
| 136 | const log = createLogger("queue"); |
| 137 | |
| 138 | for (const msg of batch.messages) { |
| 139 | const body = msg.body; |
| 140 | try { |
| 141 | let result: TasksQueueResult = { kind: "ok" }; |
| 142 | switch (body.type) { |
| 143 | case "index-page": |
| 144 | result = await handleSearchIndexMessage({ type: "index-page", pageId: body.pageId }, env); |
| 145 | break; |
| 146 | case "page-projection": |
| 147 | result = await handlePageProjection(body.pageId, env); |
| 148 | break; |
| 149 | case "workspace-sites-cleanup": |
| 150 | result = await handleWorkspaceSitesCleanup(body.workspaceId, env); |
| 151 | break; |
| 152 | case "site-cover": |
| 153 | result = await handleSiteCover(body.pageId, env); |
| 154 | break; |
| 155 | default: |
| 156 | log.warn("unknown_message_type", { type: (body as { type?: string }).type }); |
| 157 | } |
| 158 | if (result.kind === "retry") { |
| 159 | msg.retry({ delaySeconds: result.delaySeconds }); |
| 160 | continue; |
| 161 | } |
| 162 | msg.ack(); |
| 163 | } catch (e) { |
| 164 | log.error("message_failed", { ...queueMessageLogContext(body), ...errorContext(e) }); |
| 165 | msg.retry(); |
| 166 | } |
| 167 | } |
| 168 | }, |
| 169 | } satisfies ExportedHandler<Cloudflare.Env, TasksQueueMessage>; |
| 170 | |
| 171 | function queueMessageLogContext(body: TasksQueueMessage): { type: string; pageId?: string; workspaceId?: string } { |
| 172 | switch (body.type) { |
| 173 | case "index-page": |
| 174 | case "page-projection": |
| 175 | case "site-cover": |
| 176 | return { type: body.type, pageId: body.pageId }; |
| 177 | case "workspace-sites-cleanup": |
| 178 | return { type: body.type, workspaceId: body.workspaceId }; |
| 179 | } |
| 180 | } |