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
1 change: 1 addition & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
### Fixed
- Made the default home page configurable with `DEFAULT_HOME_VIEW_PAGE`, defaulting to Code Search and supporting Ask. [#1677](https://github.com/sourcebot-dev/sourcebot/pull/1677)
- Require authentication for the streaming and blocking Ask APIs in Public SaaS deployments. [#1679](https://github.com/sourcebot-dev/sourcebot/pull/1679)
- Bounded BullMQ job retention to keep Redis memory from growing with repo count, retaining only the latest job per repo, connection, and account. [#1693](https://github.com/sourcebot-dev/sourcebot/pull/1693)
Comment thread
brendan-kellam marked this conversation as resolved.

## [5.1.14] - 2026-09-17

Expand Down
9 changes: 6 additions & 3 deletions packages/backend/src/connectionSyncWorkload.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -42,9 +42,12 @@ vi.mock("@sourcebot/shared", () => ({
jobOptions: {
attempts: 2,
backoff: { type: "exponential", delayMs: 5000 },
keepJobs: {
completed: { count: 50 },
failed: { count: 50 },
retention: {
mode: "window",
Comment thread
brendan-kellam marked this conversation as resolved.
keepJobs: {
completed: { count: 50 },
failed: { count: 50 },
},
},
keepLogs: 500,
},
Expand Down
78 changes: 75 additions & 3 deletions packages/backend/src/jobManager.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@ const mocks = vi.hoisted(() => {
upsertJobScheduler: vi.fn(),
getJobSchedulerIds: vi.fn(),
removeJobScheduler: vi.fn(),
trackLatestJob: vi.fn(async () => null),
producerClose: vi.fn(),
workerClose: vi.fn(),
executionLockUsing: vi.fn(),
Expand Down Expand Up @@ -54,6 +55,7 @@ vi.mock("@sourcebot/shared", () => ({
upsertJobScheduler = mocks.upsertJobScheduler;
getJobSchedulerIds = mocks.getJobSchedulerIds;
removeJobScheduler = mocks.removeJobScheduler;
trackLatestJob = mocks.trackLatestJob;
close = mocks.producerClose;
getQueue = vi.fn(() => ({
getJobCounts: vi.fn(),
Expand Down Expand Up @@ -113,9 +115,12 @@ const createWorkload = (
jobOptions: {
attempts: 2,
backoff: { type: "exponential", delayMs: 5000 },
keepJobs: {
completed: { count: 50 },
failed: { count: 50 },
retention: {
mode: "window",
Comment thread
brendan-kellam marked this conversation as resolved.
keepJobs: {
completed: { count: 50 },
failed: { count: 50 },
},
},
keepLogs: 500,
},
Expand Down Expand Up @@ -152,6 +157,27 @@ describe("BullMQJobManager lifecycle", () => {
);
});

test("rejects a latestPerResource workload that does not publish its job id in onStarted", () => {
const manager = new BullMQJobManager({} as Redis);
const workload = createWorkload();
workload.queueSpec.jobOptions = {
...workload.queueSpec.jobOptions,
retention: { mode: "latestPerResource", maxAgeSeconds: 60 },
};

expect(() => manager.register(workload)).toThrow(
/must publish its job id to the parent resource in onStarted/,
);
expect(() =>
manager.register(
createWorkload({
queueSpec: workload.queueSpec,
onStarted: vi.fn(),
}),
),
).not.toThrow();
});

test("delegates enqueueing to BullMQClient and returns its job id", async () => {
const manager = new BullMQJobManager({} as Redis);
const workload = createWorkload();
Expand Down Expand Up @@ -297,6 +323,52 @@ describe("BullMQJobManager lifecycle", () => {
});
});

test("tracks the latest job for its resource after onStarted and before processing", async () => {
const calls: string[] = [];
mocks.trackLatestJob.mockImplementation(async () => {
calls.push("tracked");
return "job-0";
});
const workload = createWorkload({
onStarted: vi.fn(async () => {
calls.push("started");
}),
process: vi.fn(async () => {
calls.push("processed");
return { outcome: "SUCCESS" };
}),
});
const manager = new BullMQJobManager({} as Redis);
manager.register(workload);
await manager.start();

await mocks.workers[0].processor({ ...job, attemptsMade: 0 });

expect(calls).toEqual(["started", "tracked", "processed"]);
expect(mocks.trackLatestJob).toHaveBeenCalledWith(
workload.queueSpec,
data,
"job-1",
);
});

test("still processes the job when tracking the latest job fails", async () => {
mocks.trackLatestJob.mockRejectedValue(new Error("redis unavailable"));
Comment thread
brendan-kellam marked this conversation as resolved.
const workload = createWorkload();
const manager = new BullMQJobManager({} as Redis);
manager.register(workload);
await manager.start();

await expect(
mocks.workers[0].processor({ ...job, attemptsMade: 0 }),
).resolves.toEqual({ outcome: "SUCCESS" });
expect(workload.process).toHaveBeenCalledTimes(1);
expect(mocks.logger.warn).toHaveBeenCalledWith(
expect.stringContaining("Failed to track latest job"),
expect.any(Error),
);
});

test("runs onStarted and processing while the execution lock is held", async () => {
const calls: string[] = [];
const workloadSignal = new AbortController().signal;
Expand Down
41 changes: 41 additions & 0 deletions packages/backend/src/jobManager.ts
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@ import {
DataOf,
JobEnqueueOptions,
QueueName,
QueueSpec,
ResultOf,
Schedule,
scheduleToMs,
Expand Down Expand Up @@ -41,6 +42,14 @@ export class BullMQJobManager implements JobManager {
if (this.workloads.has(name)) {
throw new Error(`Workload "${name}" is already registered`);
}
if (
workload.queueSpec.jobOptions.retention.mode === "latestPerResource"
&& !workload.onStarted
) {
throw new Error(
`Workload "${name}" uses latestPerResource retention and must publish its job id to the parent resource in onStarted`,
);
}
this.workloads.set(name, workload);
}

Expand Down Expand Up @@ -150,6 +159,11 @@ export class BullMQJobManager implements JobManager {

const process = async (signal: AbortSignal) => {
await workload.onStarted?.(lifecycleContext);
// Contract: `latestPerResource` workloads publish this job's id to
// their parent's `latest...JobId` pointer in onStarted (enforced in
// `register`), so the pointer has moved before the superseded job
// is removed.
await this.trackLatestJob(spec, job);
Comment thread
cursor[bot] marked this conversation as resolved.
Comment thread
brendan-kellam marked this conversation as resolved.
return workload.process({
...lifecycleContext,
signal,
Expand Down Expand Up @@ -257,6 +271,33 @@ export class BullMQJobManager implements JobManager {
}
}

private async trackLatestJob<TName extends QueueName>(
spec: QueueSpec<TName>,
job: Job,
): Promise<void> {
if (!job.id) {
return;
}
try {
const removedJobId = await this.bullmqClient.trackLatestJob(
spec,
job.data,
job.id,
);
if (removedJobId) {
logger.debug(
`Removed job ${removedJobId} superseded by ${job.id} on "${spec.name}"`,
);
}
} catch (error) {
// Retention is best-effort; the age backstop still bounds the queue.
logger.warn(
`Failed to track latest job ${job.id} on "${spec.name}"`,
error,
);
}
}

private getWorkload<TName extends QueueName>(
workloadName: TName,
): Workload<TName> {
Expand Down
22 changes: 12 additions & 10 deletions packages/backend/src/repoIndexWorkload.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -19,18 +19,10 @@ const repoFindUnique = vi.fn();
const repoUpdate = vi.fn();
const repoUpdateMany = vi.fn();

const transaction = vi.fn(async (callback: (tx: unknown) => Promise<unknown>) =>
callback({
repo: {
findUnique: repoFindUnique,
update: repoUpdate,
},
}),
);

const db = {
$transaction: transaction,
repo: {
findUnique: repoFindUnique,
update: repoUpdate,
updateMany: repoUpdateMany,
},
} as unknown as PrismaClient;
Expand Down Expand Up @@ -121,10 +113,20 @@ describe("repoIndexWorkload", () => {
});
});

test("publishes the job as the repository's latest indexing job when it starts", async () => {
await workload.onStarted!(lifecycleContext);

expect(repoUpdateMany).toHaveBeenCalledWith({
where: { id: 42 },
data: { latestIndexingJobId: "job-1" },
});
});

test("skips an INDEX job when the repository no longer exists", async () => {
await workload.process(processContext);

expect(repoUpdate).not.toHaveBeenCalled();
expect(repoUpdateMany).not.toHaveBeenCalled();
expect(lifecycleLogger.debug).toHaveBeenCalledWith(
"Skipping INDEX job for repo 42: repository no longer exists",
);
Expand Down
63 changes: 30 additions & 33 deletions packages/backend/src/repoIndexWorkload.ts
Original file line number Diff line number Diff line change
Expand Up @@ -25,13 +25,22 @@ export const createRepoIndexWorkload = ({
queueSpec: REPO_INDEX_QUEUE,
concurrency: settings.maxRepoIndexingJobConcurrency,
executionLock: REPOSITORY_EXECUTION_LOCK,
process: async ({ data, jobId, signal }) => {
onStarted: async ({ data: { repoId }, jobId }) => {
await db.repo.updateMany({
where: {
id: repoId,
},
data: {
latestIndexingJobId: jobId,
},
});
},
process: async ({ data, signal }) => {
signal.throwIfAborted();

const start = await prepareRepoIndexJob({
db,
repoId: data.repoId,
jobId,
});

if (start.action === "skip") {
Expand Down Expand Up @@ -119,45 +128,33 @@ type RepoIndexStartDecision =
const prepareRepoIndexJob = async ({
db,
repoId,
jobId,
}: {
db: PrismaClient;
repoId: number;
jobId: string;
}): Promise<RepoIndexStartDecision> =>
db.$transaction(async (tx) => {
const repo = await tx.repo.findUnique({
where: { id: repoId },
include: {
connections: {
include: {
connection: true,
},
}): Promise<RepoIndexStartDecision> => {
const repo = await db.repo.findUnique({
where: { id: repoId },
include: {
connections: {
include: {
connection: true,
},
},
});

if (!repo) {
return {
action: "skip",
reason: "repository no longer exists",
};
}

await tx.repo.update({
where: {
id: repoId,
},
data: {
latestIndexingJobId: jobId,
},
});
},
});

if (!repo) {
return {
action: "run",
repo,
action: "skip",
reason: "repository no longer exists",
};
});
}

return {
action: "run",
repo,
};
};

const indexRepository = async (
db: PrismaClient,
Expand Down
Loading
Loading