export { AuthObject } from "@/worker/objects/auth-object"; export { FileDavObject } from "@/worker/objects/file-dav-object"; export { CalDavObject } from "@/worker/objects/cal-dav-object"; export { CardDavObject } from "@/worker/objects/card-dav-object"; import { createApp } from "@/worker/app"; import { isBlobGcMessage } from "@/worker/queues/blob-gc"; import { handleRepairJob, isRepairJobMessage } from "@/worker/queues/repair-jobs"; import type { QueueHandleResult } from "@/worker/queues/types"; import { handleBlobGc } from "@/worker/r2/gc"; import { createLogger, errorContext, setLevel } from "@/worker/util/logging"; const app = createApp(); const queueLog = createLogger("queue"); type QueueMessage = | Message | { body: unknown; ack: () => void; retry: (options?: { delaySeconds?: number }) => void }; export async function handleQueueMessage(env: Env, body: unknown, nowMs = Date.now()): Promise { if (isBlobGcMessage(body)) return await handleBlobGc(env, body, nowMs); if (isRepairJobMessage(body)) return await handleRepairJob(env, body, nowMs); queueLog.warn("unknown_message_type", { eventType: "unknown_message_type", type: body && typeof body === "object" ? (body as Record).type : typeof body, }); return { action: "ack" }; } function applyQueueResult(message: QueueMessage, result: QueueHandleResult): void { if (result.action === "retry") { message.retry({ delaySeconds: result.delaySeconds }); return; } message.ack(); } export default { fetch(request, env, ctx) { setLevel(env.LOG_LEVEL); return app.fetch(request, env, ctx); }, async queue(batch, env) { setLevel(env.LOG_LEVEL); for (const message of batch.messages) { try { applyQueueResult(message, await handleQueueMessage(env, message.body)); } catch (error) { queueLog.error("message_failed", { eventType: "message_failed", ...errorContext(error), }); message.retry({ delaySeconds: 30 }); } } }, } satisfies ExportedHandler;