Skip to content
File

Blob: src/worker/index.ts

typescript93 lines
1import { handleIncomingEmail } from "@/worker/email-handler";
2import { InboxWebSocket } from "@/worker/durable-objects/inbox-ws";
3import { createDb } from "@/worker/db";
4import { createLogger } from "@/worker/logger";
5import { app } from "@/worker/router";
6import { isAllowedWebSocketOrigin, withSecurityHeaders } from "@/worker/security";
7import { cleanupExpiredInboxes, consumeWebSocketTicket, getInboxByAddress } from "@/worker/services/inbox";
8 
9export { InboxWebSocket };
10const logger = createLogger("worker");
11 
12async function handleWebSocketUpgrade(request: Request, env: Env) {
13 const url = new URL(request.url);
14 const address = url.searchParams.get("address")?.trim().toLowerCase() ?? "";
15 const ticket = url.searchParams.get("ticket");
16 
17 if (!address || request.headers.get("upgrade") !== "websocket") {
18 logger.warn("websocket_upgrade_rejected", "Rejected websocket upgrade request", {
19 address,
20 reason: "missing_upgrade_or_address",
21 });
22 return new Response("Expected websocket upgrade", { status: 426 });
23 }
24 
25 if (!isAllowedWebSocketOrigin(request)) {
26 logger.warn("websocket_upgrade_rejected", "Rejected websocket upgrade request", {
27 address,
28 reason: "invalid_origin",
29 });
30 return new Response("Forbidden", { status: 403 });
31 }
32 
33 const ticketRecord = await consumeWebSocketTicket(env, ticket);
34 if (!ticketRecord || ticketRecord.address !== address) {
35 logger.warn("websocket_upgrade_rejected", "Rejected websocket upgrade request", {
36 address,
37 reason: "invalid_ticket",
38 });
39 return new Response("Unauthorized", { status: 401 });
40 }
41 
42 const session = ticketRecord.session;
43 
44 const inbox = await getInboxByAddress(env, address, createDb(env.DB.withSession("first-primary")));
45 if (!inbox) {
46 logger.warn("websocket_upgrade_rejected", "Rejected websocket upgrade request", {
47 address,
48 reason: "inbox_not_found",
49 });
50 return new Response("Inbox not found", { status: 404 });
51 }
52 
53 if (session.type === "user" && session.address !== address) {
54 logger.warn("websocket_upgrade_rejected", "Rejected websocket upgrade request", {
55 address,
56 reason: "forbidden_user_scope",
57 });
58 return new Response("Forbidden", { status: 403 });
59 }
60 
61 const stub = env.INBOX_WS.getByName(address);
62 return stub.fetch(request);
63}
64 
65export default {
66 async fetch(request, env, ctx) {
67 const url = new URL(request.url);
68 
69 if (url.pathname === "/ws") {
70 return handleWebSocketUpgrade(request, env);
71 }
72 
73 if (url.pathname.startsWith("/api/")) {
74 const response = await app.fetch(request, env, ctx);
75 return withSecurityHeaders(request, response);
76 }
77 
78 const response = await env.ASSETS.fetch(request);
79 return withSecurityHeaders(request, response);
80 },
81 
82 async email(message, env, ctx) {
83 await handleIncomingEmail(message, env, ctx);
84 },
85 
86 async scheduled(_controller, env) {
87 const result = await cleanupExpiredInboxes(env);
88 logger.info("scheduled_cleanup_completed", "Finished scheduled inbox cleanup", {
89 deleted: result.deleted,
90 });
91 },
92} satisfies ExportedHandler<Env>;