mirror of
https://github.com/tdurieux/anonymous_github.git
synced 2026-09-29 05:31:43 +02:00
Merge pull request #842 from tdurieux/fix/repository-refresh-concurrency
fix: serialize repository refreshes and recover stale downloads
This commit is contained in:
+86
-11
@@ -1,4 +1,5 @@
|
|||||||
import storage from "./storage";
|
import storage from "./storage";
|
||||||
|
import { randomUUID } from "crypto";
|
||||||
import { RepositoryStatus } from "./types";
|
import { RepositoryStatus } from "./types";
|
||||||
import { Readable } from "stream";
|
import { Readable } from "stream";
|
||||||
import * as sha1 from "crypto-js/sha1";
|
import * as sha1 from "crypto-js/sha1";
|
||||||
@@ -298,11 +299,70 @@ export default class Repository {
|
|||||||
return true;
|
return true;
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
private refreshToken?: string;
|
||||||
* Update the repository if a new commit exists
|
|
||||||
*
|
private refreshFilter() {
|
||||||
* @returns void
|
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 {
|
||||||
|
// 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.
|
||||||
|
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<void> {
|
async updateIfNeeded(opt?: { force: boolean }): Promise<void> {
|
||||||
this.assertNotArchived();
|
this.assertNotArchived();
|
||||||
if (
|
if (
|
||||||
@@ -311,7 +371,6 @@ export default class Repository {
|
|||||||
this._model.options.expirationDate
|
this._model.options.expirationDate
|
||||||
) {
|
) {
|
||||||
if (this._model.options.expirationDate <= new Date()) {
|
if (this._model.options.expirationDate <= new Date()) {
|
||||||
this._model.status = RepositoryStatus.EXPIRED;
|
|
||||||
await this.expire();
|
await this.expire();
|
||||||
throw new AnonymousError("repository_expired", {
|
throw new AnonymousError("repository_expired", {
|
||||||
object: this,
|
object: this,
|
||||||
@@ -345,10 +404,11 @@ export default class Repository {
|
|||||||
if (this.model.source.repositoryName !== ghRepo.fullName) {
|
if (this.model.source.repositoryName !== ghRepo.fullName) {
|
||||||
this.model.source.repositoryName = ghRepo.fullName;
|
this.model.source.repositoryName = ghRepo.fullName;
|
||||||
if (isConnected) {
|
if (isConnected) {
|
||||||
await AnonymizedRepositoryModel.updateOne(
|
const result = await AnonymizedRepositoryModel.updateOne(
|
||||||
{ _id: this._model._id },
|
{ _id: this._model._id, ...this.refreshFilter() },
|
||||||
{ $set: { "source.repositoryName": ghRepo.fullName } }
|
{ $set: { "source.repositoryName": ghRepo.fullName } }
|
||||||
).exec();
|
).exec();
|
||||||
|
this.checkRefreshWrite(result);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
const branches = await ghRepo.branches({
|
const branches = await ghRepo.branches({
|
||||||
@@ -401,19 +461,32 @@ export default class Repository {
|
|||||||
commit: newCommit,
|
commit: newCommit,
|
||||||
});
|
});
|
||||||
|
|
||||||
|
const statusDate = new Date();
|
||||||
if (isConnected) {
|
if (isConnected) {
|
||||||
await AnonymizedRepositoryModel.updateOne(
|
const result = await AnonymizedRepositoryModel.updateOne(
|
||||||
{ _id: this._model._id },
|
{ _id: this._model._id, ...this.refreshFilter() },
|
||||||
{
|
{
|
||||||
$set: {
|
$set: {
|
||||||
"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.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() });
|
||||||
|
if (!current) throw new AnonymousError("invalid_status", { httpStatus: 409 });
|
||||||
}
|
}
|
||||||
await this.resetSate(RepositoryStatus.PREPARING);
|
|
||||||
await downloadQueue.add(this.repoId, { repoId: this.repoId }, {
|
await downloadQueue.add(this.repoId, { repoId: this.repoId }, {
|
||||||
jobId: `repo-${this.repoId}`,
|
jobId: `repo-${this.repoId}`,
|
||||||
attempts: 3,
|
attempts: 3,
|
||||||
@@ -481,6 +554,7 @@ export default class Repository {
|
|||||||
const result = await AnonymizedRepositoryModel.updateOne(
|
const result = await AnonymizedRepositoryModel.updateOne(
|
||||||
{
|
{
|
||||||
_id: this._model._id,
|
_id: this._model._id,
|
||||||
|
...this.refreshFilter(),
|
||||||
...(this.protectLifecycle ? {
|
...(this.protectLifecycle ? {
|
||||||
status: { $nin: [RepositoryStatus.ARCHIVED, RepositoryStatus.REMOVING, RepositoryStatus.REMOVED,
|
status: { $nin: [RepositoryStatus.ARCHIVED, RepositoryStatus.REMOVING, RepositoryStatus.REMOVED,
|
||||||
RepositoryStatus.EXPIRING, RepositoryStatus.EXPIRED] },
|
RepositoryStatus.EXPIRING, RepositoryStatus.EXPIRED] },
|
||||||
@@ -490,6 +564,7 @@ export default class Repository {
|
|||||||
},
|
},
|
||||||
{ $set: { status, statusDate, statusMessage, ...(publishedAt ? { publishedAt } : {}) } }
|
{ $set: { status, statusDate, statusMessage, ...(publishedAt ? { publishedAt } : {}) } }
|
||||||
).exec();
|
).exec();
|
||||||
|
this.checkRefreshWrite(result);
|
||||||
if (this.protectLifecycle && result.matchedCount === 0) {
|
if (this.protectLifecycle && result.matchedCount === 0) {
|
||||||
throw new AnonymousError("repository_job_cancelled", { httpStatus: 410 });
|
throw new AnonymousError("repository_job_cancelled", { httpStatus: 410 });
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -11,6 +11,8 @@ const AnonymizedRepositorySchema = new Schema({
|
|||||||
default: "preparing",
|
default: "preparing",
|
||||||
},
|
},
|
||||||
statusDate: Date,
|
statusDate: Date,
|
||||||
|
refreshToken: { type: String, select: false },
|
||||||
|
refreshUntil: { type: Date, select: false },
|
||||||
archivedAt: Date,
|
archivedAt: Date,
|
||||||
archiveReason: String,
|
archiveReason: String,
|
||||||
archiveCachePending: Boolean,
|
archiveCachePending: Boolean,
|
||||||
|
|||||||
@@ -132,14 +132,15 @@ router.post(
|
|||||||
if (
|
if (
|
||||||
repo.status == RepositoryStatus.PREPARING ||
|
repo.status == RepositoryStatus.PREPARING ||
|
||||||
repo.status == RepositoryStatus.QUEUE ||
|
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.REMOVING ||
|
||||||
repo.status == RepositoryStatus.EXPIRING
|
repo.status == RepositoryStatus.EXPIRING
|
||||||
) {
|
) {
|
||||||
throw new AnonymousError("invalid_status", { httpStatus: 409 });
|
throw new AnonymousError("invalid_status", { httpStatus: 409 });
|
||||||
}
|
}
|
||||||
|
|
||||||
await repo.updateIfNeeded({ force: true });
|
await repo.refresh();
|
||||||
res.json({ status: repo.status });
|
res.json({ status: repo.status });
|
||||||
} catch (error) {
|
} catch (error) {
|
||||||
handleError(error, res, req);
|
handleError(error, res, req);
|
||||||
|
|||||||
@@ -19,12 +19,13 @@ describe("repository refresh and restoration", () => {
|
|||||||
afterEach(() => { while (restores.length) restores.pop()(); });
|
afterEach(() => { while (restores.length) restores.pop()(); });
|
||||||
function repository(status) {
|
function repository(status) {
|
||||||
const repo = new Repository(new Model({ repoId: "restore-me", status,
|
const repo = new Repository(new Model({ repoId: "restore-me", status,
|
||||||
|
statusDate: new Date(),
|
||||||
owner: "507f1f77bcf86cd799439011",
|
owner: "507f1f77bcf86cd799439011",
|
||||||
source: { repositoryName: "owner/repo", branch: "main", commit: "saved-sha" },
|
source: { repositoryName: "owner/repo", branch: "main", commit: "saved-sha" },
|
||||||
options: { expirationMode: "never" },
|
options: { expirationMode: "never" },
|
||||||
}));
|
}));
|
||||||
stub(utils, "getRepo", async () => repo);
|
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; });
|
stub(utils, "handleError", error => { throw error; });
|
||||||
return repo;
|
return repo;
|
||||||
}
|
}
|
||||||
@@ -59,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");
|
||||||
@@ -75,11 +76,237 @@ 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);
|
||||||
expect(repo.status).to.equal("removed");
|
expect(repo.status).to.equal("removed");
|
||||||
expect(repo.model.source.commit).to.equal("saved-sha");
|
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);
|
||||||
|
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 };
|
||||||
|
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", { 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");
|
||||||
|
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 () => {
|
||||||
|
if (removedDuringReset) stored.status = "removing";
|
||||||
|
};
|
||||||
|
let added = false;
|
||||||
|
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);
|
||||||
|
expect(stored.status).to.equal(removedDuringReset ? "removing" : "preparing");
|
||||||
|
expect(failure?.message).to.equal(removedDuringReset ? "invalid_status" : undefined);
|
||||||
|
expect(stored).not.to.have.property("refreshToken");
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
|
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);
|
||||||
|
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)); }
|
||||||
|
});
|
||||||
|
|
||||||
});
|
});
|
||||||
|
|||||||
Reference in New Issue
Block a user