mirror of
https://github.com/fluxerapp/fluxer
synced 2026-10-07 19:22:14 +09:00
fix(worker): catch up cron jobs missed by a delayed tick (#2275)
This commit is contained in:
@@ -6,6 +6,8 @@ import type {WorkerJobPayload} from '@pkgs/worker/src/contracts/WorkerTypes';
|
||||
import type {WorkerTaskName} from './WorkerLaneConfig';
|
||||
import type {WorkerService} from './WorkerService';
|
||||
|
||||
const MAX_CATCHUP_SECONDS = 60;
|
||||
|
||||
interface CronDefinition {
|
||||
id: string;
|
||||
taskType: WorkerTaskName;
|
||||
@@ -77,12 +79,22 @@ function matchesCronExpression(expression: string, date: Date): boolean {
|
||||
);
|
||||
}
|
||||
|
||||
function findLatestDueSecond(expression: string, fromSeconds: number, toSeconds: number): number | null {
|
||||
for (let second = toSeconds; second >= fromSeconds; second--) {
|
||||
if (matchesCronExpression(expression, new Date(second * 1000))) {
|
||||
return second;
|
||||
}
|
||||
}
|
||||
return null;
|
||||
}
|
||||
|
||||
export class CronScheduler {
|
||||
private readonly workerService: WorkerService;
|
||||
private readonly logger: LoggerInterface;
|
||||
private readonly kvClient: IKVProvider | null;
|
||||
private readonly definitions: Map<string, CronDefinition> = new Map();
|
||||
private intervalId: NodeJS.Timeout | null = null;
|
||||
private lastTickSecond: number | null = null;
|
||||
|
||||
constructor(workerService: WorkerService, logger: LoggerInterface, kvClient: IKVProvider | null = null) {
|
||||
this.workerService = workerService;
|
||||
@@ -124,28 +136,33 @@ export class CronScheduler {
|
||||
clearInterval(this.intervalId);
|
||||
this.intervalId = null;
|
||||
}
|
||||
this.lastTickSecond = null;
|
||||
}
|
||||
|
||||
private async tick(): Promise<void> {
|
||||
const now = new Date();
|
||||
const nowSeconds = Math.floor(now.getTime() / 1000);
|
||||
const nowSeconds = Math.floor(Date.now() / 1000);
|
||||
const previousTickSecond = this.lastTickSecond;
|
||||
this.lastTickSecond = nowSeconds;
|
||||
const fromSeconds =
|
||||
previousTickSecond === null || previousTickSecond >= nowSeconds
|
||||
? nowSeconds
|
||||
: Math.max(previousTickSecond + 1, nowSeconds - MAX_CATCHUP_SECONDS);
|
||||
for (const def of this.definitions.values()) {
|
||||
if (def.lastFired === nowSeconds) {
|
||||
const dueSecond = findLatestDueSecond(def.cronExpression, fromSeconds, nowSeconds);
|
||||
if (dueSecond === null || def.lastFired === dueSecond) {
|
||||
continue;
|
||||
}
|
||||
if (matchesCronExpression(def.cronExpression, now)) {
|
||||
def.lastFired = nowSeconds;
|
||||
try {
|
||||
const jobKey = `cron:${def.id}:${nowSeconds}`;
|
||||
const acquired = await this.acquireEnqueueLease(jobKey);
|
||||
if (!acquired) {
|
||||
continue;
|
||||
}
|
||||
await this.workerService.addJob(def.taskType, def.payload, {jobKey, skipLedger: !def.ledger});
|
||||
this.logger.debug({cronId: def.id, taskType: def.taskType}, 'Cron job fired');
|
||||
} catch (error) {
|
||||
this.logger.error({err: error, cronId: def.id, taskType: def.taskType}, 'Failed to enqueue cron job');
|
||||
def.lastFired = dueSecond;
|
||||
try {
|
||||
const jobKey = `cron:${def.id}:${dueSecond}`;
|
||||
const acquired = await this.acquireEnqueueLease(jobKey);
|
||||
if (!acquired) {
|
||||
continue;
|
||||
}
|
||||
await this.workerService.addJob(def.taskType, def.payload, {jobKey, skipLedger: !def.ledger});
|
||||
this.logger.debug({cronId: def.id, taskType: def.taskType}, 'Cron job fired');
|
||||
} catch (error) {
|
||||
this.logger.error({err: error, cronId: def.id, taskType: def.taskType}, 'Failed to enqueue cron job');
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,94 @@
|
||||
// SPDX-License-Identifier: AGPL-3.0-or-later
|
||||
|
||||
import type {LoggerInterface} from '@fluxer/logger/src/LoggerInterface';
|
||||
import type {IKVProvider} from '@pkgs/kv_client/src/IKVProvider';
|
||||
import {afterEach, describe, expect, it, vi} from 'vitest';
|
||||
import {CronScheduler} from '../CronScheduler';
|
||||
import type {WorkerService} from '../WorkerService';
|
||||
|
||||
function createLogger(): LoggerInterface {
|
||||
const logger = {
|
||||
trace: vi.fn(),
|
||||
debug: vi.fn(),
|
||||
info: vi.fn(),
|
||||
warn: vi.fn(),
|
||||
error: vi.fn(),
|
||||
child: () => logger,
|
||||
};
|
||||
return logger as unknown as LoggerInterface;
|
||||
}
|
||||
|
||||
function createScheduler(): {
|
||||
scheduler: CronScheduler;
|
||||
addJob: ReturnType<typeof vi.fn>;
|
||||
setnx: ReturnType<typeof vi.fn>;
|
||||
} {
|
||||
const addJob = vi.fn().mockResolvedValue(1n);
|
||||
const setnx = vi.fn().mockResolvedValue(true);
|
||||
const workerService = {addJob} as unknown as WorkerService;
|
||||
const kvClient = {setnx} as unknown as IKVProvider;
|
||||
return {scheduler: new CronScheduler(workerService, createLogger(), kvClient), addJob, setnx};
|
||||
}
|
||||
|
||||
describe('CronScheduler', () => {
|
||||
afterEach(() => {
|
||||
vi.useRealTimers();
|
||||
});
|
||||
|
||||
it('fires a schedule whose second was skipped by a stalled tick', async () => {
|
||||
vi.useFakeTimers();
|
||||
vi.setSystemTime(new Date('2026-01-01T00:00:00.000Z'));
|
||||
const {scheduler, addJob, setnx} = createScheduler();
|
||||
scheduler.upsert('expireAttachments', 'expireAttachments', {}, '5 * * * * *', {ledger: false});
|
||||
|
||||
scheduler.start();
|
||||
await vi.advanceTimersByTimeAsync(3000);
|
||||
expect(addJob).not.toHaveBeenCalled();
|
||||
|
||||
vi.setSystemTime(new Date('2026-01-01T00:00:10.400Z'));
|
||||
await vi.advanceTimersByTimeAsync(2000);
|
||||
scheduler.stop();
|
||||
|
||||
const skippedSecond = Math.floor(Date.parse('2026-01-01T00:00:05.000Z') / 1000);
|
||||
expect(addJob).toHaveBeenCalledTimes(1);
|
||||
expect(addJob).toHaveBeenCalledWith(
|
||||
'expireAttachments',
|
||||
{},
|
||||
{
|
||||
jobKey: `cron:expireAttachments:${skippedSecond}`,
|
||||
skipLedger: true,
|
||||
},
|
||||
);
|
||||
expect(setnx).toHaveBeenCalledTimes(1);
|
||||
});
|
||||
|
||||
it('fires a skipped schedule once instead of once per skipped second', async () => {
|
||||
vi.useFakeTimers();
|
||||
vi.setSystemTime(new Date('2026-01-01T00:00:00.000Z'));
|
||||
const {scheduler, addJob} = createScheduler();
|
||||
scheduler.upsert('flushUserActivityBuffer', 'flushUserActivityBuffer', {}, '*/2 * * * * *', {ledger: false});
|
||||
|
||||
scheduler.start();
|
||||
await vi.advanceTimersByTimeAsync(1000);
|
||||
addJob.mockClear();
|
||||
|
||||
vi.setSystemTime(new Date('2026-01-01T00:00:20.400Z'));
|
||||
await vi.advanceTimersByTimeAsync(1000);
|
||||
scheduler.stop();
|
||||
|
||||
expect(addJob).toHaveBeenCalledTimes(1);
|
||||
});
|
||||
|
||||
it('does not replay schedules from before the scheduler started', async () => {
|
||||
vi.useFakeTimers();
|
||||
vi.setSystemTime(new Date('2026-01-01T00:00:30.000Z'));
|
||||
const {scheduler, addJob} = createScheduler();
|
||||
scheduler.upsert('expireAttachments', 'expireAttachments', {}, '5 * * * * *', {ledger: false});
|
||||
|
||||
scheduler.start();
|
||||
await vi.advanceTimersByTimeAsync(2000);
|
||||
scheduler.stop();
|
||||
|
||||
expect(addJob).not.toHaveBeenCalled();
|
||||
});
|
||||
});
|
||||
Reference in New Issue
Block a user