Skip to content
File

Blob: src/worker/index.ts

typescript181 lines
1import { routePartykitRequest } from "partyserver";
2import { eq } from "drizzle-orm";
3import { app } from "@/worker/router";
4import { createSessionDb } from "@/worker/db/d1/client";
5import { pages, users } from "@/worker/db/d1/schema";
6import { handleHttpRequest } from "@/worker/lib/http-entry";
7import { sitesApp } from "@/worker/sites/router";
8import { isWriterRole, resolvePageAccessLevels, resolvePrincipal } from "@/worker/lib/permissions";
9import { verifyAccessToken } from "@/worker/lib/auth";
10import { createLogger, errorContext, setLevel } from "@/worker/lib/logger";
11import { isAllowedOrigin } from "@/worker/lib/origins";
12import { applyBaselineSecurityHeaders } from "@/worker/lib/security-headers";
13import { renderSpaShell } from "@/worker/lib/spa-shell";
14import type { TasksQueueMessage, TasksQueueResult } from "@/worker/queues/messages";
15import { handlePageProjection } from "@/worker/queues/page-projection";
16import { handleSearchIndexMessage } from "@/worker/queues/search-indexer";
17import { handleSiteCover } from "@/worker/queues/site-cover";
18import { handleWorkspaceSitesCleanup } from "@/worker/queues/workspace-sites-cleanup";
19 
20export { DocSync } from "@/worker/durable-objects/doc-sync";
21export { WorkspaceIndexer } from "@/worker/durable-objects/workspace-indexer";
22 
23const log = createLogger("websocket");
24 
25async 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 
122export 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 
171function 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}