From efae78056e122001056846a03ec7ebe6bb4c6535 Mon Sep 17 00:00:00 2001 From: Hampus Date: Sun, 30 Aug 2026 22:16:18 +0200 Subject: [PATCH] fix(worker): park far-future jobs in a KV due queue (#2144) --- .../KVScheduledJobQueueService.ts | 113 +++++++ .../src/api/middleware/ServiceSingletons.ts | 13 + .../src/api/worker/WorkerDependencies.ts | 5 + fluxer_api/src/api/worker/WorkerLaneConfig.ts | 16 +- fluxer_api/src/api/worker/WorkerMain.ts | 9 + fluxer_api/src/api/worker/WorkerRunner.ts | 42 +-- .../src/api/worker/WorkerTaskRegistry.ts | 2 + .../worker/tasks/ProcessScheduledJobQueue.ts | 80 +++++ .../worker/tests/ScheduledJobQueue.test.ts | 279 ++++++++++++++++++ .../api/worker/tests/WorkerLaneConfig.test.ts | 13 +- 10 files changed, 551 insertions(+), 21 deletions(-) create mode 100644 fluxer_api/src/api/infrastructure/KVScheduledJobQueueService.ts create mode 100644 fluxer_api/src/api/worker/tasks/ProcessScheduledJobQueue.ts create mode 100644 fluxer_api/src/api/worker/tests/ScheduledJobQueue.test.ts diff --git a/fluxer_api/src/api/infrastructure/KVScheduledJobQueueService.ts b/fluxer_api/src/api/infrastructure/KVScheduledJobQueueService.ts new file mode 100644 index 000000000..e7b3761fb --- /dev/null +++ b/fluxer_api/src/api/infrastructure/KVScheduledJobQueueService.ts @@ -0,0 +1,113 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later + +import {createHash} from 'node:crypto'; +import type {IKVProvider} from '@pkgs/kv_client/src/IKVProvider'; +import type {WorkerJobPayload} from '@pkgs/worker/src/contracts/WorkerTypes'; +import {Logger} from '../Logger'; + +export interface ParkedScheduledJob { + jobIdentity: string; + taskType: string; + payload: WorkerJobPayload; + runAtMs: number; + ledgerJobId: string | null; +} + +const QUEUE_KEY = 'scheduled_job_queue'; +const SECONDARY_KEY_PREFIX = 'scheduled_job_queue:'; + +export function buildScheduledJobIdentity( + taskType: string, + payload: WorkerJobPayload, + runAtMs: number, + ledgerJobId: string | null, +): string { + if (ledgerJobId !== null) { + return ledgerJobId; + } + return createHash('sha256') + .update(`${taskType}|${runAtMs}|${JSON.stringify(payload)}`) + .digest('hex'); +} + +export class KVScheduledJobQueueService { + constructor(private readonly kvClient: IKVProvider) {} + + private getSecondaryKey(jobIdentity: string): string { + return `${SECONDARY_KEY_PREFIX}${jobIdentity}`; + } + + private serializeQueueItem(job: ParkedScheduledJob): string { + return JSON.stringify(job); + } + + private deserializeQueueItem(value: string): ParkedScheduledJob { + const parsed: unknown = JSON.parse(value); + if (typeof parsed !== 'object' || parsed === null) { + throw new Error('parked scheduled job must be a JSON object'); + } + const {jobIdentity, taskType, payload, runAtMs, ledgerJobId} = parsed as Record; + if (typeof jobIdentity !== 'string' || typeof taskType !== 'string' || typeof runAtMs !== 'number') { + throw new Error('parked scheduled job is missing required fields'); + } + if (typeof payload !== 'object' || payload === null || Array.isArray(payload)) { + throw new Error('parked scheduled job payload must be a JSON object'); + } + return { + jobIdentity, + taskType, + payload: payload as WorkerJobPayload, + runAtMs, + ledgerJobId: typeof ledgerJobId === 'string' ? ledgerJobId : null, + }; + } + + async parkJob(job: ParkedScheduledJob, releaseAt: Date): Promise { + try { + const secondaryKey = this.getSecondaryKey(job.jobIdentity); + const value = this.serializeQueueItem(job); + await this.kvClient.removeBulkDeletion(QUEUE_KEY, secondaryKey); + await this.kvClient.scheduleBulkDeletion(QUEUE_KEY, secondaryKey, releaseAt.getTime(), value); + Logger.debug({jobIdentity: job.jobIdentity, taskType: job.taskType, releaseAt}, 'Parked scheduled job'); + } catch (error) { + Logger.error({error, jobIdentity: job.jobIdentity, taskType: job.taskType}, 'Failed to park scheduled job'); + throw error; + } + } + + async claimJob(jobIdentity: string): Promise { + try { + return await this.kvClient.removeBulkDeletion(QUEUE_KEY, this.getSecondaryKey(jobIdentity)); + } catch (error) { + Logger.error({error, jobIdentity}, 'Failed to claim parked scheduled job'); + throw error; + } + } + + async getReadyJobs(nowMs: number, limit: number): Promise> { + try { + const results = await this.kvClient.zrangebyscore(QUEUE_KEY, '-inf', nowMs, 'LIMIT', 0, limit); + const jobs: Array = []; + for (const result of results) { + try { + jobs.push(this.deserializeQueueItem(result)); + } catch (error) { + Logger.error({error, result}, 'Failed to parse parked scheduled job entry'); + } + } + return jobs; + } catch (error) { + Logger.error({error, nowMs, limit}, 'Failed to fetch ready parked scheduled jobs'); + throw error; + } + } + + async getQueueSize(): Promise { + try { + return await this.kvClient.zcard(QUEUE_KEY); + } catch (error) { + Logger.error({error}, 'Failed to get parked scheduled job queue size'); + throw error; + } + } +} diff --git a/fluxer_api/src/api/middleware/ServiceSingletons.ts b/fluxer_api/src/api/middleware/ServiceSingletons.ts index 727ada566..6a8b72ae5 100644 --- a/fluxer_api/src/api/middleware/ServiceSingletons.ts +++ b/fluxer_api/src/api/middleware/ServiceSingletons.ts @@ -58,6 +58,7 @@ import type {IUnfurlerService} from '../infrastructure/IUnfurlerService'; import {KVAccountDeletionQueueService} from '../infrastructure/KVAccountDeletionQueueService'; import {KVActivityTracker} from '../infrastructure/KVActivityTracker'; import {KVBulkMessageDeletionQueueService} from '../infrastructure/KVBulkMessageDeletionQueueService'; +import {KVScheduledJobQueueService} from '../infrastructure/KVScheduledJobQueueService'; import {NatsUnfurlerService} from '../infrastructure/NatsUnfurlerService'; import {PremiumStateReconciliationQueueService} from '../infrastructure/PremiumStateReconciliationQueueService'; import {createDownloadsStorageService, createStorageService} from '../infrastructure/StorageServiceFactory'; @@ -222,6 +223,18 @@ export function getKVBulkMessageDeletionQueue(): KVBulkMessageDeletionQueueServi return bulkMessageDeletionQueue; } +let scheduledJobQueueClient: IKVProvider | null = null; +let scheduledJobQueue: KVScheduledJobQueueService | null = null; + +export function getKVScheduledJobQueue(): KVScheduledJobQueueService { + const kvClient = getKVClient(); + if (!scheduledJobQueue || scheduledJobQueueClient !== kvClient) { + scheduledJobQueue = new KVScheduledJobQueueService(kvClient); + scheduledJobQueueClient = kvClient; + } + return scheduledJobQueue; +} + let premiumStateQueueClient: IKVProvider | null = null; let premiumStateQueue: PremiumStateReconciliationQueueService | null = null; diff --git a/fluxer_api/src/api/worker/WorkerDependencies.ts b/fluxer_api/src/api/worker/WorkerDependencies.ts index f30ca8caf..9e7163f6b 100644 --- a/fluxer_api/src/api/worker/WorkerDependencies.ts +++ b/fluxer_api/src/api/worker/WorkerDependencies.ts @@ -40,6 +40,7 @@ import type {IVoiceRoomStore} from '../infrastructure/IVoiceRoomStore'; import type {KVAccountDeletionQueueService} from '../infrastructure/KVAccountDeletionQueueService'; import type {KVActivityTracker} from '../infrastructure/KVActivityTracker'; import type {KVBulkMessageDeletionQueueService} from '../infrastructure/KVBulkMessageDeletionQueueService'; +import type {KVScheduledJobQueueService} from '../infrastructure/KVScheduledJobQueueService'; import type {PremiumStateReconciliationQueueService} from '../infrastructure/PremiumStateReconciliationQueueService'; import type {UserCacheService} from '../infrastructure/UserCacheService'; import type {InstanceConfigRepository} from '../instance/InstanceConfigRepository'; @@ -86,6 +87,7 @@ import { getKVAccountDeletionQueue, getKVActivityTracker, getKVBulkMessageDeletionQueue, + getKVScheduledJobQueue, getLimitConfigService, getNcmecSubmissionService, getOAuth2TokenRepository, @@ -162,6 +164,7 @@ export interface WorkerDependencies { activityTracker: KVActivityTracker; deletionQueueService: KVAccountDeletionQueueService; bulkMessageDeletionQueueService: KVBulkMessageDeletionQueueService; + scheduledJobQueueService: KVScheduledJobQueueService; premiumStateReconciliationQueueService: PremiumStateReconciliationQueueService; deletionEligibilityService: UserDeletionEligibilityService; voiceRoomStore: IVoiceRoomStore; @@ -226,6 +229,7 @@ export async function initializeWorkerDependencies(snowflakeService: ISnowflakeS const activityTracker = getKVActivityTracker(); const deletionQueueService = getKVAccountDeletionQueue(); const bulkMessageDeletionQueueService = getKVBulkMessageDeletionQueue(); + const scheduledJobQueueService = getKVScheduledJobQueue(); const premiumStateReconciliationQueueService = getPremiumStateReconciliationQueueService(); const deletionEligibilityService = new UserDeletionEligibilityService(kvClient); await ensureVoiceResourcesInitialized(); @@ -337,6 +341,7 @@ export async function initializeWorkerDependencies(snowflakeService: ISnowflakeS activityTracker, deletionQueueService, bulkMessageDeletionQueueService, + scheduledJobQueueService, premiumStateReconciliationQueueService, deletionEligibilityService, voiceRoomStore, diff --git a/fluxer_api/src/api/worker/WorkerLaneConfig.ts b/fluxer_api/src/api/worker/WorkerLaneConfig.ts index be711f989..a821f7394 100644 --- a/fluxer_api/src/api/worker/WorkerLaneConfig.ts +++ b/fluxer_api/src/api/worker/WorkerLaneConfig.ts @@ -72,6 +72,7 @@ const LANE_CONFIG = { 'processInactivityDeletions', 'processPendingBulkMessageDeletions', 'processPremiumStateReconciliationQueue', + 'processScheduledJobQueue', 'prunePostgresKvTtl', 'refreshSearchIndex', 'syncDiscoveryIndex', @@ -213,6 +214,13 @@ function resolveWorkerLanes(config: WorkerLaneRuntimeConfig): Array lane.ackWaitMs)); + return Math.max(1, Math.floor(minimumAckWaitMs / SCHEDULED_JOB_DRAIN_INTERVALS_PER_ACK_WAIT / 1000)); +} + function resolveCronSchedulerEnabled(mode: APIWorkerMode, configuredValue: boolean | undefined): boolean { if (configuredValue !== undefined) { return configuredValue; @@ -245,4 +253,10 @@ export function findLaneForTask(taskName: string): APIWorkerLaneName | null { } export type {WorkerLaneDefinition}; -export {resolveCronSchedulerEnabled, resolveWorkerLanes, validateLaneCompleteness, WORKER_LANES}; +export { + resolveCronSchedulerEnabled, + resolveScheduledJobDrainIntervalSeconds, + resolveWorkerLanes, + validateLaneCompleteness, + WORKER_LANES, +}; diff --git a/fluxer_api/src/api/worker/WorkerMain.ts b/fluxer_api/src/api/worker/WorkerMain.ts index b42f1c8c8..b66180072 100644 --- a/fluxer_api/src/api/worker/WorkerMain.ts +++ b/fluxer_api/src/api/worker/WorkerMain.ts @@ -25,6 +25,7 @@ import {clearWorkerDependencies, setWorkerDependencies} from './WorkerContext'; import {initializeWorkerDependencies, shutdownWorkerDependencies, type WorkerDependencies} from './WorkerDependencies'; import { resolveCronSchedulerEnabled, + resolveScheduledJobDrainIntervalSeconds, resolveWorkerLanes, validateLaneCompleteness, type WorkerLaneDefinition, @@ -50,6 +51,13 @@ function registerCronJobs(cron: CronScheduler): void { cron.upsert('processPremiumStateReconciliationQueue', 'processPremiumStateReconciliationQueue', {}, '0 * * * * *', { ledger: false, }); + cron.upsert( + 'processScheduledJobQueue', + 'processScheduledJobQueue', + {}, + `*/${resolveScheduledJobDrainIntervalSeconds()} * * * * *`, + {ledger: false}, + ); cron.upsert('processExpiredPremiumSweep', 'processExpiredPremiumSweep', {}, '0 0 * * * *', {ledger: false}); cron.upsert('processInactivityDeletions', 'processInactivityDeletions', {}, '0 0 */6 * * *', {ledger: false}); cron.upsert('expireAttachments', 'expireAttachments', {}, '0 0 */12 * * *', {ledger: false}); @@ -221,6 +229,7 @@ export async function startWorkerMain(): Promise { const runner = new WorkerRunner({ tasks: laneTasks, queue, + scheduledJobQueue: dependencies.scheduledJobQueueService, consumerName: lane.consumerName, laneName: lane.name, ledger: jobLedger, diff --git a/fluxer_api/src/api/worker/WorkerRunner.ts b/fluxer_api/src/api/worker/WorkerRunner.ts index 2fa6a28c9..37ab3acab 100644 --- a/fluxer_api/src/api/worker/WorkerRunner.ts +++ b/fluxer_api/src/api/worker/WorkerRunner.ts @@ -3,8 +3,8 @@ import {randomUUID} from 'node:crypto'; import type {IWorkerService} from '@pkgs/worker/src/contracts/IWorkerService'; import {JobCancelledError, type WorkerTaskHandler} from '@pkgs/worker/src/contracts/WorkerTask'; -import type {WorkerJobPayload} from '@pkgs/worker/src/contracts/WorkerTypes'; import type {ConsumerMessages, JsMsg} from 'nats'; +import {buildScheduledJobIdentity, type ParkedScheduledJob} from '../infrastructure/KVScheduledJobQueueService'; import type {IJobLedgerRepository} from '../jobs/IJobLedgerRepository'; import {Logger} from '../Logger'; import {getWorkerService} from '../middleware/ServiceRegistry'; @@ -38,22 +38,17 @@ interface WorkerRunnerDlqMeta { interface WorkerRunnerQueue { getConnectionManager(): WorkerRunnerConnectionManager; getStreamName(): string; - enqueue( - taskType: string, - payload: WorkerJobPayload, - options?: { - runAt?: Date; - maxAttempts?: number; - priority?: number; - jobKey?: string; - }, - ): Promise; publishToDlq(taskType: string, originalPayload: Record, meta: WorkerRunnerDlqMeta): Promise; } +interface WorkerRunnerScheduledJobQueue { + parkJob(job: ParkedScheduledJob, releaseAt: Date): Promise; +} + interface WorkerRunnerOptions { tasks: Record; queue: WorkerRunnerQueue; + scheduledJobQueue: WorkerRunnerScheduledJobQueue; consumerName: string; laneName: string; ledger: IJobLedgerRepository; @@ -66,6 +61,7 @@ interface WorkerRunnerOptions { export class WorkerRunner { private readonly tasks: Record; private readonly queue: WorkerRunnerQueue; + private readonly scheduledJobQueue: WorkerRunnerScheduledJobQueue; private readonly consumerName: string; private readonly laneName: string; private readonly workerId: string; @@ -81,6 +77,7 @@ export class WorkerRunner { constructor(options: WorkerRunnerOptions) { this.tasks = options.tasks; this.queue = options.queue; + this.scheduledJobQueue = options.scheduledJobQueue; this.consumerName = options.consumerName; this.laneName = options.laneName; this.workerId = options.workerId ?? `worker-${options.laneName}-${randomUUID()}`; @@ -204,20 +201,27 @@ export class WorkerRunner { const delayMs = runAtMs - Date.now(); if (delayMs > 0) { const deliveryCount = msg.info.deliveryCount; - const shouldReEnqueue = delayMs > this.ackWaitMs || deliveryCount >= this.maxDeliver - 1; - if (shouldReEnqueue) { + const shouldPark = delayMs > this.ackWaitMs || deliveryCount >= this.maxDeliver - 1; + if (shouldPark) { + const ledgerJobIdValue = ledgerJobId === null ? null : ledgerJobId.toString(); try { - await this.queue.enqueue(taskType, jobPayload, {runAt: new Date(runAtMs)}); + await this.scheduledJobQueue.parkJob( + { + jobIdentity: buildScheduledJobIdentity(taskType, jobPayload, runAtMs, ledgerJobIdValue), + taskType, + payload: jobPayload, + runAtMs, + ledgerJobId: ledgerJobIdValue, + }, + new Date(runAtMs - this.ackWaitMs), + ); msg.ack(); Logger.debug( {taskType, seq: msg.seq, runAt, deliveryCount}, - 'Re-enqueued scheduled job to free ack slot', + 'Parked scheduled job in the due queue to free ack slot', ); } catch (error) { - Logger.error( - {taskType, seq: msg.seq, err: error}, - 'Failed to re-enqueue scheduled job, falling back to NAK', - ); + Logger.error({taskType, seq: msg.seq, err: error}, 'Failed to park scheduled job, falling back to NAK'); msg.nak(Math.min(delayMs, this.ackWaitMs - 5000)); } } else { diff --git a/fluxer_api/src/api/worker/WorkerTaskRegistry.ts b/fluxer_api/src/api/worker/WorkerTaskRegistry.ts index 28066ea7c..cc875676a 100644 --- a/fluxer_api/src/api/worker/WorkerTaskRegistry.ts +++ b/fluxer_api/src/api/worker/WorkerTaskRegistry.ts @@ -30,6 +30,7 @@ import processExpiredPremiumSweep from './tasks/ProcessExpiredPremiumSweep'; import processInactivityDeletions from './tasks/ProcessInactivityDeletions'; import processPendingBulkMessageDeletions from './tasks/ProcessPendingBulkMessageDeletions'; import processPremiumStateReconciliationQueue from './tasks/ProcessPremiumStateReconciliationQueue'; +import processScheduledJobQueue from './tasks/ProcessScheduledJobQueue'; import processStripeWebhook from './tasks/ProcessStripeWebhook'; import prunePostgresKvTtl from './tasks/PrunePostgresKvTtl'; import reconcileUserPayments from './tasks/ReconcileUserPayments'; @@ -75,6 +76,7 @@ export const workerTasks: Record = { processInactivityDeletions, processPendingBulkMessageDeletions, processPremiumStateReconciliationQueue, + processScheduledJobQueue, reconcileUserPayments, prunePostgresKvTtl, refreshSearchIndex, diff --git a/fluxer_api/src/api/worker/tasks/ProcessScheduledJobQueue.ts b/fluxer_api/src/api/worker/tasks/ProcessScheduledJobQueue.ts new file mode 100644 index 000000000..6a2f30999 --- /dev/null +++ b/fluxer_api/src/api/worker/tasks/ProcessScheduledJobQueue.ts @@ -0,0 +1,80 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later + +import type {WorkerTaskHandler} from '@pkgs/worker/src/contracts/WorkerTask'; +import type {WorkerJobPayload} from '@pkgs/worker/src/contracts/WorkerTypes'; +import type {ParkedScheduledJob} from '../../infrastructure/KVScheduledJobQueueService'; +import {Logger} from '../../Logger'; +import {getWorkerDependencies} from '../WorkerContext'; +import {findLaneForTask, type WorkerTaskName} from '../WorkerLaneConfig'; + +const MAX_JOBS_PER_RUN = 500; + +function buildReleasePayload(job: ParkedScheduledJob): WorkerJobPayload { + if (job.ledgerJobId === null) { + return job.payload; + } + return {...job.payload, __jobId: job.ledgerJobId}; +} + +const processScheduledJobQueue: WorkerTaskHandler = async (_payload, helpers) => { + const {scheduledJobQueueService, workerService} = getWorkerDependencies(); + const nowMs = Date.now(); + const readyJobs = await scheduledJobQueueService.getReadyJobs(nowMs, MAX_JOBS_PER_RUN); + if (readyJobs.length === 0) { + helpers.logger.debug('No parked scheduled jobs are due'); + return; + } + let releasedCount = 0; + let claimedElsewhereCount = 0; + let droppedCount = 0; + let failedCount = 0; + for (const job of readyJobs) { + try { + const claimed = await scheduledJobQueueService.claimJob(job.jobIdentity); + if (!claimed) { + claimedElsewhereCount += 1; + continue; + } + if (findLaneForTask(job.taskType) === null) { + droppedCount += 1; + Logger.error( + {jobIdentity: job.jobIdentity, taskType: job.taskType}, + 'Dropping parked scheduled job with an unknown task type', + ); + continue; + } + await workerService.addJob(job.taskType as WorkerTaskName, buildReleasePayload(job), { + runAt: new Date(job.runAtMs), + jobKey: job.jobIdentity, + skipLedger: true, + }); + releasedCount += 1; + } catch (error) { + failedCount += 1; + Logger.error( + {error, jobIdentity: job.jobIdentity, taskType: job.taskType}, + 'Failed to release parked scheduled job, re-parking for the next drain', + ); + try { + await scheduledJobQueueService.parkJob(job, new Date(nowMs)); + } catch (parkError) { + Logger.error( + {error: parkError, jobIdentity: job.jobIdentity, taskType: job.taskType}, + 'Failed to re-park scheduled job after a release failure', + ); + } + } + } + helpers.logger.info( + { + total: readyJobs.length, + releasedCount, + claimedElsewhereCount, + droppedCount, + failedCount, + }, + 'Released parked scheduled jobs', + ); +}; + +export default processScheduledJobQueue; diff --git a/fluxer_api/src/api/worker/tests/ScheduledJobQueue.test.ts b/fluxer_api/src/api/worker/tests/ScheduledJobQueue.test.ts new file mode 100644 index 000000000..b9a9417a1 --- /dev/null +++ b/fluxer_api/src/api/worker/tests/ScheduledJobQueue.test.ts @@ -0,0 +1,279 @@ +// SPDX-License-Identifier: AGPL-3.0-or-later + +import type {IWorkerService} from '@pkgs/worker/src/contracts/IWorkerService'; +import type {WorkerTaskHelpers} from '@pkgs/worker/src/contracts/WorkerTask'; +import type {WorkerJobOptions, WorkerJobPayload} from '@pkgs/worker/src/contracts/WorkerTypes'; +import type {JsMsg} from 'nats'; +import {afterEach, beforeAll, describe, expect, it, vi} from 'vitest'; +import {KVScheduledJobQueueService} from '../../infrastructure/KVScheduledJobQueueService'; +import type {IJobLedgerRepository} from '../../jobs/IJobLedgerRepository'; +import {setInjectedWorkerService} from '../../middleware/ServiceRegistry'; +import {MockKVProvider} from '../../test/mocks/MockKVProvider'; +import {NoopLogger} from '../../test/mocks/NoopLogger'; +import {NoopWorkerService} from '../../test/NoopWorkerService'; +import processScheduledJobQueue from '../tasks/ProcessScheduledJobQueue'; +import {setWorkerDependenciesForTest} from '../WorkerContext'; +import {WorkerRunner} from '../WorkerRunner'; + +const LANE_ACK_WAIT_MS = 60000; +const LANE_MAX_DELIVER = 25; +const TASK_TYPE = 'sendScheduledMessage'; + +interface RecordedJob { + taskType: string; + payload: WorkerJobPayload; + options: WorkerJobOptions | undefined; +} + +class RecordingWorkerService implements IWorkerService { + readonly jobs: Array = []; + + async addJob( + taskType: string, + payload: TPayload, + options?: WorkerJobOptions, + ): Promise { + this.jobs.push({taskType, payload, options}); + return BigInt(this.jobs.length); + } + + async cancelJob(_jobId: bigint): Promise { + return false; + } + + async retryDeadLetterJob(_jobId: bigint): Promise { + return false; + } +} + +const queueStub = { + getConnectionManager: () => { + throw new Error('WorkerRunner tests never consume messages'); + }, + getStreamName: () => 'JOBS', + enqueue: vi.fn(), + publishToDlq: vi.fn(), +}; + +class ConcurrentlyClaimedKVProvider extends MockKVProvider { + override async removeBulkDeletion(queueKey: string, secondaryKey: string): Promise { + await super.removeBulkDeletion(queueKey, secondaryKey); + return false; + } +} + +class TestWorkerRunner extends WorkerRunner { + async runJob(taskType: string, msg: JsMsg): Promise { + return await this.processJob(taskType, msg); + } +} + +function createRunner(scheduledJobQueue: KVScheduledJobQueueService): TestWorkerRunner { + return new TestWorkerRunner({ + tasks: {[TASK_TYPE]: async () => {}}, + queue: queueStub, + scheduledJobQueue, + consumerName: 'workers_lifecycle', + laneName: 'lifecycle', + ledger: {} as IJobLedgerRepository, + concurrency: 8, + maxDeliver: LANE_MAX_DELIVER, + ackWaitMs: LANE_ACK_WAIT_MS, + }); +} + +function createJobMessage(options: { + runAt: Date; + payload: Record; + seq?: number; + deliveryCount?: number; +}) { + const envelope = { + payload: options.payload, + run_at: options.runAt.toISOString(), + max_attempts: 5, + priority: 0, + created_at: new Date().toISOString(), + }; + return { + seq: options.seq ?? 1, + subject: `jobs.${TASK_TYPE}`, + redelivered: (options.deliveryCount ?? 1) > 1, + data: new TextEncoder().encode(JSON.stringify(envelope)), + info: {deliveryCount: options.deliveryCount ?? 1}, + ack: vi.fn(), + nak: vi.fn(), + term: vi.fn(), + }; +} + +function createHelpers(): WorkerTaskHelpers { + return { + logger: new NoopLogger(), + jobId: 0n, + addJob: async () => 0n, + reportProgress: async () => {}, + shouldCancel: async () => false, + setContextLink: async () => {}, + }; +} + +async function drain(scheduledJobQueue: KVScheduledJobQueueService, workerService: RecordingWorkerService) { + setWorkerDependenciesForTest({scheduledJobQueueService: scheduledJobQueue, workerService}); + await processScheduledJobQueue({}, createHelpers()); +} + +describe('Scheduled job due queue', () => { + beforeAll(() => { + setInjectedWorkerService(new NoopWorkerService()); + }); + + afterEach(() => { + vi.useRealTimers(); + queueStub.enqueue.mockClear(); + queueStub.publishToDlq.mockClear(); + }); + + it('naks a job due inside one ack-wait instead of parking it', async () => { + const scheduledJobQueue = new KVScheduledJobQueueService(new MockKVProvider()); + const runner = createRunner(scheduledJobQueue); + const msg = createJobMessage({ + runAt: new Date(Date.now() + LANE_ACK_WAIT_MS - 30000), + payload: {scheduledMessageId: '77', __jobId: '4242'}, + }); + + await runner.runJob(TASK_TYPE, msg as unknown as JsMsg); + + expect(msg.nak).toHaveBeenCalledTimes(1); + expect(msg.ack).not.toHaveBeenCalled(); + const delayMs = msg.nak.mock.calls[0]![0] as number; + expect(delayMs).toBeGreaterThan(0); + expect(delayMs).toBeLessThanOrEqual(LANE_ACK_WAIT_MS - 30000); + expect(await scheduledJobQueue.getQueueSize()).toBe(0); + }); + + it('parks a job due beyond one ack-wait without republishing it', async () => { + const scheduledJobQueue = new KVScheduledJobQueueService(new MockKVProvider()); + const runner = createRunner(scheduledJobQueue); + const runAt = new Date(Date.now() + 60 * 60 * 1000); + const msg = createJobMessage({runAt, payload: {scheduledMessageId: '77', __jobId: '4242'}}); + + await runner.runJob(TASK_TYPE, msg as unknown as JsMsg); + + expect(msg.ack).toHaveBeenCalledTimes(1); + expect(msg.nak).not.toHaveBeenCalled(); + expect(queueStub.enqueue).not.toHaveBeenCalled(); + expect(await scheduledJobQueue.getQueueSize()).toBe(1); + expect(await scheduledJobQueue.getReadyJobs(Date.now(), 10)).toHaveLength(0); + const due = await scheduledJobQueue.getReadyJobs(runAt.getTime() - LANE_ACK_WAIT_MS, 10); + expect(due).toHaveLength(1); + expect(due[0]).toEqual({ + jobIdentity: '4242', + taskType: TASK_TYPE, + payload: {scheduledMessageId: '77'}, + runAtMs: runAt.getTime(), + ledgerJobId: '4242', + }); + }); + + it('releases a parked job exactly once when its due time arrives', async () => { + vi.useFakeTimers({toFake: ['Date']}); + const start = new Date('2026-08-30T12:00:00.000Z'); + vi.setSystemTime(start); + const scheduledJobQueue = new KVScheduledJobQueueService(new MockKVProvider()); + const runner = createRunner(scheduledJobQueue); + const runAt = new Date(start.getTime() + 60 * 60 * 1000); + const msg = createJobMessage({runAt, payload: {scheduledMessageId: '77', __jobId: '4242'}}); + await runner.runJob(TASK_TYPE, msg as unknown as JsMsg); + const workerService = new RecordingWorkerService(); + + await drain(scheduledJobQueue, workerService); + expect(workerService.jobs).toHaveLength(0); + + vi.setSystemTime(new Date(runAt.getTime() - LANE_ACK_WAIT_MS + 1000)); + await drain(scheduledJobQueue, workerService); + + expect(workerService.jobs).toHaveLength(1); + expect(workerService.jobs[0]!.taskType).toBe(TASK_TYPE); + expect(workerService.jobs[0]!.payload).toEqual({scheduledMessageId: '77', __jobId: '4242'}); + expect(workerService.jobs[0]!.options?.runAt?.getTime()).toBe(runAt.getTime()); + expect(workerService.jobs[0]!.options?.jobKey).toBe('4242'); + expect(workerService.jobs[0]!.options?.skipLedger).toBe(true); + expect(await scheduledJobQueue.getQueueSize()).toBe(0); + + await drain(scheduledJobQueue, workerService); + expect(workerService.jobs).toHaveLength(1); + }); + + it('runs a job parked twice only once', async () => { + vi.useFakeTimers({toFake: ['Date']}); + const start = new Date('2026-08-30T12:00:00.000Z'); + vi.setSystemTime(start); + const scheduledJobQueue = new KVScheduledJobQueueService(new MockKVProvider()); + const runner = createRunner(scheduledJobQueue); + const runAt = new Date(start.getTime() + 60 * 60 * 1000); + const payload = {scheduledMessageId: '77', __jobId: '4242'}; + + await runner.runJob(TASK_TYPE, createJobMessage({runAt, payload, seq: 1}) as unknown as JsMsg); + await runner.runJob(TASK_TYPE, createJobMessage({runAt, payload, seq: 2, deliveryCount: 2}) as unknown as JsMsg); + + expect(await scheduledJobQueue.getQueueSize()).toBe(1); + + vi.setSystemTime(new Date(runAt.getTime() - LANE_ACK_WAIT_MS + 1000)); + const workerService = new RecordingWorkerService(); + await drain(scheduledJobQueue, workerService); + + expect(workerService.jobs).toHaveLength(1); + expect(await scheduledJobQueue.getQueueSize()).toBe(0); + }); + + it('parks a nearly delivery-exhausted job so it gets a fresh delivery budget', async () => { + const scheduledJobQueue = new KVScheduledJobQueueService(new MockKVProvider()); + const runner = createRunner(scheduledJobQueue); + const runAt = new Date(Date.now() + 20000); + const msg = createJobMessage({ + runAt, + payload: {scheduledMessageId: '77', __jobId: '4242'}, + deliveryCount: LANE_MAX_DELIVER - 1, + }); + + await runner.runJob(TASK_TYPE, msg as unknown as JsMsg); + + expect(msg.ack).toHaveBeenCalledTimes(1); + expect(msg.nak).not.toHaveBeenCalled(); + expect(queueStub.enqueue).not.toHaveBeenCalled(); + const workerService = new RecordingWorkerService(); + await drain(scheduledJobQueue, workerService); + expect(workerService.jobs).toHaveLength(1); + expect(workerService.jobs[0]!.options?.runAt?.getTime()).toBe(runAt.getTime()); + }); + + it('does nothing when the due queue is empty', async () => { + const scheduledJobQueue = new KVScheduledJobQueueService(new MockKVProvider()); + const workerService = new RecordingWorkerService(); + + await expect(drain(scheduledJobQueue, workerService)).resolves.toBeUndefined(); + + expect(workerService.jobs).toHaveLength(0); + expect(await scheduledJobQueue.getQueueSize()).toBe(0); + }); + + it('does not release a job another drain claimed first', async () => { + const scheduledJobQueue = new KVScheduledJobQueueService(new ConcurrentlyClaimedKVProvider()); + const runner = createRunner(scheduledJobQueue); + const runAt = new Date(Date.now() + 20000); + const msg = createJobMessage({ + runAt, + payload: {scheduledMessageId: '77', __jobId: '4242'}, + deliveryCount: LANE_MAX_DELIVER - 1, + }); + await runner.runJob(TASK_TYPE, msg as unknown as JsMsg); + expect(await scheduledJobQueue.getReadyJobs(Date.now(), 10)).toHaveLength(1); + + const workerService = new RecordingWorkerService(); + await drain(scheduledJobQueue, workerService); + + expect(workerService.jobs).toHaveLength(0); + expect(await scheduledJobQueue.getQueueSize()).toBe(0); + }); +}); diff --git a/fluxer_api/src/api/worker/tests/WorkerLaneConfig.test.ts b/fluxer_api/src/api/worker/tests/WorkerLaneConfig.test.ts index 661912944..71dd5a663 100644 --- a/fluxer_api/src/api/worker/tests/WorkerLaneConfig.test.ts +++ b/fluxer_api/src/api/worker/tests/WorkerLaneConfig.test.ts @@ -1,7 +1,12 @@ // SPDX-License-Identifier: AGPL-3.0-or-later import {describe, expect, it} from 'vitest'; -import {resolveCronSchedulerEnabled, resolveWorkerLanes, WORKER_LANES} from '../WorkerLaneConfig'; +import { + resolveCronSchedulerEnabled, + resolveScheduledJobDrainIntervalSeconds, + resolveWorkerLanes, + WORKER_LANES, +} from '../WorkerLaneConfig'; describe('WorkerLaneConfig', () => { it('returns all lanes in all_lanes mode', () => { @@ -97,4 +102,10 @@ describe('WorkerLaneConfig', () => { expect(resolveCronSchedulerEnabled('single_lane', true)).toBe(true); expect(resolveCronSchedulerEnabled('all_lanes', false)).toBe(false); }); + it('drains parked scheduled jobs well inside the shortest lane ack-wait', () => { + const intervalSeconds = resolveScheduledJobDrainIntervalSeconds(); + const shortestAckWaitMs = Math.min(...WORKER_LANES.map((lane) => lane.ackWaitMs)); + expect(intervalSeconds).toBeGreaterThanOrEqual(1); + expect(intervalSeconds * 1000 * 2).toBeLessThanOrEqual(shortestAckWaitMs); + }); });