Skip to content
File

Blob: tests/worker/runtime/queues/maintenance.workers.test.ts

typescript192 lines
1import { env } from "cloudflare:workers";
2import { describe, expect, it, vi } from "vitest";
3 
4import worker from "@/worker/index";
5import { parseIfHeader, type ParsedIfHeader } from "@/worker/dav/if-header";
6import { normalizeFilesPath, type NormalizedDavPath } from "@/worker/dav/paths";
7import { createControlPlaneDb } from "@/worker/db/d1/client";
8import { getSubjectById } from "@/worker/db/d1/repository";
9import { blobGcMessage } from "@/worker/queues/blob-gc";
10import { handleRepairJob, repairJobMessage } from "@/worker/queues/repair-jobs";
11import { handleBlobGc } from "@/worker/r2/gc";
12import { createDavFixture } from "@tests/worker/helpers/dav";
13 
14function filePath(pathname: string): NormalizedDavPath {
15 const result = normalizeFilesPath(pathname);
16 if (!result.ok) throw new Error(`Invalid test path: ${pathname}`);
17 return result.path;
18}
19 
20function emptyIfHeader(): ParsedIfHeader {
21 const result = parseIfHeader(null);
22 if (!result.ok) throw new Error(result.message);
23 return result.header;
24}
25 
26async function storageIdFor(subjectId: string): Promise<string> {
27 const subject = await getSubjectById(createControlPlaneDb(env.DAV_CONTROL_PLANE), subjectId);
28 if (!subject) throw new Error(`Missing subject ${subjectId}`);
29 return subject.storageId;
30}
31 
32describe("Worker maintenance operations", () => {
33 it("aborts expired pending uploads and garbage-collects their R2 blobs idempotently", async () => {
34 const fixture = await createDavFixture();
35 const storageId = await storageIdFor(fixture.subjectId);
36 const object = env.FILE_DAV.getByName(storageId);
37 const nowMs = Date.now();
38 const begin = await object.beginWrite({
39 subjectId: fixture.subjectId,
40 storageId,
41 path: filePath(`/files/stale-${crypto.randomUUID()}.txt`),
42 contentLength: 5,
43 contentType: "text/plain",
44 contentLanguage: null,
45 ifMatch: null,
46 ifNoneMatch: null,
47 ifHeader: emptyIfHeader(),
48 lockTokens: [],
49 nowMs,
50 maxFileBytes: 100,
51 });
52 expect(begin.ok).toBe(true);
53 if (!begin.ok) throw new Error(begin.message);
54 
55 await env.FILE_BLOBS.put(begin.blobKey, new TextEncoder().encode("stale"));
56 expect(await env.FILE_BLOBS.get(begin.blobKey)).not.toBeNull();
57 
58 const beforeExpiry = await object.cleanupExpiredPendingUploads({
59 subjectId: fixture.subjectId,
60 storageId,
61 nowMs,
62 limit: 10,
63 });
64 expect(beforeExpiry.cleaned).toBe(0);
65 
66 const afterExpiry = await object.cleanupExpiredPendingUploads({
67 subjectId: fixture.subjectId,
68 storageId,
69 nowMs: nowMs + 15 * 60_000,
70 limit: 10,
71 });
72 expect(afterExpiry.cleaned).toBe(1);
73 expect(afterExpiry.blobGc).toEqual([{ blobId: begin.blobId, blobKey: begin.blobKey }]);
74 
75 const message = blobGcMessage({
76 subjectId: fixture.subjectId,
77 storageId,
78 blobId: begin.blobId,
79 blobKey: begin.blobKey,
80 notBeforeMs: nowMs,
81 });
82 await expect(handleBlobGc(env, message, nowMs + 15 * 60_000)).resolves.toEqual({ action: "ack" });
83 expect(await env.FILE_BLOBS.get(begin.blobKey)).toBeNull();
84 await expect(handleBlobGc(env, message, nowMs + 15 * 60_000)).resolves.toEqual({ action: "ack" });
85 });
86 
87 it("deletes expired locks and leaves active locks alone", async () => {
88 const fixture = await createDavFixture();
89 const storageId = await storageIdFor(fixture.subjectId);
90 const object = env.FILE_DAV.getByName(storageId);
91 const nowMs = Date.now();
92 const locked = await object.lock({
93 subjectId: fixture.subjectId,
94 storageId,
95 path: filePath(`/files/lock-${crypto.randomUUID()}.txt`),
96 ownerXml: '<D:owner xmlns:D="DAV:"/>',
97 scope: "exclusive",
98 depth: "0",
99 timeoutSeconds: 1,
100 ifHeader: { lists: [] },
101 lockTokens: [],
102 nowMs,
103 refresh: false,
104 });
105 expect(locked.ok).toBe(true);
106 
107 await expect(
108 object.cleanupExpiredLocks({ subjectId: fixture.subjectId, storageId, nowMs, limit: 10 }),
109 ).resolves.toMatchObject({ cleaned: 0 });
110 await expect(
111 object.cleanupExpiredLocks({ subjectId: fixture.subjectId, storageId, nowMs: nowMs + 1_000, limit: 10 }),
112 ).resolves.toMatchObject({ cleaned: 1 });
113 await expect(
114 object.cleanupExpiredLocks({ subjectId: fixture.subjectId, storageId, nowMs: nowMs + 1_001, limit: 10 }),
115 ).resolves.toMatchObject({ cleaned: 0 });
116 });
117 
118 it("handles storage repair jobs as idempotent no-ops when nothing is stale", async () => {
119 const fixture = await createDavFixture();
120 const storageId = await storageIdFor(fixture.subjectId);
121 const sendBatch = vi.fn();
122 const repairEnv = {
123 FILE_DAV: env.FILE_DAV,
124 CAL_DAV: env.CAL_DAV,
125 CARD_DAV: env.CARD_DAV,
126 BLOB_GC: { sendBatch },
127 } as unknown as Env;
128 const message = repairJobMessage({
129 subjectId: fixture.subjectId,
130 storageId,
131 reason: "test-noop",
132 createdAtMs: Date.now(),
133 });
134 
135 await expect(handleRepairJob(repairEnv, message)).resolves.toMatchObject({
136 action: "ack",
137 cleanedPendingUploads: 0,
138 cleanedLocks: 0,
139 queuedBlobGc: 0,
140 });
141 await expect(handleRepairJob(repairEnv, message)).resolves.toMatchObject({
142 action: "ack",
143 cleanedPendingUploads: 0,
144 cleanedLocks: 0,
145 queuedBlobGc: 0,
146 });
147 expect(sendBatch).not.toHaveBeenCalled();
148 });
149 
150 it("retries queue messages when a handler throws and delays future blob GC messages", async () => {
151 const nowMs = Date.now();
152 const futureMessage = blobGcMessage({
153 subjectId: "sub",
154 storageId: "stg",
155 blobId: "blob",
156 blobKey: "files/stg/blob",
157 notBeforeMs: nowMs + 60_000,
158 });
159 await expect(handleBlobGc(env, futureMessage, nowMs)).resolves.toEqual({ action: "retry", delaySeconds: 60 });
160 
161 if (!worker.queue) throw new Error("Worker queue handler is not exported");
162 const retry = vi.fn();
163 const ack = vi.fn();
164 const failingMessage = {
165 ...futureMessage,
166 not_before_ms: nowMs - 1,
167 };
168 await worker.queue(
169 {
170 messages: [
171 {
172 body: failingMessage,
173 ack,
174 retry,
175 },
176 ],
177 } as unknown as MessageBatch<unknown>,
178 {
179 LOG_LEVEL: "error",
180 FILE_DAV: {
181 getByName() {
182 throw new Error("storage unavailable");
183 },
184 },
185 } as unknown as Env,
186 );
187 
188 expect(ack).not.toHaveBeenCalled();
189 expect(retry).toHaveBeenCalledWith({ delaySeconds: 30 });
190 });
191});