Fix repository cleanup after removal and expiration (#782)

This commit is contained in:
Thomas Durieux
2026-08-20 12:52:19 +02:00
committed by GitHub
parent 46e956a779
commit 35c3b1804b
8 changed files with 370 additions and 69 deletions
+5 -5
View File
@@ -331,11 +331,11 @@ export default class AnonymizedFile {
});
}
const content = await this.repository.source?.getFileContent(this);
if (
!this.repository.model.isReseted ||
this.repository.status != RepositoryStatus.READY
) {
this.repository.model.isReseted = false;
const cacheWasReset = this.repository.model.isReseted;
if (cacheWasReset) {
await this.repository.markCachePresent();
}
if (cacheWasReset || this.repository.status != RepositoryStatus.READY) {
await this.repository.updateStatus(RepositoryStatus.READY);
}
return content;
+14
View File
@@ -543,6 +543,20 @@ export default class Repository {
}
}
/**
* Record that repository content has been cached again after a reset.
*/
async markCachePresent() {
if (!this.model.isReseted) return;
this.model.isReseted = false;
if (isConnected) {
await AnonymizedRepositoryModel.updateOne(
{ _id: this._model._id },
{ $set: { isReseted: false } }
).exec();
}
}
/**
* Compute the size of the repository in term of storage and number of files.
*
+101 -6
View File
@@ -1,4 +1,4 @@
import { Queue, Worker } from "bullmq";
import { JobsOptions, Queue, Worker } from "bullmq";
import config from "../config";
import AnonymizedRepositoryModel from "../core/model/anonymizedRepositories/anonymizedRepositories.model";
import { RepositoryStatus } from "../core/types";
@@ -22,6 +22,22 @@ const IN_FLIGHT_STATUSES: RepositoryStatus[] = [
RepositoryStatus.DOWNLOAD,
];
const LIVE_JOB_STATES = new Set([
"active",
"waiting",
"delayed",
"prioritized",
"waiting-children",
]);
const REMOVAL_JOB_OPTIONS: JobsOptions = {
attempts: 3,
backoff: { type: "exponential", delay: 1000 },
removeOnComplete: true,
// Keep a bounded failure history so operators can inspect and retry jobs.
removeOnFail: { count: 1000 },
};
async function markErrorIfInFlight(repoId: string, message: string) {
try {
await AnonymizedRepositoryModel.updateOne(
@@ -44,6 +60,26 @@ async function markErrorIfInFlight(repoId: string, message: string) {
}
}
async function markErrorIfRemoving(repoId: string, message: string) {
try {
await AnonymizedRepositoryModel.updateOne(
{ repoId, status: RepositoryStatus.REMOVING },
{
$set: {
status: RepositoryStatus.ERROR,
statusDate: new Date(),
statusMessage: message || "removal_failed",
},
}
).exec();
} catch (e) {
logger.error("markErrorIfRemoving failed", {
...serializeError(e),
repoId,
});
}
}
/**
* Recover repositories left in an in-flight status (preparing/queue/download)
* with no live BullMQ job — typically caused by a worker process crash or
@@ -84,6 +120,58 @@ export let cacheQueue: Queue<RepoJobData>;
export let removeQueue: Queue<RepoJobData>;
export let downloadQueue: Queue<RepoJobData>;
type RemovalQueue = Pick<Queue<RepoJobData>, "add" | "getJob">;
/**
* Add an idempotent repository-removal job. A live job wins; a terminal job
* with the same stable id is replaced so a retry is not silently discarded.
*/
export async function addRemovalJob(
repoId: string,
queue: RemovalQueue = removeQueue
): Promise<boolean> {
const jobId = `repo-${repoId}`;
const existing = await queue.getJob(jobId);
if (existing) {
const state = await existing.getState();
if (LIVE_JOB_STATES.has(state)) return false;
await existing.remove().catch(() => undefined);
}
await queue.add(repoId, { repoId }, { ...REMOVAL_JOB_OPTIONS, jobId });
return true;
}
/**
* 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.
*/
export async function recoverStuckRemoving() {
if (!removeQueue) return;
try {
const stuck = await AnonymizedRepositoryModel.find(
{ status: RepositoryStatus.REMOVING },
{ repoId: 1 }
).lean();
for (const doc of stuck) {
try {
const queued = await addRemovalJob(doc.repoId);
if (queued) {
logger.info("requeued interrupted removal", { repoId: doc.repoId });
}
} catch (e) {
logger.warn("removal recovery failed", {
...serializeError(e),
repoId: doc.repoId,
});
await markErrorIfRemoving(doc.repoId, "removal_interrupted");
}
}
} catch (e) {
logger.error("recoverStuckRemoving failed", serializeError(e));
}
}
// avoid to load the queue outside the main server
export function startWorker() {
const connection = {
@@ -103,10 +191,7 @@ export function startWorker() {
host: config.REDIS_HOSTNAME,
port: config.REDIS_PORT,
},
defaultJobOptions: {
removeOnComplete: true,
removeOnFail: true,
},
defaultJobOptions: REMOVAL_JOB_OPTIONS,
});
downloadQueue = new Queue<RepoJobData>("repository download", {
connection,
@@ -144,8 +229,18 @@ export function startWorker() {
recordMetric("remove", "completed", (job.finishedOn || Date.now()) - (job.processedOn || job.timestamp));
await job.remove();
});
removeWorker.on("failed", async (job) => {
removeWorker.on("failed", async (job, err) => {
if (job) recordMetric("remove", "failed", Date.now() - (job.processedOn || job.timestamp));
const repoId = job?.data?.repoId;
logger.error("removal failed", {
...serializeError(err),
repoId,
});
if (!repoId) return;
if (job && typeof job.attemptsMade === "number" && job.opts?.attempts) {
if (job.attemptsMade < job.opts.attempts) return;
}
await markErrorIfRemoving(repoId, err?.message || "removal_failed");
});
const downloadWorker = new Worker<RepoJobData>(
+19 -11
View File
@@ -17,23 +17,31 @@ export async function processRemoveRepository(
) {
const { connect, getRepository }: Database =
database || require("../../server/database");
let repo: Awaited<ReturnType<typeof getRepositoryImport>> | undefined;
try {
await connect();
logger.info("removing repository", { repoId: job.data.repoId });
const repo = await getRepository(job.data.repoId);
repo = await getRepository(job.data.repoId);
await repo.updateStatus(RepositoryStatus.REMOVING, "");
try {
await repo.remove();
} catch (error) {
if (error instanceof Error) {
await repo.updateStatus(RepositoryStatus.ERROR, error.message);
} else if (typeof error === "string") {
await repo.updateStatus(RepositoryStatus.ERROR, error);
}
throw error;
}
await repo.remove();
logger.info("repository removed", { repoId: job.data.repoId });
} catch (error) {
if (repo) {
const message =
error instanceof Error
? error.message
: typeof error === "string"
? error
: "removal_failed";
try {
await repo.updateStatus(RepositoryStatus.ERROR, message);
} catch (statusError) {
logger.error("failed to record repository removal error", {
...serializeError(statusError),
repoId: job.data.repoId,
});
}
}
logger.error("repository removal failed", {
...serializeError(error),
repoId: job.data.repoId,
+12 -1
View File
@@ -17,9 +17,14 @@ import router from "./routes";
import {
conferenceStatusCheck,
repositoryStatusCheck,
runRepositoryStatusCheck,
dailyStatsSnapshot,
} from "./schedule";
import { startWorker, recoverStuckPreparing } from "../queue";
import {
startWorker,
recoverStuckPreparing,
recoverStuckRemoving,
} from "../queue";
import {
computeStats,
ensureTodaySnapshot,
@@ -363,6 +368,12 @@ export default async function start() {
recoverStuckPreparing().catch((err) =>
logger.error("recoverStuckPreparing failed", serializeError(err))
);
recoverStuckRemoving().catch((err) =>
logger.error("recoverStuckRemoving failed", serializeError(err))
);
runRepositoryStatusCheck().catch((err) =>
logger.error("initial repository status check failed", serializeError(err))
);
}
start();
+9 -2
View File
@@ -17,7 +17,7 @@ import { IAnonymizedRepositoryDocument } from "../../core/model/anonymizedReposi
import UserModel from "../../core/model/users/users.model";
import ConferenceModel from "../../core/model/conference/conferences.model";
import AnonymousError from "../../core/AnonymousError";
import { downloadQueue, removeQueue } from "../../queue";
import { addRemovalJob, downloadQueue } from "../../queue";
import RepositoryModel from "../../core/model/repositories/repositories.model";
import User from "../../core/User";
import { RepositoryStatus } from "../../core/types";
@@ -241,7 +241,14 @@ router.delete(
const user = await getUser(req);
isOwnerOrAdmin([repo.owner.id], user);
await repo.updateStatus(RepositoryStatus.REMOVING);
await removeQueue.add(repo.repoId, { repoId: repo.repoId }, { jobId: `repo-${repo.repoId}` });
try {
await addRemovalJob(repo.repoId);
} catch (error) {
const message =
error instanceof Error ? error.message : "removal_enqueue_failed";
await repo.updateStatus(RepositoryStatus.ERROR, message);
throw error;
}
return res.json({ status: repo.status });
} catch (error) {
handleError(error, res, req);
+109 -44
View File
@@ -30,55 +30,120 @@ export function conferenceStatusCheck() {
export function repositoryStatusCheck() {
// check every 6 hours the status of the repositories
schedule.scheduleJob("0 */6 * * *", async () => {
logger.info("checking repository status and unused repositories");
const now = new Date();
const fourMonthAgo = new Date(now);
fourMonthAgo.setMonth(fourMonthAgo.getMonth() - 4);
const cursor = AnonymizedRepositoryModel.find({
status: RepositoryStatus.READY,
isReseted: false,
$or: [
{
"options.expirationMode": { $in: ["redirect", "remove"] },
"options.expirationDate": { $lte: now },
},
{ lastView: { $lt: fourMonthAgo } },
],
}).cursor();
const batch: Promise<void>[] = [];
for await (const data of cursor) {
batch.push(
(async () => {
const repo = new Repository(data);
try {
await repo.check();
} catch {
logger.info("repository expired", { repoId: repo.repoId });
}
await runRepositoryStatusCheck();
});
}
if (repo.model.lastView < fourMonthAgo) {
try {
await repo.removeCache();
} catch (error) {
logger.error("repository cache removal failed", {
...serializeError(error),
repoId: repo.repoId,
});
return;
}
logger.info("removed cache for unused repository", {
export function repositoryMaintenanceQuery(now: Date) {
const fourMonthAgo = new Date(now);
fourMonthAgo.setMonth(fourMonthAgo.getMonth() - 4);
return {
status: RepositoryStatus.READY,
$or: [
{
"options.expirationMode": { $in: ["redirect", "remove"] },
"options.expirationDate": { $lte: now },
},
{
isReseted: { $ne: true },
lastView: { $lt: fourMonthAgo },
},
],
};
}
export async function runRepositoryStatusCheck(now = new Date()) {
logger.info("checking repository status and unused repositories");
const fourMonthAgo = new Date(now);
fourMonthAgo.setMonth(fourMonthAgo.getMonth() - 4);
const batch: Promise<void>[] = [];
const flushBatch = async () => {
await Promise.all(batch);
batch.length = 0;
};
const cursor = AnonymizedRepositoryModel.find(
repositoryMaintenanceQuery(now)
).cursor();
for await (const data of cursor) {
batch.push(
(async () => {
const repo = new Repository(data);
const shouldExpire =
repo.options.expirationMode !== "never" &&
repo.options.expirationDate != null &&
repo.options.expirationDate <= now;
if (shouldExpire) {
try {
await repo.expire();
logger.info("repository expired", { repoId: repo.repoId });
} catch (error) {
logger.error("repository expiration failed", {
...serializeError(error),
repoId: repo.repoId,
});
}
})()
);
if (batch.length >= 10) {
await Promise.all(batch);
batch.length = 0;
}
return;
}
if (
repo.model.isReseted !== true &&
repo.model.lastView < fourMonthAgo
) {
try {
await repo.removeCache();
logger.info("removed cache for unused repository", {
repoId: repo.repoId,
});
} catch (error) {
logger.error("repository cache removal failed", {
...serializeError(error),
repoId: repo.repoId,
});
}
}
})()
);
if (batch.length >= 10) {
await flushBatch();
}
await Promise.all(batch);
});
}
await flushBatch();
// Repair terminal records left with data by an older or interrupted
// expiration. This makes the cleanup idempotent across deployments.
const dirtyTerminalCursor = AnonymizedRepositoryModel.find({
$or: [
{ status: RepositoryStatus.EXPIRING },
{
status: RepositoryStatus.EXPIRED,
isReseted: { $ne: true },
},
],
}).cursor();
for await (const data of dirtyTerminalCursor) {
batch.push(
(async () => {
const repo = new Repository(data);
try {
await repo.resetSate();
await repo.updateStatus(RepositoryStatus.EXPIRED);
logger.info("recovered expired repository cleanup", {
repoId: repo.repoId,
});
} catch (error) {
logger.error("expired repository cleanup failed", {
...serializeError(error),
repoId: repo.repoId,
});
}
})()
);
if (batch.length >= 10) {
await flushBatch();
}
}
await flushBatch();
}
export function dailyStatsSnapshot() {