mirror of
https://github.com/fluxerapp/fluxer
synced 2026-10-07 19:22:14 +09:00
fix(worker): park far-future jobs in a KV due queue (#2144)
This commit is contained in:
@@ -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<string, unknown>;
|
||||
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<void> {
|
||||
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<boolean> {
|
||||
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<Array<ParkedScheduledJob>> {
|
||||
try {
|
||||
const results = await this.kvClient.zrangebyscore(QUEUE_KEY, '-inf', nowMs, 'LIMIT', 0, limit);
|
||||
const jobs: Array<ParkedScheduledJob> = [];
|
||||
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<number> {
|
||||
try {
|
||||
return await this.kvClient.zcard(QUEUE_KEY);
|
||||
} catch (error) {
|
||||
Logger.error({error}, 'Failed to get parked scheduled job queue size');
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -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;
|
||||
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -72,6 +72,7 @@ const LANE_CONFIG = {
|
||||
'processInactivityDeletions',
|
||||
'processPendingBulkMessageDeletions',
|
||||
'processPremiumStateReconciliationQueue',
|
||||
'processScheduledJobQueue',
|
||||
'prunePostgresKvTtl',
|
||||
'refreshSearchIndex',
|
||||
'syncDiscoveryIndex',
|
||||
@@ -213,6 +214,13 @@ function resolveWorkerLanes(config: WorkerLaneRuntimeConfig): Array<WorkerLaneDe
|
||||
});
|
||||
}
|
||||
|
||||
const SCHEDULED_JOB_DRAIN_INTERVALS_PER_ACK_WAIT = 2;
|
||||
|
||||
function resolveScheduledJobDrainIntervalSeconds(): number {
|
||||
const minimumAckWaitMs = Math.min(...WORKER_LANES.map((lane) => 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,
|
||||
};
|
||||
|
||||
@@ -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<void> {
|
||||
const runner = new WorkerRunner({
|
||||
tasks: laneTasks,
|
||||
queue,
|
||||
scheduledJobQueue: dependencies.scheduledJobQueueService,
|
||||
consumerName: lane.consumerName,
|
||||
laneName: lane.name,
|
||||
ledger: jobLedger,
|
||||
|
||||
@@ -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<string>;
|
||||
publishToDlq(taskType: string, originalPayload: Record<string, unknown>, meta: WorkerRunnerDlqMeta): Promise<void>;
|
||||
}
|
||||
|
||||
interface WorkerRunnerScheduledJobQueue {
|
||||
parkJob(job: ParkedScheduledJob, releaseAt: Date): Promise<void>;
|
||||
}
|
||||
|
||||
interface WorkerRunnerOptions {
|
||||
tasks: Record<string, WorkerTaskHandler>;
|
||||
queue: WorkerRunnerQueue;
|
||||
scheduledJobQueue: WorkerRunnerScheduledJobQueue;
|
||||
consumerName: string;
|
||||
laneName: string;
|
||||
ledger: IJobLedgerRepository;
|
||||
@@ -66,6 +61,7 @@ interface WorkerRunnerOptions {
|
||||
export class WorkerRunner {
|
||||
private readonly tasks: Record<string, WorkerTaskHandler>;
|
||||
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 {
|
||||
|
||||
@@ -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<WorkerTaskName, WorkerTaskHandler> = {
|
||||
processInactivityDeletions,
|
||||
processPendingBulkMessageDeletions,
|
||||
processPremiumStateReconciliationQueue,
|
||||
processScheduledJobQueue,
|
||||
reconcileUserPayments,
|
||||
prunePostgresKvTtl,
|
||||
refreshSearchIndex,
|
||||
|
||||
@@ -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;
|
||||
@@ -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<RecordedJob> = [];
|
||||
|
||||
async addJob<TPayload extends WorkerJobPayload = WorkerJobPayload>(
|
||||
taskType: string,
|
||||
payload: TPayload,
|
||||
options?: WorkerJobOptions,
|
||||
): Promise<bigint> {
|
||||
this.jobs.push({taskType, payload, options});
|
||||
return BigInt(this.jobs.length);
|
||||
}
|
||||
|
||||
async cancelJob(_jobId: bigint): Promise<boolean> {
|
||||
return false;
|
||||
}
|
||||
|
||||
async retryDeadLetterJob(_jobId: bigint): Promise<boolean> {
|
||||
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<boolean> {
|
||||
await super.removeBulkDeletion(queueKey, secondaryKey);
|
||||
return false;
|
||||
}
|
||||
}
|
||||
|
||||
class TestWorkerRunner extends WorkerRunner {
|
||||
async runJob(taskType: string, msg: JsMsg): Promise<boolean> {
|
||||
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<string, unknown>;
|
||||
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);
|
||||
});
|
||||
});
|
||||
@@ -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);
|
||||
});
|
||||
});
|
||||
|
||||
Reference in New Issue
Block a user