import { env } from "cloudflare:workers"; import { describe, expect, it, vi } from "vitest"; import worker from "@/worker/index"; import { parseIfHeader, type ParsedIfHeader } from "@/worker/dav/if-header"; import { normalizeFilesPath, type NormalizedDavPath } from "@/worker/dav/paths"; import { createControlPlaneDb } from "@/worker/db/d1/client"; import { getSubjectById } from "@/worker/db/d1/repository"; import { blobGcMessage } from "@/worker/queues/blob-gc"; import { handleRepairJob, repairJobMessage } from "@/worker/queues/repair-jobs"; import { handleBlobGc } from "@/worker/r2/gc"; import { createDavFixture } from "@tests/worker/helpers/dav"; function filePath(pathname: string): NormalizedDavPath { const result = normalizeFilesPath(pathname); if (!result.ok) throw new Error(`Invalid test path: ${pathname}`); return result.path; } function emptyIfHeader(): ParsedIfHeader { const result = parseIfHeader(null); if (!result.ok) throw new Error(result.message); return result.header; } async function storageIdFor(subjectId: string): Promise { const subject = await getSubjectById(createControlPlaneDb(env.DAV_CONTROL_PLANE), subjectId); if (!subject) throw new Error(`Missing subject ${subjectId}`); return subject.storageId; } describe("Worker maintenance operations", () => { it("aborts expired pending uploads and garbage-collects their R2 blobs idempotently", async () => { const fixture = await createDavFixture(); const storageId = await storageIdFor(fixture.subjectId); const object = env.FILE_DAV.getByName(storageId); const nowMs = Date.now(); const begin = await object.beginWrite({ subjectId: fixture.subjectId, storageId, path: filePath(`/files/stale-${crypto.randomUUID()}.txt`), contentLength: 5, contentType: "text/plain", contentLanguage: null, ifMatch: null, ifNoneMatch: null, ifHeader: emptyIfHeader(), lockTokens: [], nowMs, maxFileBytes: 100, }); expect(begin.ok).toBe(true); if (!begin.ok) throw new Error(begin.message); await env.FILE_BLOBS.put(begin.blobKey, new TextEncoder().encode("stale")); expect(await env.FILE_BLOBS.get(begin.blobKey)).not.toBeNull(); const beforeExpiry = await object.cleanupExpiredPendingUploads({ subjectId: fixture.subjectId, storageId, nowMs, limit: 10, }); expect(beforeExpiry.cleaned).toBe(0); const afterExpiry = await object.cleanupExpiredPendingUploads({ subjectId: fixture.subjectId, storageId, nowMs: nowMs + 15 * 60_000, limit: 10, }); expect(afterExpiry.cleaned).toBe(1); expect(afterExpiry.blobGc).toEqual([{ blobId: begin.blobId, blobKey: begin.blobKey }]); const message = blobGcMessage({ subjectId: fixture.subjectId, storageId, blobId: begin.blobId, blobKey: begin.blobKey, notBeforeMs: nowMs, }); await expect(handleBlobGc(env, message, nowMs + 15 * 60_000)).resolves.toEqual({ action: "ack" }); expect(await env.FILE_BLOBS.get(begin.blobKey)).toBeNull(); await expect(handleBlobGc(env, message, nowMs + 15 * 60_000)).resolves.toEqual({ action: "ack" }); }); it("deletes expired locks and leaves active locks alone", async () => { const fixture = await createDavFixture(); const storageId = await storageIdFor(fixture.subjectId); const object = env.FILE_DAV.getByName(storageId); const nowMs = Date.now(); const locked = await object.lock({ subjectId: fixture.subjectId, storageId, path: filePath(`/files/lock-${crypto.randomUUID()}.txt`), ownerXml: '', scope: "exclusive", depth: "0", timeoutSeconds: 1, ifHeader: { lists: [] }, lockTokens: [], nowMs, refresh: false, }); expect(locked.ok).toBe(true); await expect( object.cleanupExpiredLocks({ subjectId: fixture.subjectId, storageId, nowMs, limit: 10 }), ).resolves.toMatchObject({ cleaned: 0 }); await expect( object.cleanupExpiredLocks({ subjectId: fixture.subjectId, storageId, nowMs: nowMs + 1_000, limit: 10 }), ).resolves.toMatchObject({ cleaned: 1 }); await expect( object.cleanupExpiredLocks({ subjectId: fixture.subjectId, storageId, nowMs: nowMs + 1_001, limit: 10 }), ).resolves.toMatchObject({ cleaned: 0 }); }); it("handles storage repair jobs as idempotent no-ops when nothing is stale", async () => { const fixture = await createDavFixture(); const storageId = await storageIdFor(fixture.subjectId); const sendBatch = vi.fn(); const repairEnv = { FILE_DAV: env.FILE_DAV, CAL_DAV: env.CAL_DAV, CARD_DAV: env.CARD_DAV, BLOB_GC: { sendBatch }, } as unknown as Env; const message = repairJobMessage({ subjectId: fixture.subjectId, storageId, reason: "test-noop", createdAtMs: Date.now(), }); await expect(handleRepairJob(repairEnv, message)).resolves.toMatchObject({ action: "ack", cleanedPendingUploads: 0, cleanedLocks: 0, queuedBlobGc: 0, }); await expect(handleRepairJob(repairEnv, message)).resolves.toMatchObject({ action: "ack", cleanedPendingUploads: 0, cleanedLocks: 0, queuedBlobGc: 0, }); expect(sendBatch).not.toHaveBeenCalled(); }); it("retries queue messages when a handler throws and delays future blob GC messages", async () => { const nowMs = Date.now(); const futureMessage = blobGcMessage({ subjectId: "sub", storageId: "stg", blobId: "blob", blobKey: "files/stg/blob", notBeforeMs: nowMs + 60_000, }); await expect(handleBlobGc(env, futureMessage, nowMs)).resolves.toEqual({ action: "retry", delaySeconds: 60 }); if (!worker.queue) throw new Error("Worker queue handler is not exported"); const retry = vi.fn(); const ack = vi.fn(); const failingMessage = { ...futureMessage, not_before_ms: nowMs - 1, }; await worker.queue( { messages: [ { body: failingMessage, ack, retry, }, ], } as unknown as MessageBatch, { LOG_LEVEL: "error", FILE_DAV: { getByName() { throw new Error("storage unavailable"); }, }, } as unknown as Env, ); expect(ack).not.toHaveBeenCalled(); expect(retry).toHaveBeenCalledWith({ delaySeconds: 30 }); }); });