File
Blob: src/worker/index.ts
| 1 | export { AuthObject } from "@/worker/objects/auth-object"; |
| 2 | export { FileDavObject } from "@/worker/objects/file-dav-object"; |
| 3 | export { CalDavObject } from "@/worker/objects/cal-dav-object"; |
| 4 | export { CardDavObject } from "@/worker/objects/card-dav-object"; |
| 5 | |
| 6 | import { createApp } from "@/worker/app"; |
| 7 | import { isBlobGcMessage } from "@/worker/queues/blob-gc"; |
| 8 | import { handleRepairJob, isRepairJobMessage } from "@/worker/queues/repair-jobs"; |
| 9 | import type { QueueHandleResult } from "@/worker/queues/types"; |
| 10 | import { handleBlobGc } from "@/worker/r2/gc"; |
| 11 | import { createLogger, errorContext, setLevel } from "@/worker/util/logging"; |
| 12 | |
| 13 | const app = createApp(); |
| 14 | const queueLog = createLogger("queue"); |
| 15 | |
| 16 | type QueueMessage = |
| 17 | | Message<unknown> |
| 18 | | { body: unknown; ack: () => void; retry: (options?: { delaySeconds?: number }) => void }; |
| 19 | |
| 20 | export async function handleQueueMessage(env: Env, body: unknown, nowMs = Date.now()): Promise<QueueHandleResult> { |
| 21 | if (isBlobGcMessage(body)) return await handleBlobGc(env, body, nowMs); |
| 22 | if (isRepairJobMessage(body)) return await handleRepairJob(env, body, nowMs); |
| 23 | queueLog.warn("unknown_message_type", { |
| 24 | eventType: "unknown_message_type", |
| 25 | type: body && typeof body === "object" ? (body as Record<string, unknown>).type : typeof body, |
| 26 | }); |
| 27 | return { action: "ack" }; |
| 28 | } |
| 29 | |
| 30 | function applyQueueResult(message: QueueMessage, result: QueueHandleResult): void { |
| 31 | if (result.action === "retry") { |
| 32 | message.retry({ delaySeconds: result.delaySeconds }); |
| 33 | return; |
| 34 | } |
| 35 | message.ack(); |
| 36 | } |
| 37 | |
| 38 | export default { |
| 39 | fetch(request, env, ctx) { |
| 40 | setLevel(env.LOG_LEVEL); |
| 41 | return app.fetch(request, env, ctx); |
| 42 | }, |
| 43 | async queue(batch, env) { |
| 44 | setLevel(env.LOG_LEVEL); |
| 45 | for (const message of batch.messages) { |
| 46 | try { |
| 47 | applyQueueResult(message, await handleQueueMessage(env, message.body)); |
| 48 | } catch (error) { |
| 49 | queueLog.error("message_failed", { |
| 50 | eventType: "message_failed", |
| 51 | ...errorContext(error), |
| 52 | }); |
| 53 | message.retry({ delaySeconds: 30 }); |
| 54 | } |
| 55 | } |
| 56 | }, |
| 57 | } satisfies ExportedHandler<Env>; |