File
Blob: src/worker/queues/repair-jobs.ts
| 1 | import { blobGcMessage } from "@/worker/queues/blob-gc"; |
| 2 | import { enqueueBlobGc } from "@/worker/r2/gc"; |
| 3 | |
| 4 | export interface RepairJobMessage { |
| 5 | type: "storage_repair"; |
| 6 | subject_id: string; |
| 7 | storage_id: string; |
| 8 | reason: string; |
| 9 | created_at_ms: number; |
| 10 | } |
| 11 | |
| 12 | export function repairJobMessage(input: { |
| 13 | subjectId: string; |
| 14 | storageId: string; |
| 15 | reason: string; |
| 16 | createdAtMs: number; |
| 17 | }): RepairJobMessage { |
| 18 | return { |
| 19 | type: "storage_repair", |
| 20 | subject_id: input.subjectId, |
| 21 | storage_id: input.storageId, |
| 22 | reason: input.reason, |
| 23 | created_at_ms: input.createdAtMs, |
| 24 | }; |
| 25 | } |
| 26 | |
| 27 | export function isRepairJobMessage(value: unknown): value is RepairJobMessage { |
| 28 | if (!value || typeof value !== "object") return false; |
| 29 | const record = value as Record<string, unknown>; |
| 30 | return ( |
| 31 | record.type === "storage_repair" && |
| 32 | typeof record.subject_id === "string" && |
| 33 | record.subject_id.length > 0 && |
| 34 | typeof record.storage_id === "string" && |
| 35 | record.storage_id.length > 0 && |
| 36 | typeof record.reason === "string" && |
| 37 | record.reason.length > 0 && |
| 38 | typeof record.created_at_ms === "number" && |
| 39 | Number.isFinite(record.created_at_ms) |
| 40 | ); |
| 41 | } |
| 42 | |
| 43 | export async function handleRepairJob(env: Env, message: RepairJobMessage, nowMs = Date.now()) { |
| 44 | const subject = { subjectId: message.subject_id, storageId: message.storage_id, nowMs }; |
| 45 | const fileRepair = await env.FILE_DAV.getByName(message.storage_id).repair({ ...subject, limit: 100 }); |
| 46 | await env.CAL_DAV.getByName(message.storage_id).ensureInitialized(subject); |
| 47 | await env.CARD_DAV.getByName(message.storage_id).ensureInitialized(subject); |
| 48 | |
| 49 | await enqueueBlobGc( |
| 50 | env, |
| 51 | fileRepair.blobGc.map((blob) => |
| 52 | blobGcMessage({ |
| 53 | subjectId: message.subject_id, |
| 54 | storageId: message.storage_id, |
| 55 | blobId: blob.blobId, |
| 56 | blobKey: blob.blobKey, |
| 57 | notBeforeMs: nowMs, |
| 58 | }), |
| 59 | ), |
| 60 | ); |
| 61 | |
| 62 | return { |
| 63 | action: "ack" as const, |
| 64 | cleanedPendingUploads: fileRepair.cleanedPendingUploads, |
| 65 | cleanedLocks: fileRepair.cleanedLocks, |
| 66 | queuedBlobGc: fileRepair.blobGc.length, |
| 67 | }; |
| 68 | } |