fix: preserve active downloads and atomically prepare refreshed snapshots

This commit is contained in:
Thomas Durieux
2026-09-27 06:41:39 +00:00
parent 0002e8d1a1
commit 083544bf68
2 changed files with 94 additions and 8 deletions
+28 -1
View File
@@ -333,7 +333,27 @@ export default class Repository {
if (!claimed.matchedCount) throw new AnonymousError("invalid_status", { httpStatus: 409 }); if (!claimed.matchedCount) throw new AnonymousError("invalid_status", { httpStatus: 409 });
this.refreshToken = token; this.refreshToken = token;
try { try {
// An old timestamp can still belong to a live worker. Reusing its job ID
// would discard the replacement and cancel the old worker's generation.
const job = await downloadQueue.getJob(`repo-${this.repoId}`);
if (job) {
const state = await job.getState();
if (state !== "completed" && state !== "failed") {
throw new AnonymousError("invalid_status", { httpStatus: 409 });
}
await job.remove();
}
await this.updateIfNeeded({ force: true }); await this.updateIfNeeded({ force: true });
} catch (error) {
if (this.status === RepositoryStatus.PREPARING) {
// A failed reset/enqueue must remain retryable, even if the lease
// expired. Never change a replacement lease or a concurrent removal.
await AnonymizedRepositoryModel.updateOne({
_id: this.model._id, refreshToken: token,
status: RepositoryStatus.PREPARING, statusDate: this.model.statusDate,
}, { $set: { status: RepositoryStatus.ERROR, statusDate: new Date(), statusMessage: "preparation_interrupted" } }).exec();
}
throw error;
} finally { } finally {
this.refreshToken = undefined; this.refreshToken = undefined;
// Only this lease may be released, including after a failed GitHub lookup. // Only this lease may be released, including after a failed GitHub lookup.
@@ -441,6 +461,7 @@ export default class Repository {
commit: newCommit, commit: newCommit,
}); });
const statusDate = new Date();
if (isConnected) { if (isConnected) {
const result = await AnonymizedRepositoryModel.updateOne( const result = await AnonymizedRepositoryModel.updateOne(
{ _id: this._model._id, ...this.refreshFilter() }, { _id: this._model._id, ...this.refreshFilter() },
@@ -449,12 +470,18 @@ export default class Repository {
"source.commit": newCommit, "source.commit": newCommit,
"source.commitDate": this._model.source.commitDate, "source.commitDate": this._model.source.commitDate,
anonymizeDate: this._model.anonymizeDate, anonymizeDate: this._model.anonymizeDate,
status: RepositoryStatus.PREPARING,
statusDate,
statusMessage: null,
}, },
} }
).exec(); ).exec();
this.checkRefreshWrite(result); this.checkRefreshWrite(result);
} }
await this.resetSate(RepositoryStatus.PREPARING); this.model.status = RepositoryStatus.PREPARING;
this.model.statusDate = statusDate;
this.model.statusMessage = undefined;
await this.resetSate();
if (isConnected && this.refreshToken) { if (isConnected && this.refreshToken) {
// Removal or expiry may have started while deleting the old cache. // Removal or expiry may have started while deleting the old cache.
const current = await AnonymizedRepositoryModel.exists({ _id: this.model._id, ...this.refreshFilter() }); const current = await AnonymizedRepositoryModel.exists({ _id: this.model._id, ...this.refreshFilter() });
+66 -7
View File
@@ -60,9 +60,9 @@ describe("repository refresh and restoration", () => {
branches: async () => [{ name: "main", commit }], branches: async () => [{ name: "main", commit }],
getCommitInfo: async () => ({ commit: {} }), getCommitInfo: async () => ({ commit: {} }),
})); }));
repo.resetSate = async status => { repo.model.status = status; }; repo.resetSate = async () => {};
let added; let added;
stub(queue, "downloadQueue", { add: async (...args) => { added = args; } }); stub(queue, "downloadQueue", { getJob: async () => undefined, add: async (...args) => { added = args; } });
let response; let response;
await refresh({}, { json: body => { response = body; } }); await refresh({}, { json: body => { response = body; } });
expect(response.status).to.equal("preparing"); expect(response.status).to.equal("preparing");
@@ -76,7 +76,7 @@ describe("repository refresh and restoration", () => {
const repo = repository("removed"); const repo = repository("removed");
const failure = new Error("token_expired"); const failure = new Error("token_expired");
repo.getToken = async () => { throw failure; }; repo.getToken = async () => { throw failure; };
stub(queue, "downloadQueue", { add: async () => { throw new Error("must not enqueue"); } }); stub(queue, "downloadQueue", { getJob: async () => undefined, add: async () => { throw new Error("must not enqueue"); } });
let caught; let caught;
try { await refresh({}, {}); } catch (error) { caught = error; } try { await refresh({}, {}); } catch (error) { caught = error; }
expect(caught).to.equal(failure); expect(caught).to.equal(failure);
@@ -96,6 +96,7 @@ describe("repository refresh and restoration", () => {
function database(repo) { function database(repo) {
const stored = repo.model.toObject(); const stored = repo.model.toObject();
stub(db, "isConnected", true); stub(db, "isConnected", true);
stub(queue, "downloadQueue", { getJob: async () => undefined });
const matches = filter => require("sift").default(filter)(stored); const matches = filter => require("sift").default(filter)(stored);
stub(Model, "updateOne", (filter, update) => ({ exec: async () => { stub(Model, "updateOne", (filter, update) => ({ exec: async () => {
if (!matches(filter)) return { matchedCount: 0 }; if (!matches(filter)) return { matchedCount: 0 };
@@ -186,7 +187,7 @@ describe("repository refresh and restoration", () => {
getCommitInfo: async () => ({ commit: {} }), getCommitInfo: async () => ({ commit: {} }),
})); }));
repo.resetSate = async () => { throw new Error("must not delete cache"); }; repo.resetSate = async () => { throw new Error("must not delete cache"); };
stub(queue, "downloadQueue", { add: async () => { throw new Error("must not enqueue"); } }); stub(queue, "downloadQueue", { getJob: async () => undefined, add: async () => { throw new Error("must not enqueue"); } });
let failure; let failure;
try { await repo.refresh(); } catch (error) { failure = error; } try { await repo.refresh(); } catch (error) { failure = error; }
expect(failure?.message).to.equal("invalid_status"); expect(failure?.message).to.equal("invalid_status");
@@ -205,12 +206,11 @@ describe("repository refresh and restoration", () => {
branches: async () => [{ name: "main", commit: "saved-sha" }], branches: async () => [{ name: "main", commit: "saved-sha" }],
getCommitInfo: async () => ({ commit: {} }), getCommitInfo: async () => ({ commit: {} }),
})); }));
repo.resetSate = async status => { repo.resetSate = async () => {
await repo.updateStatus(status);
if (removedDuringReset) stored.status = "removing"; if (removedDuringReset) stored.status = "removing";
}; };
let added = false; let added = false;
stub(queue, "downloadQueue", { add: async () => { added = true; } }); stub(queue, "downloadQueue", { getJob: async () => undefined, add: async () => { added = true; } });
let failure; let failure;
try { await repo.refresh(); } catch (error) { failure = error; } try { await repo.refresh(); } catch (error) { failure = error; }
expect(added).to.equal(!removedDuringReset); expect(added).to.equal(!removedDuringReset);
@@ -220,6 +220,65 @@ describe("repository refresh and restoration", () => {
}); });
} }
for (const state of ["active", "waiting", "delayed", "prioritized", "waiting-children"]) {
it(`rejects a stale download with a ${state} job before touching its snapshot`, async () => {
const repo = repository("download");
repo.model.statusDate = new Date(Date.now() - 6 * 60_000);
const stored = database(repo);
stub(queue, "downloadQueue", { getJob: async () => ({ getState: async () => state }) });
repo.updateIfNeeded = async () => { throw new Error("must not refresh"); };
let failure;
try { await repo.refresh(); } catch (error) { failure = error; }
expect(failure?.message).to.equal("invalid_status");
expect(stored.status).to.equal("download");
expect(stored.source.commit).to.equal("saved-sha");
});
}
for (const state of ["completed", "failed"]) {
it(`removes a ${state} download job before reusing its ID`, async () => {
const repo = repository("download");
database(repo);
let removed = false;
stub(queue, "downloadQueue", { getJob: async () => ({
getState: async () => state, remove: async () => { removed = true; },
}) });
repo.updateIfNeeded = async () => { expect(removed).to.equal(true); };
await repo.refresh();
});
}
it("keeps a snapshot retryable if the lease expires after the commit update", async () => {
const repo = repository("ready");
const stored = database(repo);
repo.getToken = async () => "token";
stub(github, "getRepositoryFromGitHub", async () => ({
fullName: "owner/repo", model: {},
branches: async () => [{ name: "main", commit: "new-sha" }],
getCommitInfo: async () => ({ commit: {} }),
}));
repo.resetSate = async () => {
expect(stored.source.commit).to.equal("new-sha");
expect(stored.status).to.equal("preparing");
stored.refreshUntil = new Date(0);
};
let failure;
try { await repo.refresh(); } catch (error) { failure = error; }
expect(failure?.message).to.equal("invalid_status");
expect(stored.status).to.equal("error");
expect(stored).not.to.have.property("refreshToken");
const retry = new Repository(new Model(stored));
retry.getToken = async () => "token";
let cleared = false, queued = false;
retry.resetSate = async () => { cleared = true; };
stub(queue, "downloadQueue", { getJob: async () => undefined, add: async () => { queued = true; } });
await retry.refresh();
expect(cleared).to.equal(true);
expect(queued).to.equal(true);
expect(stored.status).to.equal("preparing");
});
it("serves the dashboard polling URL with repository status", async () => { it("serves the dashboard polling URL with repository status", async () => {
const repo = repository("ready"); const repo = repository("ready");
stub(db, "getRepository", async () => repo); stub(db, "getRepository", async () => repo);