mirror of
https://github.com/tdurieux/anonymous_github.git
synced 2026-07-23 05:20:53 +02:00
fix: preserve state across backend updates (#756)
This commit is contained in:
@@ -1,26 +1,34 @@
|
||||
import { SandboxedJob } from "bullmq";
|
||||
import { getRepository as getRepositoryImport } from "../../server/database";
|
||||
import { RepoJobData } from "../index";
|
||||
import { createLogger } from "../../core/logger";
|
||||
import { createLogger, serializeError } from "../../core/logger";
|
||||
|
||||
const logger = createLogger("queue:cache");
|
||||
|
||||
export default async function (job: SandboxedJob<RepoJobData, void>) {
|
||||
const {
|
||||
connect,
|
||||
getRepository,
|
||||
}: {
|
||||
connect: () => Promise<void>;
|
||||
getRepository: typeof getRepositoryImport;
|
||||
} = require("../../server/database");
|
||||
interface Database {
|
||||
connect: () => Promise<void>;
|
||||
getRepository: typeof getRepositoryImport;
|
||||
}
|
||||
|
||||
export async function processRemoveCache(
|
||||
job: SandboxedJob<RepoJobData, void>,
|
||||
database?: Database
|
||||
) {
|
||||
const { connect, getRepository }: Database =
|
||||
database || require("../../server/database");
|
||||
try {
|
||||
await connect();
|
||||
logger.info("removing cache", { repoId: job.data.repoId });
|
||||
const repo = await getRepository(job.data.repoId);
|
||||
await repo.removeCache();
|
||||
} catch {
|
||||
// error already handled
|
||||
} finally {
|
||||
logger.info("cache removed", { repoId: job.data.repoId });
|
||||
} catch (error) {
|
||||
logger.error("cache removal failed", {
|
||||
...serializeError(error),
|
||||
repoId: job.data.repoId,
|
||||
});
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
|
||||
export default processRemoveCache;
|
||||
|
||||
@@ -2,18 +2,21 @@ import { SandboxedJob } from "bullmq";
|
||||
import { getRepository as getRepositoryImport } from "../../server/database";
|
||||
import { RepositoryStatus } from "../../core/types";
|
||||
import { RepoJobData } from "../index";
|
||||
import { createLogger } from "../../core/logger";
|
||||
import { createLogger, serializeError } from "../../core/logger";
|
||||
|
||||
const logger = createLogger("queue:remove");
|
||||
|
||||
export default async function (job: SandboxedJob<RepoJobData, void>) {
|
||||
const {
|
||||
connect,
|
||||
getRepository,
|
||||
}: {
|
||||
connect: () => Promise<void>;
|
||||
getRepository: typeof getRepositoryImport;
|
||||
} = require("../../server/database");
|
||||
interface Database {
|
||||
connect: () => Promise<void>;
|
||||
getRepository: typeof getRepositoryImport;
|
||||
}
|
||||
|
||||
export async function processRemoveRepository(
|
||||
job: SandboxedJob<RepoJobData, void>,
|
||||
database?: Database
|
||||
) {
|
||||
const { connect, getRepository }: Database =
|
||||
database || require("../../server/database");
|
||||
try {
|
||||
await connect();
|
||||
logger.info("removing repository", { repoId: job.data.repoId });
|
||||
@@ -29,9 +32,14 @@ export default async function (job: SandboxedJob<RepoJobData, void>) {
|
||||
}
|
||||
throw error;
|
||||
}
|
||||
} catch {
|
||||
// error already handled
|
||||
} finally {
|
||||
logger.info("repository removed", { repoId: job.data.repoId });
|
||||
} catch (error) {
|
||||
logger.error("repository removal failed", {
|
||||
...serializeError(error),
|
||||
repoId: job.data.repoId,
|
||||
});
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
|
||||
export default processRemoveRepository;
|
||||
|
||||
@@ -145,6 +145,21 @@ function validateConferenceForm(conf: any) {
|
||||
}
|
||||
}
|
||||
|
||||
// eslint-disable-next-line @typescript-eslint/no-explicit-any
|
||||
export function applyConferenceForm(
|
||||
model: IConferenceDocument,
|
||||
body: any,
|
||||
isNew: boolean
|
||||
) {
|
||||
model.name = body.name;
|
||||
model.startDate = new Date(body.startDate);
|
||||
model.endDate = new Date(body.endDate);
|
||||
model.status = "ready";
|
||||
model.url = body.url;
|
||||
model.options = body.options;
|
||||
if (isNew) model.repositories = [];
|
||||
}
|
||||
|
||||
router.post(
|
||||
"/:conferenceID?",
|
||||
async (req: express.Request, res: express.Response) => {
|
||||
@@ -164,13 +179,7 @@ router.post(
|
||||
isOwnerOrAdmin(model.owners, user);
|
||||
}
|
||||
validateConferenceForm(req.body);
|
||||
model.name = req.body.name;
|
||||
model.startDate = new Date(req.body.startDate);
|
||||
model.endDate = new Date(req.body.endDate);
|
||||
model.status = "ready";
|
||||
model.url = req.body.url;
|
||||
model.repositories = [];
|
||||
model.options = req.body.options;
|
||||
applyConferenceForm(model, req.body, !req.params.conferenceID);
|
||||
|
||||
if (!req.params.conferenceID) {
|
||||
model.owners.push(user.model.id);
|
||||
|
||||
@@ -424,6 +424,18 @@ function updateRepoModel(
|
||||
};
|
||||
}
|
||||
|
||||
// eslint-disable-next-line @typescript-eslint/no-explicit-any
|
||||
export function hasRepositorySourceChanged(
|
||||
model: IAnonymizedRepositoryDocument,
|
||||
repoUpdate: any
|
||||
): boolean {
|
||||
return (
|
||||
repoUpdate.source.commit != model.source.commit ||
|
||||
repoUpdate.source.branch != model.source.branch ||
|
||||
repoUpdate.fullName != model.source.repositoryName
|
||||
);
|
||||
}
|
||||
|
||||
// update a repository
|
||||
router.post(
|
||||
"/:repoId/",
|
||||
@@ -441,47 +453,42 @@ router.post(
|
||||
|
||||
validateNewRepo(repoUpdate);
|
||||
|
||||
const r = gh(repoUpdate.fullName);
|
||||
if (!r?.owner || !r?.name) {
|
||||
await repo.resetSate(RepositoryStatus.ERROR, "repo_not_found");
|
||||
throw new AnonymousError("repo_not_found", {
|
||||
object: req.body,
|
||||
httpStatus: 404,
|
||||
});
|
||||
}
|
||||
|
||||
// Only the source repository/commit/branch backs the cached FileModel —
|
||||
// anonymization options (terms, image/link toggles, etc.) are applied on
|
||||
// the fly per request. Re-running the download queue is therefore only
|
||||
// needed when the underlying snapshot moves. Other edits (e.g. turning
|
||||
// off auto-update — see #360) just persist and return.
|
||||
const sourceChanged =
|
||||
repoUpdate.source.commit != repo.model.source.commit ||
|
||||
repoUpdate.source.branch != repo.model.source.branch ||
|
||||
repoUpdate.fullName != repo.model.source.repositoryName;
|
||||
const sourceChanged = hasRepositorySourceChanged(repo.model, repoUpdate);
|
||||
|
||||
updateRepoModel(repo.model, repoUpdate);
|
||||
const repository = await getRepositoryFromGitHub({
|
||||
accessToken: user.accessToken,
|
||||
owner: r.owner,
|
||||
repo: r.name,
|
||||
});
|
||||
|
||||
if (!repository) {
|
||||
await repo.resetSate(RepositoryStatus.ERROR, "repo_not_found");
|
||||
throw new AnonymousError("repo_not_found", {
|
||||
object: req.body,
|
||||
httpStatus: 404,
|
||||
});
|
||||
}
|
||||
|
||||
await repository.getCommitInfo(repoUpdate.source.commit, {
|
||||
accessToken: user.accessToken,
|
||||
});
|
||||
repo.model.source.repositoryId = repository.model.id;
|
||||
repo.model.source.repositoryName = repository.fullName || repoUpdate.fullName;
|
||||
|
||||
if (sourceChanged) {
|
||||
const parsedRepository = gh(repoUpdate.fullName);
|
||||
if (!parsedRepository?.owner || !parsedRepository?.name) {
|
||||
await repo.resetSate(RepositoryStatus.ERROR, "repo_not_found");
|
||||
throw new AnonymousError("repo_not_found", {
|
||||
object: req.body,
|
||||
httpStatus: 404,
|
||||
});
|
||||
}
|
||||
const repository = await getRepositoryFromGitHub({
|
||||
accessToken: user.accessToken,
|
||||
owner: parsedRepository.owner,
|
||||
repo: parsedRepository.name,
|
||||
});
|
||||
if (!repository) {
|
||||
await repo.resetSate(RepositoryStatus.ERROR, "repo_not_found");
|
||||
throw new AnonymousError("repo_not_found", {
|
||||
object: req.body,
|
||||
httpStatus: 404,
|
||||
});
|
||||
}
|
||||
await repository.getCommitInfo(repoUpdate.source.commit, {
|
||||
accessToken: user.accessToken,
|
||||
});
|
||||
repo.model.source.repositoryId = repository.model.id;
|
||||
repo.model.source.repositoryName =
|
||||
repository.fullName || repoUpdate.fullName;
|
||||
repo.model.anonymizeDate = new Date();
|
||||
await repo.remove();
|
||||
}
|
||||
@@ -546,9 +553,17 @@ router.post(
|
||||
},
|
||||
}
|
||||
).exec();
|
||||
if (!sourceChanged) {
|
||||
return res.json({ status: repo.status });
|
||||
}
|
||||
|
||||
await repo.updateStatus(RepositoryStatus.PREPARING);
|
||||
res.json({ status: repo.status });
|
||||
await downloadQueue.add(repo.repoId, { repoId: repo.repoId }, { jobId: `repo-${repo.repoId}` });
|
||||
await downloadQueue.add(
|
||||
repo.repoId,
|
||||
{ repoId: repo.repoId },
|
||||
{ jobId: `repo-${repo.repoId}` }
|
||||
);
|
||||
} catch (error) {
|
||||
return handleError(error, res, req);
|
||||
}
|
||||
|
||||
@@ -0,0 +1,127 @@
|
||||
const { expect } = require("chai");
|
||||
require("ts-node/register/transpile-only");
|
||||
|
||||
const ConferenceModel = require("../src/core/model/conference/conferences.model")
|
||||
.default;
|
||||
const {
|
||||
applyConferenceForm,
|
||||
} = require("../src/server/routes/conference");
|
||||
const {
|
||||
hasRepositorySourceChanged,
|
||||
} = require("../src/server/routes/repository-private");
|
||||
const {
|
||||
processRemoveRepository,
|
||||
} = require("../src/queue/processes/removeRepository");
|
||||
const {
|
||||
processRemoveCache,
|
||||
} = require("../src/queue/processes/removeCache");
|
||||
|
||||
describe("conference edits", function () {
|
||||
const form = {
|
||||
name: "Updated",
|
||||
startDate: "2026-01-01",
|
||||
endDate: "2026-02-01",
|
||||
url: "https://example.test",
|
||||
options: { expirationMode: "never" },
|
||||
};
|
||||
|
||||
it("preserves existing repository membership", function () {
|
||||
const existing = [{ id: "507f1f77bcf86cd799439011", addDate: new Date() }];
|
||||
const model = new ConferenceModel({ repositories: existing });
|
||||
applyConferenceForm(model, form, false);
|
||||
expect(model.repositories).to.have.length(1);
|
||||
expect(model.repositories[0].id.toString()).to.equal(existing[0].id);
|
||||
});
|
||||
|
||||
it("initializes repository membership for a new conference", function () {
|
||||
const model = new ConferenceModel();
|
||||
applyConferenceForm(model, form, true);
|
||||
expect(model.repositories).to.deep.equal([]);
|
||||
});
|
||||
});
|
||||
|
||||
describe("repository update source detection", function () {
|
||||
const model = {
|
||||
source: {
|
||||
commit: "abc123",
|
||||
branch: "main",
|
||||
repositoryName: "owner/repo",
|
||||
},
|
||||
};
|
||||
|
||||
it("does not redownload for an option-only edit", function () {
|
||||
expect(
|
||||
hasRepositorySourceChanged(model, {
|
||||
fullName: "owner/repo",
|
||||
source: { commit: "abc123", branch: "main" },
|
||||
options: { image: false },
|
||||
})
|
||||
).to.equal(false);
|
||||
});
|
||||
|
||||
it("redownloads when the commit, branch, or repository changes", function () {
|
||||
expect(
|
||||
hasRepositorySourceChanged(model, {
|
||||
fullName: "owner/repo",
|
||||
source: { commit: "def456", branch: "main" },
|
||||
})
|
||||
).to.equal(true);
|
||||
expect(
|
||||
hasRepositorySourceChanged(model, {
|
||||
fullName: "owner/repo",
|
||||
source: { commit: "abc123", branch: "next" },
|
||||
})
|
||||
).to.equal(true);
|
||||
expect(
|
||||
hasRepositorySourceChanged(model, {
|
||||
fullName: "other/repo",
|
||||
source: { commit: "abc123", branch: "main" },
|
||||
})
|
||||
).to.equal(true);
|
||||
});
|
||||
});
|
||||
|
||||
describe("removal workers", function () {
|
||||
const job = { data: { repoId: "repo-1" } };
|
||||
|
||||
it("rejects the repository job after recording a removal error", async function () {
|
||||
const statuses = [];
|
||||
const failure = new Error("storage unavailable");
|
||||
const repo = {
|
||||
updateStatus: async (status, message) => statuses.push([status, message]),
|
||||
remove: async () => {
|
||||
throw failure;
|
||||
},
|
||||
};
|
||||
|
||||
let caught;
|
||||
try {
|
||||
await processRemoveRepository(job, {
|
||||
connect: async () => undefined,
|
||||
getRepository: async () => repo,
|
||||
});
|
||||
} catch (error) {
|
||||
caught = error;
|
||||
}
|
||||
expect(caught).to.equal(failure);
|
||||
expect(statuses[statuses.length - 1][1]).to.equal(failure.message);
|
||||
});
|
||||
|
||||
it("rejects cache jobs when cache removal fails", async function () {
|
||||
const failure = new Error("storage unavailable");
|
||||
let caught;
|
||||
try {
|
||||
await processRemoveCache(job, {
|
||||
connect: async () => undefined,
|
||||
getRepository: async () => ({
|
||||
removeCache: async () => {
|
||||
throw failure;
|
||||
},
|
||||
}),
|
||||
});
|
||||
} catch (error) {
|
||||
caught = error;
|
||||
}
|
||||
expect(caught).to.equal(failure);
|
||||
});
|
||||
});
|
||||
Reference in New Issue
Block a user