Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
97 changes: 86 additions & 11 deletions src/core/Repository.ts
Original file line number Diff line number Diff line change
@@ -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";
Expand Down Expand Up @@ -298,11 +299,70 @@ 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 {
// 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> {
this.assertNotArchived();
if (
Expand All @@ -311,7 +371,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,
Expand Down Expand Up @@ -345,10 +404,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({
Expand Down Expand Up @@ -401,19 +461,32 @@ export default class Repository {
commit: newCommit,
});

const statusDate = new Date();
if (isConnected) {
await AnonymizedRepositoryModel.updateOne(
{ _id: this._model._id },
const result = await AnonymizedRepositoryModel.updateOne(
{ _id: this._model._id, ...this.refreshFilter() },
{
$set: {
"source.commit": newCommit,
"source.commitDate": this._model.source.commitDate,
anonymizeDate: this._model.anonymizeDate,
status: RepositoryStatus.PREPARING,
statusDate,
statusMessage: null,
},
}
).exec();
this.checkRefreshWrite(result);
Comment thread
tdurieux marked this conversation as resolved.
}
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 }, {
jobId: `repo-${this.repoId}`,
attempts: 3,
Expand Down Expand Up @@ -481,6 +554,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] },
Expand All @@ -490,6 +564,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 });
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
5 changes: 3 additions & 2 deletions src/server/routes/repository-private.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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)) ||
Comment thread
tdurieux marked this conversation as resolved.
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);
Expand Down
Loading
Loading