Skip to content
File

Blob: src/worker/queues/repair-jobs.ts

typescript69 lines
1import { blobGcMessage } from "@/worker/queues/blob-gc";
2import { enqueueBlobGc } from "@/worker/r2/gc";
3 
4export interface RepairJobMessage {
5 type: "storage_repair";
6 subject_id: string;
7 storage_id: string;
8 reason: string;
9 created_at_ms: number;
10}
11 
12export 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 
27export 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 
43export 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}