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)); } + }); + });