mirror of
https://github.com/tdurieux/anonymous_github.git
synced 2026-09-13 06:08:58 +02:00
fix: persist and recover repository removals for banned owners
This commit is contained in:
@@ -1,3 +1,4 @@
|
|||||||
|
import UserModel from "../core/model/users/users.model";
|
||||||
import { JobsOptions, Queue, Worker } from "bullmq";
|
import { JobsOptions, Queue, Worker } from "bullmq";
|
||||||
import config from "../config";
|
import config from "../config";
|
||||||
import AnonymizedRepositoryModel from "../core/model/anonymizedRepositories/anonymizedRepositories.model";
|
import AnonymizedRepositoryModel from "../core/model/anonymizedRepositories/anonymizedRepositories.model";
|
||||||
@@ -152,6 +153,11 @@ export async function recoverStuckRemoving() {
|
|||||||
// Claim the pending removal before queueing it. A restored repository, or
|
// Claim the pending removal before queueing it. A restored repository, or
|
||||||
// an error from a later operation, must not revive an old deletion request.
|
// an error from a later operation, must not revive an old deletion request.
|
||||||
const failedJobs = await removeQueue.getJobs(["failed"]);
|
const failedJobs = await removeQueue.getJobs(["failed"]);
|
||||||
|
// Older ban jobs could fail before persisting REMOVING. A banned owner
|
||||||
|
// still has a current removal request even if that old job left READY.
|
||||||
|
const bannedOwners = failedJobs.length
|
||||||
|
? await UserModel.distinct("_id", { status: "banned" }).exec()
|
||||||
|
: [];
|
||||||
for (const job of failedJobs) {
|
for (const job of failedJobs) {
|
||||||
const repoId = job.data?.repoId;
|
const repoId = job.data?.repoId;
|
||||||
if (!repoId) continue;
|
if (!repoId) continue;
|
||||||
@@ -161,6 +167,10 @@ export async function recoverStuckRemoving() {
|
|||||||
repoId,
|
repoId,
|
||||||
$or: [
|
$or: [
|
||||||
{ status: RepositoryStatus.REMOVING },
|
{ status: RepositoryStatus.REMOVING },
|
||||||
|
...(bannedOwners.length ? [{
|
||||||
|
owner: { $in: bannedOwners },
|
||||||
|
status: { $ne: RepositoryStatus.REMOVED },
|
||||||
|
}] : []),
|
||||||
...(job.timestamp
|
...(job.timestamp
|
||||||
? [{
|
? [{
|
||||||
status: RepositoryStatus.ERROR,
|
status: RepositoryStatus.ERROR,
|
||||||
|
|||||||
@@ -6,7 +6,7 @@ import AnonymousError from "../../core/AnonymousError";
|
|||||||
import AnonymizedRepositoryModel from "../../core/model/anonymizedRepositories/anonymizedRepositories.model";
|
import AnonymizedRepositoryModel from "../../core/model/anonymizedRepositories/anonymizedRepositories.model";
|
||||||
import ConferenceModel from "../../core/model/conference/conferences.model";
|
import ConferenceModel from "../../core/model/conference/conferences.model";
|
||||||
import UserModel from "../../core/model/users/users.model";
|
import UserModel from "../../core/model/users/users.model";
|
||||||
import { cacheQueue, downloadQueue, removeQueue } from "../../queue";
|
import { addRemovalJob, cacheQueue, downloadQueue, removeQueue } from "../../queue";
|
||||||
import { queryMetrics } from "../../queue/queueMetrics";
|
import { queryMetrics } from "../../queue/queueMetrics";
|
||||||
import {
|
import {
|
||||||
computeStats,
|
computeStats,
|
||||||
@@ -1186,7 +1186,11 @@ router.post(
|
|||||||
let queued = 0;
|
let queued = 0;
|
||||||
for (const repo of repos) {
|
for (const repo of repos) {
|
||||||
try {
|
try {
|
||||||
await removeQueue.add(repo.repoId, { repoId: repo.repoId }, { jobId: `repo-${repo.repoId}` });
|
await AnonymizedRepositoryModel.updateOne(
|
||||||
|
{ repoId: repo.repoId },
|
||||||
|
{ $set: { status: "removing", statusDate: new Date() } }
|
||||||
|
).exec();
|
||||||
|
await addRemovalJob(repo.repoId);
|
||||||
queued++;
|
queued++;
|
||||||
} catch {
|
} catch {
|
||||||
// job may already exist in the queue
|
// job may already exist in the queue
|
||||||
|
|||||||
@@ -53,6 +53,24 @@ describe("removal recovery", function () {
|
|||||||
routeUtils.handleError = originals.handleError;
|
routeUtils.handleError = originals.handleError;
|
||||||
});
|
});
|
||||||
|
|
||||||
|
it("retries an older ban job that failed before marking the repository removing", async function () {
|
||||||
|
let queued = false;
|
||||||
|
UserModel.distinct = () => ({ exec: async () => ["banned-owner"] });
|
||||||
|
queueModule.removeQueue = {
|
||||||
|
getJobs: async () => [{ data: { repoId: "repo" }, timestamp: 1000,
|
||||||
|
remove: async () => { throw new Error("must not discard ban"); } }],
|
||||||
|
getJob: async () => undefined,
|
||||||
|
add: async () => { queued = true; },
|
||||||
|
};
|
||||||
|
AnonymizedRepositoryModel.findOneAndUpdate = filter => ({
|
||||||
|
collation() { return this; },
|
||||||
|
exec: async () => require("sift").default(filter)({ repoId: "repo", owner: "banned-owner", status: "ready" }) ? { repoId: "repo" } : null,
|
||||||
|
});
|
||||||
|
AnonymizedRepositoryModel.find = () => ({ lean: async () => [] });
|
||||||
|
await queueModule.recoverStuckRemoving();
|
||||||
|
expect(queued).to.equal(true);
|
||||||
|
});
|
||||||
|
|
||||||
for (const status of ["ready", "preparing", "removed", "error"]) {
|
for (const status of ["ready", "preparing", "removed", "error"]) {
|
||||||
it(`discards an old failed removal after restoration to ${status}`, async function () {
|
it(`discards an old failed removal after restoration to ${status}`, async function () {
|
||||||
let discarded = false;
|
let discarded = false;
|
||||||
|
|||||||
Reference in New Issue
Block a user