diff --git a/CHANGELOG.md b/CHANGELOG.md index c1f683d90..36ce36050 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -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) ## [5.1.14] - 2026-09-17 diff --git a/packages/backend/src/connectionSyncWorkload.test.ts b/packages/backend/src/connectionSyncWorkload.test.ts index 3b78b780a..c28027536 100644 --- a/packages/backend/src/connectionSyncWorkload.test.ts +++ b/packages/backend/src/connectionSyncWorkload.test.ts @@ -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", + keepJobs: { + completed: { count: 50 }, + failed: { count: 50 }, + }, }, keepLogs: 500, }, diff --git a/packages/backend/src/jobManager.test.ts b/packages/backend/src/jobManager.test.ts index 161e993ef..f396154e3 100644 --- a/packages/backend/src/jobManager.test.ts +++ b/packages/backend/src/jobManager.test.ts @@ -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(), @@ -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(), @@ -113,9 +115,12 @@ const createWorkload = ( jobOptions: { attempts: 2, backoff: { type: "exponential", delayMs: 5000 }, - keepJobs: { - completed: { count: 50 }, - failed: { count: 50 }, + retention: { + mode: "window", + keepJobs: { + completed: { count: 50 }, + failed: { count: 50 }, + }, }, keepLogs: 500, }, @@ -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(); @@ -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")); + 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; diff --git a/packages/backend/src/jobManager.ts b/packages/backend/src/jobManager.ts index f2c33e542..44643c087 100644 --- a/packages/backend/src/jobManager.ts +++ b/packages/backend/src/jobManager.ts @@ -6,6 +6,7 @@ import { DataOf, JobEnqueueOptions, QueueName, + QueueSpec, ResultOf, Schedule, scheduleToMs, @@ -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); } @@ -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); return workload.process({ ...lifecycleContext, signal, @@ -257,6 +271,33 @@ export class BullMQJobManager implements JobManager { } } + private async trackLatestJob( + spec: QueueSpec, + job: Job, + ): Promise { + 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( workloadName: TName, ): Workload { diff --git a/packages/backend/src/repoIndexWorkload.test.ts b/packages/backend/src/repoIndexWorkload.test.ts index ae3d84603..693901586 100644 --- a/packages/backend/src/repoIndexWorkload.test.ts +++ b/packages/backend/src/repoIndexWorkload.test.ts @@ -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) => - callback({ - repo: { - findUnique: repoFindUnique, - update: repoUpdate, - }, - }), -); - const db = { - $transaction: transaction, repo: { + findUnique: repoFindUnique, + update: repoUpdate, updateMany: repoUpdateMany, }, } as unknown as PrismaClient; @@ -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", ); diff --git a/packages/backend/src/repoIndexWorkload.ts b/packages/backend/src/repoIndexWorkload.ts index 2a15d90aa..534aee599 100644 --- a/packages/backend/src/repoIndexWorkload.ts +++ b/packages/backend/src/repoIndexWorkload.ts @@ -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") { @@ -119,45 +128,33 @@ type RepoIndexStartDecision = const prepareRepoIndexJob = async ({ db, repoId, - jobId, }: { db: PrismaClient; repoId: number; - jobId: string; -}): Promise => - db.$transaction(async (tx) => { - const repo = await tx.repo.findUnique({ - where: { id: repoId }, - include: { - connections: { - include: { - connection: true, - }, +}): Promise => { + 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, diff --git a/packages/shared/src/bullmqClient.test.ts b/packages/shared/src/bullmqClient.test.ts index f373cf5b5..a20c560e3 100644 --- a/packages/shared/src/bullmqClient.test.ts +++ b/packages/shared/src/bullmqClient.test.ts @@ -12,6 +12,9 @@ const mocks = vi.hoisted(() => ({ { key: "scheduler-2" }, ]), removeJobScheduler: vi.fn(async () => true), + removeJob: vi.fn(async () => undefined), + multiExec: vi.fn(async () => [[null, null], [null, "OK"]]), + multiSet: vi.fn(), })); vi.mock("bullmq", () => ({ @@ -32,7 +35,23 @@ vi.mock("./jobLogger.js", () => ({ })); import { BullMQClient } from "./bullmqClient.js"; -import { CONNECTION_QUEUE, type QueueSpec } from "./queue.js"; +import { + ATTACHMENT_PRUNE_QUEUE, + CONNECTION_QUEUE, + type QueueSpec, +} from "./queue.js"; + +const createRedis = () => { + const multi = { + get: vi.fn(() => multi), + set: vi.fn((...args: unknown[]) => { + mocks.multiSet(...args); + return multi; + }), + exec: mocks.multiExec, + }; + return { multi: vi.fn(() => multi) } as unknown as Redis; +}; describe("BullMQClient", () => { beforeEach(() => { @@ -202,14 +221,77 @@ describe("BullMQClient", () => { delay: 30_000, jitter: 0.5, }, - removeOnComplete: { age: 1_209_600 }, - removeOnFail: { age: 1_209_600 }, + removeOnComplete: { age: 604_800 }, + removeOnFail: { age: 604_800 }, keepLogs: 500, }, }), ); }); + describe("trackLatestJob", () => { + test("records the new job and removes the one it supersedes", async () => { + mocks.multiExec.mockResolvedValue([[null, "job-1"], [null, "OK"]]); + mocks.getJob.mockResolvedValue({ id: "job-1", remove: mocks.removeJob }); + const client = new BullMQClient(createRedis()); + + await expect( + client.trackLatestJob(CONNECTION_QUEUE, { connectionId: 42 }, "job-2"), + ).resolves.toBe("job-1"); + + expect(mocks.multiSet).toHaveBeenCalledWith( + "sourcebot:latest-job:connection-sync:connection:42", + "job-2", + "EX", + 604_800, + ); + expect(mocks.getJob).toHaveBeenCalledWith("job-1"); + expect(mocks.removeJob).toHaveBeenCalledTimes(1); + }); + + test("keeps the record when there is no previous job", async () => { + mocks.multiExec.mockResolvedValue([[null, null], [null, "OK"]]); + const client = new BullMQClient(createRedis()); + + await expect( + client.trackLatestJob(CONNECTION_QUEUE, { connectionId: 42 }, "job-2"), + ).resolves.toBeNull(); + expect(mocks.getJob).not.toHaveBeenCalled(); + }); + + test("does not remove the job on a retry of the same job", async () => { + mocks.multiExec.mockResolvedValue([[null, "job-2"], [null, "OK"]]); + const client = new BullMQClient(createRedis()); + + await expect( + client.trackLatestJob(CONNECTION_QUEUE, { connectionId: 42 }, "job-2"), + ).resolves.toBeNull(); + expect(mocks.getJob).not.toHaveBeenCalled(); + expect(mocks.removeJob).not.toHaveBeenCalled(); + }); + + test("leaves a superseded job that cannot be removed to the age backstop", async () => { + mocks.multiExec.mockResolvedValue([[null, "job-1"], [null, "OK"]]); + mocks.removeJob.mockRejectedValueOnce(new Error("Job job-1 is locked")); + mocks.getJob.mockResolvedValue({ id: "job-1", remove: mocks.removeJob }); + const client = new BullMQClient(createRedis()); + + await expect( + client.trackLatestJob(CONNECTION_QUEUE, { connectionId: 42 }, "job-2"), + ).resolves.toBeNull(); + }); + + test("is a no-op for queues with window retention", async () => { + const redis = createRedis(); + const client = new BullMQClient(redis); + + await expect( + client.trackLatestJob(ATTACHMENT_PRUNE_QUEUE, {}, "job-2"), + ).resolves.toBeNull(); + expect(redis.multi).not.toHaveBeenCalled(); + }); + }); + test("adds enqueue priority to immediate jobs", async () => { const client = new BullMQClient({} as Redis); @@ -309,8 +391,8 @@ describe("BullMQClient", () => { delay: 30_000, jitter: 0.5, }, - removeOnComplete: { age: 1_209_600 }, - removeOnFail: { age: 1_209_600 }, + removeOnComplete: { age: 604_800 }, + removeOnFail: { age: 604_800 }, keepLogs: 500, }, }, @@ -346,8 +428,8 @@ describe("BullMQClient", () => { delay: 30_000, jitter: 0.5, }, - removeOnComplete: { age: 1_209_600 }, - removeOnFail: { age: 1_209_600 }, + removeOnComplete: { age: 604_800 }, + removeOnFail: { age: 604_800 }, keepLogs: 500, }, }, diff --git a/packages/shared/src/bullmqClient.ts b/packages/shared/src/bullmqClient.ts index 19178f752..053df4e79 100644 --- a/packages/shared/src/bullmqClient.ts +++ b/packages/shared/src/bullmqClient.ts @@ -2,6 +2,7 @@ import { Queue } from "bullmq"; import { randomUUID } from "crypto"; import { Redis } from "ioredis"; import { isDeepStrictEqual } from "node:util"; +import { toBullMQKeepJobs } from "./queue.js"; import type { DataOf, JobEnqueueOptions, @@ -38,6 +39,23 @@ type WorkloadQueue = Queue< string >; +const LATEST_JOB_KEY_PREFIX = "sourcebot:latest-job"; + +// QueueSpec is distributive so queue names, data, and result schemas stay +// correlated when TName is a union. Re-establish the shared generic here +// before invoking the optional method. +const getDeduplication = ( + spec: QueueSpec, + data: DataOf, +): { id: string; keepLastIfActive?: boolean } | undefined => { + const { deduplication }: { + deduplication?( + data: DataOf, + ): { id: string; keepLastIfActive?: boolean }; + } = spec; + return deduplication?.(data); +}; + const normalizeJobState = (state: string): WorkloadJobStatus | null => { switch (state) { case "waiting": @@ -150,15 +168,8 @@ export class BullMQClient { data: DataOf, options: JobEnqueueOptions = {}, ): Promise { - // QueueSpec is distributive so queue names, data, and result schemas stay - // correlated when TName is a union. Re-establish the shared generic here - // before invoking the optional method. - const { deduplication: getDeduplication }: { - deduplication?( - data: DataOf, - ): { id: string; keepLastIfActive?: boolean }; - } = spec; - const deduplication = getDeduplication?.(data); + const deduplication = getDeduplication(spec, data); + const keepJobs = toBullMQKeepJobs(spec.jobOptions.retention); const queue = this.getQueue(spec); const requestedJobId = randomUUID(); @@ -176,8 +187,8 @@ export class BullMQClient { ? { jitter: spec.jobOptions.backoff.jitter } : {}), }, - removeOnComplete: spec.jobOptions.keepJobs.completed, - removeOnFail: spec.jobOptions.keepJobs.failed, + removeOnComplete: keepJobs.completed, + removeOnFail: keepJobs.failed, keepLogs: spec.jobOptions.keepLogs, }); @@ -199,6 +210,7 @@ export class BullMQClient { ): Promise { const queue = this.getQueue(spec); const intervalMs = scheduleToMs(schedule); + const keepJobs = toBullMQKeepJobs(spec.jobOptions.retention); const template = { name: spec.name, data, @@ -214,8 +226,8 @@ export class BullMQClient { ? { jitter: spec.jobOptions.backoff.jitter } : {}), }, - removeOnComplete: spec.jobOptions.keepJobs.completed, - removeOnFail: spec.jobOptions.keepJobs.failed, + removeOnComplete: keepJobs.completed, + removeOnFail: keepJobs.failed, keepLogs: spec.jobOptions.keepLogs, }, }; @@ -258,6 +270,57 @@ export class BullMQClient { return job.id; } + /** + * For queues with `latestPerResource` retention, records `jobId` as the + * latest job for the resource `data` belongs to and removes the job it + * supersedes. Returns the removed job id, or null when nothing was removed. + * + * Retrying the same job leaves the record untouched. A superseded job that + * is still locked (active) is left for the age backstop to reclaim. + */ + async trackLatestJob( + spec: QueueSpec, + data: DataOf, + jobId: string, + ): Promise { + const { retention } = spec.jobOptions; + if (retention.mode !== "latestPerResource") { + return null; + } + + const resourceId = getDeduplication(spec, data)?.id; + if (!resourceId) { + throw new Error( + `Workload "${spec.name}" uses latestPerResource retention but declares no deduplication id`, + ); + } + + const key = `${LATEST_JOB_KEY_PREFIX}:${spec.name}:${resourceId}`; + // The record expires with the same age backstop as the job itself, so a + // resource that stops running leaves nothing behind. + const results = await this.connection + .multi() + .get(key) + .set(key, jobId, "EX", retention.maxAgeSeconds) + .exec(); + const previousJobId = results?.[0]?.[1]; + if (typeof previousJobId !== "string" || previousJobId === jobId) { + return null; + } + + const previousJob = await this.getQueue(spec).getJob(previousJobId); + if (!previousJob) { + return null; + } + + try { + await previousJob.remove(); + } catch { + return null; + } + return previousJobId; + } + async getJobSchedulerIds( spec: QueueSpec, ): Promise { diff --git a/packages/shared/src/queue.ts b/packages/shared/src/queue.ts index a80dddcac..295253aee 100644 --- a/packages/shared/src/queue.ts +++ b/packages/shared/src/queue.ts @@ -31,13 +31,34 @@ export type JobOptions = { delayMs: number; jitter?: number; }; - keepJobs: { - completed: KeepJobs; - failed: KeepJobs; - }; + retention: JobRetention; keepLogs: number; }; +export type JobRetention = + | { + // Keeps the N most recently finished jobs in the queue (or those younger + // than `age`), regardless of which resource they belong to. + mode: "window"; + keepJobs: { + completed: KeepJobs; + failed: KeepJobs; + }; + } + | { + // Keeps only the most recent job for each resource, keyed by the queue's + // deduplication id. When a job starts, the job it supersedes is removed. + // Workloads on these queues must publish the job id to their parent's + // `latest...JobId` pointer in `onStarted`, so the pointer never + // references a removed job. + // `maxAgeSeconds` is an age-only backstop that reclaims jobs whose + // resource stopped running (e.g. a deleted repo). It must exceed the + // longest scheduler interval, or a resource's latest job can be + // reclaimed before its successor starts. + mode: "latestPerResource"; + maxAgeSeconds: number; + }; + export type JobEnqueueOptions = { priority?: number; }; @@ -48,10 +69,9 @@ export const JOB_PRIORITIES = { SCHEDULED: 10, } as const; -// BullMQ evaluates age-based cleanup only when another job reaches the same -// terminal state. A lone completed or failed job therefore remains available -// past this age; once a newer same-state job finishes, it replaces the old one. -const TWO_WEEKS_IN_SECONDS = 14 * 24 * 60 * 60; +const ONE_DAY_IN_SECONDS = 24 * 60 * 60; +const ONE_WEEK_IN_SECONDS = 7 * ONE_DAY_IN_SECONDS; +const TWO_WEEKS_IN_SECONDS = 14 * ONE_DAY_IN_SECONDS; export const DEFAULT_JOB_OPTIONS: JobOptions = { attempts: 2, @@ -60,13 +80,48 @@ export const DEFAULT_JOB_OPTIONS: JobOptions = { delayMs: 30_000, jitter: 0.5, }, - keepJobs: { - completed: { age: TWO_WEEKS_IN_SECONDS }, - failed: { age: TWO_WEEKS_IN_SECONDS }, + retention: { + mode: "window", + keepJobs: { + completed: { + age: ONE_DAY_IN_SECONDS, + count: 5_000, + }, + failed: { + age: TWO_WEEKS_IN_SECONDS, + count: 10_000, + }, + }, }, keepLogs: DEFAULT_JOB_LOGS_MAX_ENTRIES, }; +// For queues whose latest job per resource is resolved by the web app through +// a `latest...JobId` pointer. Window retention cannot guarantee that job +// survives, since the window is shared by every resource in the queue. +export const LATEST_PER_RESOURCE_JOB_OPTIONS: JobOptions = { + ...DEFAULT_JOB_OPTIONS, + retention: { + mode: "latestPerResource", + maxAgeSeconds: ONE_WEEK_IN_SECONDS, + }, +}; + +/** + * Translates a queue's retention policy into BullMQ's per-job removal options. + */ +export const toBullMQKeepJobs = ( + retention: JobRetention, +): { completed: KeepJobs; failed: KeepJobs } => { + if (retention.mode === "window") { + return retention.keepJobs; + } + return { + completed: { age: retention.maxAgeSeconds }, + failed: { age: retention.maxAgeSeconds }, + }; +}; + export type QueueName = keyof QueueRegistry; export type DataOf = QueueRegistry[TName]["data"]; export type ResultOf = QueueRegistry[TName]["result"]; @@ -143,13 +198,13 @@ export const AUDIT_LOG_PRUNE_QUEUE: QueueSpec<"audit-log-prune"> = { export const CONNECTION_QUEUE: QueueSpec<"connection-sync"> = { name: "connection-sync", resultSchema: connectionSyncResultSchema, - jobOptions: DEFAULT_JOB_OPTIONS, + jobOptions: LATEST_PER_RESOURCE_JOB_OPTIONS, deduplication: (data) => ({ id: `connection:${data.connectionId}` }), }; export const REPO_INDEX_QUEUE: QueueSpec<"repo-index"> = { name: "repo-index", - jobOptions: DEFAULT_JOB_OPTIONS, + jobOptions: LATEST_PER_RESOURCE_JOB_OPTIONS, deduplication: ({ repoId }) => ({ id: `repo:${repoId}`, keepLastIfActive: true, @@ -167,14 +222,14 @@ export const REPO_CLEANUP_QUEUE: QueueSpec<"repo-cleanup"> = { export const ACCOUNT_PERMISSION_SYNC_QUEUE: QueueSpec<"account-permission-sync"> = { name: "account-permission-sync", - jobOptions: DEFAULT_JOB_OPTIONS, + jobOptions: LATEST_PER_RESOURCE_JOB_OPTIONS, deduplication: (data) => ({ id: `account:${data.accountId}` }), }; export const REPO_PERMISSION_SYNC_QUEUE: QueueSpec<"repo-permission-sync"> = { name: "repo-permission-sync", resultSchema: repoPermissionSyncResultSchema, - jobOptions: DEFAULT_JOB_OPTIONS, + jobOptions: LATEST_PER_RESOURCE_JOB_OPTIONS, deduplication: (data) => ({ id: `repo:${data.repoId}` }), };