From 0002e8d1a117b68304c100a066704e101f6b9dce Mon Sep 17 00:00:00 2001 From: Thomas Durieux <5577568+tdurieux@users.noreply.github.com> Date: Sun, 27 Sep 2026 06:34:05 +0000 Subject: [PATCH 1/2] fix: serialize repository refreshes and recover stale downloads --- src/core/Repository.ts | 68 +++++-- .../anonymizedRepositories.schema.ts | 2 + src/server/routes/repository-private.ts | 5 +- test/repository-refresh.test.js | 170 +++++++++++++++++- 4 files changed, 232 insertions(+), 13 deletions(-) diff --git a/src/core/Repository.ts b/src/core/Repository.ts index fa56c02..57fb113 100644 --- a/src/core/Repository.ts +++ b/src/core/Repository.ts @@ -1,4 +1,5 @@ import storage from "./storage"; +import { randomUUID } from "crypto"; import { RepositoryStatus } from "./types"; import { Readable } from "stream"; import * as sha1 from "crypto-js/sha1"; @@ -298,11 +299,50 @@ export default class Repository { return true; } - /** - * Update the repository if a new commit exists - * - * @returns void - */ + private refreshToken?: string; + + private refreshFilter() { + return this.refreshToken ? { + refreshToken: this.refreshToken, + refreshUntil: { $gt: new Date() }, + status: this.model.status, + statusDate: this.model.statusDate || { $exists: false }, + "githubAccess.revision": this.model.githubAccess?.revision || { $exists: false }, + } : {}; + } + + private checkRefreshWrite(result: { matchedCount: number }) { + if (this.refreshToken && !result.matchedCount) { + throw new AnonymousError("invalid_status", { httpStatus: 409 }); + } + } + + /** Serialize dashboard refreshes across server processes without hiding a ready snapshot. */ + async refresh() { + this.assertNotArchived(); + if (!isConnected) return this.updateIfNeeded({ force: true }); + const token = randomUUID(); + const now = new Date(); + const claimed = await AnonymizedRepositoryModel.updateOne({ + _id: this.model._id, + status: this.model.status, + statusDate: this.model.statusDate || { $exists: false }, + "githubAccess.revision": this.model.githubAccess?.revision || { $exists: false }, + $or: [{ refreshUntil: { $exists: false } }, { refreshUntil: { $lte: now } }], + }, { $set: { refreshToken: token, refreshUntil: new Date(now.getTime() + 5 * 60_000) } }).exec(); + if (!claimed.matchedCount) throw new AnonymousError("invalid_status", { httpStatus: 409 }); + this.refreshToken = token; + try { + await this.updateIfNeeded({ force: true }); + } finally { + this.refreshToken = undefined; + // Only this lease may be released, including after a failed GitHub lookup. + await AnonymizedRepositoryModel.updateOne({ _id: this.model._id, refreshToken: token }, + { $unset: { refreshToken: 1, refreshUntil: 1 } }).exec(); + } + } + + /** Update the repository if a new commit exists. */ async updateIfNeeded(opt?: { force: boolean }): Promise { this.assertNotArchived(); if ( @@ -311,7 +351,6 @@ export default class Repository { this._model.options.expirationDate ) { if (this._model.options.expirationDate <= new Date()) { - this._model.status = RepositoryStatus.EXPIRED; await this.expire(); throw new AnonymousError("repository_expired", { object: this, @@ -345,10 +384,11 @@ export default class Repository { if (this.model.source.repositoryName !== ghRepo.fullName) { this.model.source.repositoryName = ghRepo.fullName; if (isConnected) { - await AnonymizedRepositoryModel.updateOne( - { _id: this._model._id }, + const result = await AnonymizedRepositoryModel.updateOne( + { _id: this._model._id, ...this.refreshFilter() }, { $set: { "source.repositoryName": ghRepo.fullName } } ).exec(); + this.checkRefreshWrite(result); } } const branches = await ghRepo.branches({ @@ -402,8 +442,8 @@ export default class Repository { }); if (isConnected) { - await AnonymizedRepositoryModel.updateOne( - { _id: this._model._id }, + const result = await AnonymizedRepositoryModel.updateOne( + { _id: this._model._id, ...this.refreshFilter() }, { $set: { "source.commit": newCommit, @@ -412,8 +452,14 @@ export default class Repository { }, } ).exec(); + this.checkRefreshWrite(result); } await this.resetSate(RepositoryStatus.PREPARING); + if (isConnected && this.refreshToken) { + // Removal or expiry may have started while deleting the old cache. + const current = await AnonymizedRepositoryModel.exists({ _id: this.model._id, ...this.refreshFilter() }); + if (!current) throw new AnonymousError("invalid_status", { httpStatus: 409 }); + } await downloadQueue.add(this.repoId, { repoId: this.repoId }, { jobId: `repo-${this.repoId}`, attempts: 3, @@ -481,6 +527,7 @@ export default class Repository { const result = await AnonymizedRepositoryModel.updateOne( { _id: this._model._id, + ...this.refreshFilter(), ...(this.protectLifecycle ? { status: { $nin: [RepositoryStatus.ARCHIVED, RepositoryStatus.REMOVING, RepositoryStatus.REMOVED, RepositoryStatus.EXPIRING, RepositoryStatus.EXPIRED] }, @@ -490,6 +537,7 @@ export default class Repository { }, { $set: { status, statusDate, statusMessage, ...(publishedAt ? { publishedAt } : {}) } } ).exec(); + this.checkRefreshWrite(result); if (this.protectLifecycle && result.matchedCount === 0) { throw new AnonymousError("repository_job_cancelled", { httpStatus: 410 }); } diff --git a/src/core/model/anonymizedRepositories/anonymizedRepositories.schema.ts b/src/core/model/anonymizedRepositories/anonymizedRepositories.schema.ts index b42f49d..e7ac0c3 100644 --- a/src/core/model/anonymizedRepositories/anonymizedRepositories.schema.ts +++ b/src/core/model/anonymizedRepositories/anonymizedRepositories.schema.ts @@ -11,6 +11,8 @@ const AnonymizedRepositorySchema = new Schema({ default: "preparing", }, statusDate: Date, + refreshToken: { type: String, select: false }, + refreshUntil: { type: Date, select: false }, archivedAt: Date, archiveReason: String, archiveCachePending: Boolean, diff --git a/src/server/routes/repository-private.ts b/src/server/routes/repository-private.ts index b248a45..caa8e58 100644 --- a/src/server/routes/repository-private.ts +++ b/src/server/routes/repository-private.ts @@ -132,14 +132,15 @@ router.post( if ( repo.status == RepositoryStatus.PREPARING || repo.status == RepositoryStatus.QUEUE || - repo.status == RepositoryStatus.DOWNLOAD || + (repo.status == RepositoryStatus.DOWNLOAD && + repo.model.statusDate > new Date(Date.now() - 5 * 60_000)) || repo.status == RepositoryStatus.REMOVING || repo.status == RepositoryStatus.EXPIRING ) { throw new AnonymousError("invalid_status", { httpStatus: 409 }); } - await repo.updateIfNeeded({ force: true }); + await repo.refresh(); res.json({ status: repo.status }); } catch (error) { handleError(error, res, req); diff --git a/test/repository-refresh.test.js b/test/repository-refresh.test.js index c6d1900..04c303e 100644 --- a/test/repository-refresh.test.js +++ b/test/repository-refresh.test.js @@ -19,12 +19,13 @@ describe("repository refresh and restoration", () => { afterEach(() => { while (restores.length) restores.pop()(); }); function repository(status) { const repo = new Repository(new Model({ repoId: "restore-me", status, + statusDate: new Date(), owner: "507f1f77bcf86cd799439011", source: { repositoryName: "owner/repo", branch: "main", commit: "saved-sha" }, options: { expirationMode: "never" }, })); stub(utils, "getRepo", async () => repo); - stub(utils, "getUser", async () => ({ isAdmin: true })); + stub(utils, "getUser", async () => ({ isAdmin: true, model: { id: "admin" } })); stub(utils, "handleError", error => { throw error; }); return repo; } @@ -82,4 +83,171 @@ describe("repository refresh and restoration", () => { expect(repo.status).to.equal("removed"); expect(repo.model.source.commit).to.equal("saved-sha"); }); + + it("allows a stale download to be retried", async () => { + const repo = repository("download"); + repo.model.statusDate = new Date(Date.now() - 6 * 60_000); + let retried = false; + repo.refresh = async () => { retried = true; }; + await refresh({}, { json: () => {} }); + expect(retried).to.equal(true); + }); + + function database(repo) { + const stored = repo.model.toObject(); + stub(db, "isConnected", true); + const matches = filter => require("sift").default(filter)(stored); + stub(Model, "updateOne", (filter, update) => ({ exec: async () => { + if (!matches(filter)) return { matchedCount: 0 }; + for (const [key, value] of Object.entries(update.$set || {})) { + const parts = key.split("."); + let target = stored; + while (parts.length > 1) { const part = parts.shift(); target = target[part] ||= {}; } + target[parts[0]] = value; + } + for (const key of Object.keys(update.$unset || {})) delete stored[key]; + return { matchedCount: 1 }; + } })); + stub(Model, "exists", async filter => matches(filter) ? { _id: stored._id } : null); + return stored; + } + + it("claims a lease before GitHub work and rejects a concurrent refresh", async () => { + const first = repository("ready"); + const second = new Repository(new Model(first.model.toObject())); + const stored = database(first); + let release, started; + const entered = new Promise(resolve => { started = resolve; }); + first.updateIfNeeded = async () => { + started(); + await new Promise(resolve => { release = resolve; }); + }; + second.updateIfNeeded = async () => { throw new Error("must not start another refresh"); }; + const running = first.refresh(); + await entered; + try { + let failure; + try { await second.refresh(); } catch (error) { failure = error; } + expect(failure?.message).to.equal("invalid_status"); + expect(stored.status).to.equal("ready"); + } finally { release(); await running; } + expect(stored).not.to.have.property("refreshToken"); + }); + + it("releases a failed lease so reconnecting can be retried", async () => { + const repo = repository("removed"); + const stored = database(repo); + const failure = new Error("token_expired"); + repo.getToken = async () => { throw failure; }; + let caught; + try { await repo.refresh(); } catch (error) { caught = error; } + expect(caught).to.equal(failure); + expect(stored).not.to.have.property("refreshToken"); + expect(stored.status).to.equal("removed"); + repo.updateIfNeeded = async () => {}; + await repo.refresh(); + }); + + it("reclaims an expired lease without letting the old request write or release it", async () => { + const old = repository("ready"); + const stored = database(old); + const replacement = new Repository(new Model(old.model.toObject())); + let release, started; + const entered = new Promise(resolve => { started = resolve; }); + old.updateIfNeeded = async () => { + started(); + await new Promise(resolve => { release = resolve; }); + await old.updateStatus("preparing"); + }; + const running = old.refresh().catch(error => error); + await entered; + stored.refreshUntil = new Date(0); + replacement.updateIfNeeded = async () => { + const replacementToken = stored.refreshToken; + release(); + expect((await running).message).to.equal("invalid_status"); + expect(stored.refreshToken).to.equal(replacementToken); + }; + await replacement.refresh(); + expect(stored.status).to.equal("ready"); + }); + + for (const status of ["removing", "expiring", "archived"]) { + it(`does not overwrite ${status} requested during a GitHub lookup`, async () => { + const repo = repository("ready"); + const stored = database(repo); + repo.getToken = async () => "token"; + stub(github, "getRepositoryFromGitHub", async () => ({ + fullName: "owner/repo", model: {}, + branches: async () => { + stored.status = status; + return [{ name: "main", commit: "new-sha" }]; + }, + getCommitInfo: async () => ({ commit: {} }), + })); + repo.resetSate = async () => { throw new Error("must not delete cache"); }; + stub(queue, "downloadQueue", { add: async () => { throw new Error("must not enqueue"); } }); + let failure; + try { await repo.refresh(); } catch (error) { failure = error; } + expect(failure?.message).to.equal("invalid_status"); + expect(stored.status).to.equal(status); + expect(stored.source.commit).to.equal("saved-sha"); + }); + } + + for (const removedDuringReset of [false, true]) { + it(`rebuilds under a lease, removal during reset: ${removedDuringReset}`, async () => { + const repo = repository("removed"); + const stored = database(repo); + repo.getToken = async () => "token"; + stub(github, "getRepositoryFromGitHub", async () => ({ + fullName: "owner/repo", model: {}, + branches: async () => [{ name: "main", commit: "saved-sha" }], + getCommitInfo: async () => ({ commit: {} }), + })); + repo.resetSate = async status => { + await repo.updateStatus(status); + if (removedDuringReset) stored.status = "removing"; + }; + let added = false; + stub(queue, "downloadQueue", { add: async () => { added = true; } }); + let failure; + try { await repo.refresh(); } catch (error) { failure = error; } + expect(added).to.equal(!removedDuringReset); + expect(stored.status).to.equal(removedDuringReset ? "removing" : "preparing"); + expect(failure?.message).to.equal(removedDuringReset ? "invalid_status" : undefined); + expect(stored).not.to.have.property("refreshToken"); + }); + } + + it("serves the dashboard polling URL with repository status", async () => { + const repo = repository("ready"); + stub(db, "getRepository", async () => repo); + stub(require("../src/core/GitHubUtils"), "getToken", async () => { throw new Error("token_expired"); }); + const express = require("express"); + const app = express(); + app.use((req, _res, next) => { req.isAuthenticated = () => true; next(); }); + app.use("/api/repo", require("../src/server/routes/repository-public").default); + app.use("/api/repo", require("../src/server/routes/file").default); + app.use("/api/repo", router); + const server = await new Promise(resolve => { + const listener = app.listen(0, "127.0.0.1", () => resolve(listener)); + }); + try { + const response = await new Promise((resolve, reject) => { + require("http").get({ host: "127.0.0.1", port: server.address().port, path: "/api/repo/restore-me" }, res => { + let body = ""; + res.on("data", chunk => { body += chunk; }); + res.on("end", () => { + try { resolve({ status: res.statusCode, data: JSON.parse(body) }); } + catch (error) { reject(error); } + }); + }).on("error", reject); + }); + expect(response.status).to.equal(200); + expect(response.data.status).to.equal("ready"); + expect(response.data.connectionError).to.equal("token_expired"); + } finally { await new Promise(resolve => server.close(resolve)); } + }); + }); From 083544bf68737263c9d54182fd80be249fb0c9e4 Mon Sep 17 00:00:00 2001 From: Thomas Durieux <5577568+tdurieux@users.noreply.github.com> Date: Sun, 27 Sep 2026 06:41:39 +0000 Subject: [PATCH 2/2] fix: preserve active downloads and atomically prepare refreshed snapshots --- src/core/Repository.ts | 29 ++++++++++++- test/repository-refresh.test.js | 73 +++++++++++++++++++++++++++++---- 2 files changed, 94 insertions(+), 8 deletions(-) diff --git a/src/core/Repository.ts b/src/core/Repository.ts index 57fb113..77b6760 100644 --- a/src/core/Repository.ts +++ b/src/core/Repository.ts @@ -333,7 +333,27 @@ export default class Repository { if (!claimed.matchedCount) throw new AnonymousError("invalid_status", { httpStatus: 409 }); this.refreshToken = token; 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 }); + } 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 { this.refreshToken = undefined; // Only this lease may be released, including after a failed GitHub lookup. @@ -441,6 +461,7 @@ export default class Repository { commit: newCommit, }); + const statusDate = new Date(); if (isConnected) { const result = await AnonymizedRepositoryModel.updateOne( { _id: this._model._id, ...this.refreshFilter() }, @@ -449,12 +470,18 @@ export default class Repository { "source.commit": newCommit, "source.commitDate": this._model.source.commitDate, anonymizeDate: this._model.anonymizeDate, + status: RepositoryStatus.PREPARING, + statusDate, + statusMessage: null, }, } ).exec(); 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) { // Removal or expiry may have started while deleting the old cache. const current = await AnonymizedRepositoryModel.exists({ _id: this.model._id, ...this.refreshFilter() }); diff --git a/test/repository-refresh.test.js b/test/repository-refresh.test.js index 04c303e..7014877 100644 --- a/test/repository-refresh.test.js +++ b/test/repository-refresh.test.js @@ -60,9 +60,9 @@ describe("repository refresh and restoration", () => { branches: async () => [{ name: "main", commit }], getCommitInfo: async () => ({ commit: {} }), })); - repo.resetSate = async status => { repo.model.status = status; }; + repo.resetSate = async () => {}; let added; - stub(queue, "downloadQueue", { add: async (...args) => { added = args; } }); + stub(queue, "downloadQueue", { getJob: async () => undefined, add: async (...args) => { added = args; } }); let response; await refresh({}, { json: body => { response = body; } }); expect(response.status).to.equal("preparing"); @@ -76,7 +76,7 @@ describe("repository refresh and restoration", () => { const repo = repository("removed"); const failure = new Error("token_expired"); 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; try { await refresh({}, {}); } catch (error) { caught = error; } expect(caught).to.equal(failure); @@ -96,6 +96,7 @@ describe("repository refresh and restoration", () => { function database(repo) { const stored = repo.model.toObject(); stub(db, "isConnected", true); + stub(queue, "downloadQueue", { getJob: async () => undefined }); const matches = filter => require("sift").default(filter)(stored); stub(Model, "updateOne", (filter, update) => ({ exec: async () => { if (!matches(filter)) return { matchedCount: 0 }; @@ -186,7 +187,7 @@ describe("repository refresh and restoration", () => { getCommitInfo: async () => ({ commit: {} }), })); 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; try { await repo.refresh(); } catch (error) { failure = error; } expect(failure?.message).to.equal("invalid_status"); @@ -205,12 +206,11 @@ describe("repository refresh and restoration", () => { branches: async () => [{ name: "main", commit: "saved-sha" }], getCommitInfo: async () => ({ commit: {} }), })); - repo.resetSate = async status => { - await repo.updateStatus(status); + repo.resetSate = async () => { if (removedDuringReset) stored.status = "removing"; }; let added = false; - stub(queue, "downloadQueue", { add: async () => { added = true; } }); + stub(queue, "downloadQueue", { getJob: async () => undefined, add: async () => { added = true; } }); let failure; try { await repo.refresh(); } catch (error) { failure = error; } 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 () => { const repo = repository("ready"); stub(db, "getRepository", async () => repo);