From 12b0bb52c017b85f9657a430822d0dcf6b359241 Mon Sep 17 00:00:00 2001 From: tdurieux Date: Sun, 6 Sep 2026 09:16:13 +0200 Subject: [PATCH] fix: ignore stale removal jobs after repository restoration --- src/queue/index.ts | 28 ++++++++-- test/backend-reliability.test.js | 87 ++++++++++++++++++++++++++++++++ 2 files changed, 112 insertions(+), 3 deletions(-) diff --git a/src/queue/index.ts b/src/queue/index.ts index 6fffdd5..13c1f02 100644 --- a/src/queue/index.ts +++ b/src/queue/index.ts @@ -144,18 +144,40 @@ export async function addRemovalJob( /** * Requeue removals whose database status survived but whose BullMQ job did * not, for example after a crash between updating MongoDB and queueing Redis. - * Repository removal is idempotent, so replaying a terminal job is safe. + * Failed jobs are replayed only while their removal request is still current. */ export async function recoverStuckRemoving() { if (!removeQueue) return; try { - // A failed job is durable proof that removal was requested. Retry it even - // if the final failure handler already changed the repository to ERROR. + // Claim the pending removal before queueing it. A restored repository, or + // an error from a later operation, must not revive an old deletion request. const failedJobs = await removeQueue.getJobs(["failed"]); for (const job of failedJobs) { const repoId = job.data?.repoId; if (!repoId) continue; try { + const pending = await AnonymizedRepositoryModel.findOneAndUpdate( + { + repoId, + $or: [ + { status: RepositoryStatus.REMOVING }, + ...(job.timestamp + ? [{ + status: RepositoryStatus.ERROR, + $or: [ + { anonymizeDate: { $lte: new Date(job.timestamp) } }, + { anonymizeDate: { $exists: false } }, + ], + }] + : []), + ], + }, + { $set: { status: RepositoryStatus.REMOVING } } + ).collation({ locale: "en", strength: 2 }).exec(); + if (!pending) { + await job.remove(); + continue; + } await addRemovalJob(repoId); logger.info("requeued failed removal", { repoId }); } catch (e) { diff --git a/test/backend-reliability.test.js b/test/backend-reliability.test.js index bd4acc2..33c0997 100644 --- a/test/backend-reliability.test.js +++ b/test/backend-reliability.test.js @@ -24,6 +24,93 @@ const { } = require("../src/server/schedule"); const Repository = require("../src/core/Repository").default; const AnonymizedRepositoryModel = require("../src/core/model/anonymizedRepositories/anonymizedRepositories.model").default; +const queueModule = require("../src/queue"); +const routeUtils = require("../src/server/routes/route-utils"); + +const UserModel = require("../src/core/model/users/users.model").default; + +describe("removal recovery", function () { + let originals; + beforeEach(function () { + originals = { + distinct: UserModel.distinct, + queue: queueModule.removeQueue, + find: AnonymizedRepositoryModel.find, + findOneAndUpdate: AnonymizedRepositoryModel.findOneAndUpdate, + getRepo: routeUtils.getRepo, + getUser: routeUtils.getUser, + handleError: routeUtils.handleError, + }; + UserModel.distinct = () => ({ exec: async () => [] }); + }); + afterEach(function () { + UserModel.distinct = originals.distinct; + queueModule.removeQueue = originals.queue; + AnonymizedRepositoryModel.find = originals.find; + AnonymizedRepositoryModel.findOneAndUpdate = originals.findOneAndUpdate; + routeUtils.getRepo = originals.getRepo; + routeUtils.getUser = originals.getUser; + routeUtils.handleError = originals.handleError; + }); + + for (const status of ["ready", "preparing", "removed", "error"]) { + it(`discards an old failed removal after restoration to ${status}`, async function () { + let discarded = false; + let queued = false; + const job = { + data: { repoId: "repo-1" }, timestamp: 1000, + remove: async () => { discarded = true; }, + }; + queueModule.removeQueue = { + getJobs: async () => [job], + getJob: async () => undefined, + add: async () => { queued = true; }, + }; + AnonymizedRepositoryModel.findOneAndUpdate = (filter) => ({ + collation() { return this; }, + exec: async () => { + // A restored snapshot was created after the failed deletion request. + const eligible = require("sift").default(filter)({ + repoId: "repo-1", status, anonymizeDate: new Date(2000), + }); + return eligible ? { repoId: "repo-1" } : null; + }, + }); + AnonymizedRepositoryModel.find = () => ({ lean: async () => [] }); + await queueModule.recoverStuckRemoving(); + expect(discarded).to.equal(true); + expect(queued).to.equal(false); + }); + } + + for (const status of ["error", "removing"]) { + it(`retries a failed removal still in ${status}`, async function () { + const repo = { + repoId: "repo-1", status, anonymizeDate: new Date(500), + statusDate: new Date(3000), + }; + let added; + const job = { data: { repoId: repo.repoId }, timestamp: 1000, finishedOn: 2000 }; + queueModule.removeQueue = { + getJobs: async () => [job], + getJob: async () => undefined, + add: async (...args) => { added = args; }, + }; + AnonymizedRepositoryModel.findOneAndUpdate = (filter, update) => ({ + collation() { return this; }, + exec: async () => { + if (!require("sift").default(filter)(repo)) return null; + Object.assign(repo, update.$set); + return repo; + }, + }); + AnonymizedRepositoryModel.find = () => ({ lean: async () => [] }); + await queueModule.recoverStuckRemoving(); + expect(repo.status).to.equal("removing"); + expect(added[1]).to.deep.equal({ repoId: repo.repoId }); + }); + } +}); describe("conference edits", function () { const form = {