diff --git a/packages/pgconductor-js/src/database-client.ts b/packages/pgconductor-js/src/database-client.ts index c388fa2..d67b74a 100644 --- a/packages/pgconductor-js/src/database-client.ts +++ b/packages/pgconductor-js/src/database-client.ts @@ -325,17 +325,20 @@ export class DatabaseClient { return new Date(); } - // In tests, query DB to respect fake_now + return this.getDatabaseTime(options); + } + + async getDatabaseTime(options?: QueryMethodOptions): Promise { const result = await this.query( (sql) => sql<{ now: Date }[]>` select pgconductor._private_current_time() as now `, - { label: "getCurrentTime", ...options }, + { label: "getDatabaseTime", ...options }, ); const row = result[0]; if (!row) { - throw new Error("getCurrentTime returned no rows"); + throw new Error("getDatabaseTime returned no rows"); } return row.now; } diff --git a/packages/pgconductor-js/src/lib/clock.ts b/packages/pgconductor-js/src/lib/clock.ts new file mode 100644 index 0000000..1c63402 --- /dev/null +++ b/packages/pgconductor-js/src/lib/clock.ts @@ -0,0 +1,80 @@ +import type { Logger } from "./logger"; + +const DATABASE_CLOCK_REFRESH_INTERVAL_MS = 10 * 60 * 1000; + +export type ClockOptions = { + sampleDatabaseTime: (signal?: AbortSignal) => Promise; + logger: Logger; + localClock?: () => Date; + refreshIntervalMs?: number; +}; + +/** + * A worker-local clock corrected to the database server's clock. + * The offset is database time minus the midpoint of the local request times. + */ +export class Clock { + private offsetMs = 0; + private refreshTimer: ReturnType | null = null; + private running = false; + + constructor({ + sampleDatabaseTime, + logger, + localClock = () => new Date(), + refreshIntervalMs = DATABASE_CLOCK_REFRESH_INTERVAL_MS, + }: ClockOptions) { + this.sampleDatabaseTime = sampleDatabaseTime; + this.logger = logger; + this.localClock = localClock; + this.refreshIntervalMs = refreshIntervalMs; + } + + private readonly sampleDatabaseTime: (signal?: AbortSignal) => Promise; + private readonly logger: Logger; + private readonly localClock: () => Date; + private readonly refreshIntervalMs: number; + + now(): Date { + return new Date(this.localClock().getTime() + this.offsetMs); + } + + get offset(): number { + return this.offsetMs; + } + + async start(signal?: AbortSignal): Promise { + this.stop(); + this.running = true; + await this.refresh(signal); + if (this.running) this.scheduleNextRefresh(); + } + + stop(): void { + this.running = false; + if (this.refreshTimer !== null) { + clearTimeout(this.refreshTimer); + this.refreshTimer = null; + } + } + + async refresh(signal?: AbortSignal): Promise { + try { + const requestStart = this.localClock(); + const databaseTime = await this.sampleDatabaseTime(signal); + const responseEnd = this.localClock(); + const midpoint = (requestStart.getTime() + responseEnd.getTime()) / 2; + this.offsetMs = databaseTime.getTime() - midpoint; + } catch (error) { + this.logger.warn("Database clock offset refresh failed; retaining previous offset", error); + } + } + + private scheduleNextRefresh(): void { + this.refreshTimer = setTimeout(async () => { + this.refreshTimer = null; + await this.refresh(); + if (this.running) this.scheduleNextRefresh(); + }, this.refreshIntervalMs); + } +} diff --git a/packages/pgconductor-js/src/lib/cron.ts b/packages/pgconductor-js/src/lib/cron.ts new file mode 100644 index 0000000..8beedbe --- /dev/null +++ b/packages/pgconductor-js/src/lib/cron.ts @@ -0,0 +1,5 @@ +import CronExpressionParser from "cron-parser"; + +export function nextCronOccurrence(expression: string, now: Date): Date { + return CronExpressionParser.parse(expression, { currentDate: now }).next().toDate(); +} diff --git a/packages/pgconductor-js/src/task-context.ts b/packages/pgconductor-js/src/task-context.ts index 509febe..184940c 100644 --- a/packages/pgconductor-js/src/task-context.ts +++ b/packages/pgconductor-js/src/task-context.ts @@ -1,5 +1,6 @@ -import CronExpressionParser from "cron-parser"; import type { DatabaseClient, JsonValue, Execution, Payload } from "./database-client"; +import { nextCronOccurrence } from "./lib/cron"; +import type { Clock } from "./lib/clock"; import type { TaskDefinition, TaskName, @@ -76,6 +77,7 @@ export function createTaskSignal( export type TaskContextOptions = { abortController: TypedAbortController; db: DatabaseClient; + clock: Clock; execution: Execution; logger: Logger; window?: [string, string]; @@ -295,8 +297,7 @@ export class TaskContext< throw new Error("cron expression is required"); } - const interval = CronExpressionParser.parse(options.cron); - const nextTimestamp = interval.next().toDate(); + const nextTimestamp = nextCronOccurrence(options.cron, this.opts.clock.now()); const queue = task.queue || "default"; await this.opts.db.scheduleCronExecution( diff --git a/packages/pgconductor-js/src/worker.ts b/packages/pgconductor-js/src/worker.ts index d18a18d..6337d45 100644 --- a/packages/pgconductor-js/src/worker.ts +++ b/packages/pgconductor-js/src/worker.ts @@ -17,7 +17,8 @@ import { mapConcurrent } from "./lib/map-concurrent"; import { Deferred } from "./lib/deferred"; import { type PollableAsyncIterable } from "./lib/async-queue"; import { BatchingAsyncQueue, type BatchGroup } from "./lib/batching-async-queue"; -import CronExpressionParser from "cron-parser"; +import { nextCronOccurrence } from "./lib/cron"; +import { Clock } from "./lib/clock"; import { createTaskSignal, isTaskAbortReason, @@ -131,6 +132,7 @@ export class Worker< private readonly flushBatchSize: number; private readonly flushIntervalMs: number; private readonly pollIntervalMs: number; + private readonly clock: Clock; private _startDeferred: Deferred | null = null; private _stopDeferred: Deferred | null = null; @@ -161,6 +163,10 @@ export class Worker< this.flushIntervalMs = fullConfig.flushIntervalMs; this.fetchBatchSize = fullConfig.fetchBatchSize; this.flushBatchSize = fullConfig.flushBatchSize; + this.clock = new Clock({ + sampleDatabaseTime: (signal) => this.db.getDatabaseTime({ signal }), + logger: this.logger, + }); } /** @@ -221,7 +227,8 @@ export class Worker< this._stopDeferred = new Deferred(); this._abortController = new AbortController(); - // Synchronous registration + // Sample before calculating or registering cron schedules. + await this.clock.start(this.abortController.signal); await this.register(); // Worker is now started @@ -247,6 +254,7 @@ export class Worker< this.logger.error("Worker pipeline error:", err); } finally { queue.close(); + this.clock.stop(); this._stopDeferred?.resolve(); this._startDeferred = null; this._stopDeferred = null; @@ -317,8 +325,7 @@ export class Worker< task.triggers .filter((t): t is { cron: string; name: string } => "cron" in t) .map((trigger) => { - const interval = CronExpressionParser.parse(trigger.cron); - const nextTimestamp = interval.next().toDate(); + const nextTimestamp = nextCronOccurrence(trigger.cron, this.clock.now()); const timestampSeconds = Math.floor(nextTimestamp.getTime() / 1000); return { task_key: task.name, @@ -585,6 +592,7 @@ export class Worker< TaskContext.create( { db: this.db, + clock: this.clock, abortController: taskAbortController, execution: exec, logger: makeChildLogger(this.logger, { @@ -807,8 +815,7 @@ export class Worker< } const scheduleName = parts[1]; - const interval = CronExpressionParser.parse(execution.cron_expression); - const nextTimestamp = interval.next().toDate(); + const nextTimestamp = nextCronOccurrence(execution.cron_expression, this.clock.now()); const timestampSeconds = Math.floor(nextTimestamp.getTime() / 1000); const nextDedupeKey = `scheduled::${scheduleName}::${timestampSeconds}`; diff --git a/packages/pgconductor-js/tests/mocks/database-client.mock.ts b/packages/pgconductor-js/tests/mocks/database-client.mock.ts index db8b679..b799b5f 100644 --- a/packages/pgconductor-js/tests/mocks/database-client.mock.ts +++ b/packages/pgconductor-js/tests/mocks/database-client.mock.ts @@ -30,6 +30,7 @@ export class MockDatabaseClient implements IDatabaseClient { clearWaitingState = mock(async () => {}); cancelExecution = mock(async () => true); getCurrentTime = mock(async () => new Date()); + getDatabaseTime = mock(async () => new Date()); setFakeTime = mock(async () => {}); clearFakeTime = mock(async () => {}); subscribeEvent = mock(async () => "mock-subscription-id"); diff --git a/packages/pgconductor-js/tests/mocks/in-memory-database-client.ts b/packages/pgconductor-js/tests/mocks/in-memory-database-client.ts index 7170e10..f58de28 100644 --- a/packages/pgconductor-js/tests/mocks/in-memory-database-client.ts +++ b/packages/pgconductor-js/tests/mocks/in-memory-database-client.ts @@ -158,11 +158,15 @@ export class InMemoryDatabaseClient implements IDatabaseClient { return this.getInternalTime(); } - // Async method matching DatabaseClient interface + // Async methods matching DatabaseClient interface async getCurrentTime(): Promise { return this.getInternalTime(); } + async getDatabaseTime(): Promise { + return this.getInternalTime(); + } + // ============================================================================ // Orchestrator Management // ============================================================================ diff --git a/packages/pgconductor-js/tests/unit/clock.test.ts b/packages/pgconductor-js/tests/unit/clock.test.ts new file mode 100644 index 0000000..03979d1 --- /dev/null +++ b/packages/pgconductor-js/tests/unit/clock.test.ts @@ -0,0 +1,149 @@ +import { beforeEach, describe, expect, mock, test } from "bun:test"; +import { Clock } from "../../src/lib/clock"; +import { nextCronOccurrence } from "../../src/lib/cron"; + +const logger = { + info: mock(), + warn: mock(), + error: mock(), + debug: mock(), +}; + +const date = (ms: number): Date => new Date(ms); + +describe("Clock", () => { + beforeEach(() => { + logger.warn.mockClear(); + }); + + test("corrects positive and negative local clock skew", async () => { + let now = 1_000; + const positive = new Clock({ + sampleDatabaseTime: async () => date(2_000), + logger, + localClock: () => date(now), + }); + + await positive.refresh(); + now = 1_100; + expect(positive.now()).toEqual(date(2_100)); + + now = 1_000; + const negative = new Clock({ + sampleDatabaseTime: async () => date(0), + logger, + localClock: () => date(now), + }); + + await negative.refresh(); + now = 1_100; + expect(negative.now()).toEqual(date(100)); + }); + + test("uses the request midpoint to account for latency", async () => { + const localTimes = [date(1_000), date(1_300)]; + const clock = new Clock({ + sampleDatabaseTime: async () => date(1_100), + logger, + localClock: () => { + const value = localTimes.shift(); + if (!value) throw new Error("local clock exhausted"); + return value; + }, + }); + + await clock.refresh(); + expect(clock.offset).toBe(-50); + }); + + test("gives skewed workers the same cron slot", async () => { + const databaseNow = date(Date.UTC(2025, 0, 1, 12, 0, 0)); + const positive = new Clock({ + sampleDatabaseTime: async () => databaseNow, + logger, + localClock: () => date(databaseNow.getTime() - 5 * 60 * 1000), + }); + const negative = new Clock({ + sampleDatabaseTime: async () => databaseNow, + logger, + localClock: () => date(databaseNow.getTime() + 5 * 60 * 1000), + }); + + await positive.refresh(); + await negative.refresh(); + + expect(nextCronOccurrence("0 */5 * * * *", positive.now())).toEqual( + nextCronOccurrence("0 */5 * * * *", negative.now()), + ); + }); + + test("retains the local clock when the initial sample fails", async () => { + const clock = new Clock({ + sampleDatabaseTime: async () => { + throw new Error("database unavailable"); + }, + logger, + localClock: () => date(1_000), + }); + + await clock.start(); + clock.stop(); + + expect(clock.now()).toEqual(date(1_000)); + expect(logger.warn).toHaveBeenCalledWith( + "Database clock offset refresh failed; retaining previous offset", + expect.any(Error), + ); + }); + + test("retains the previous offset when a refresh fails", async () => { + let fail = false; + const clock = new Clock({ + sampleDatabaseTime: async () => { + if (fail) throw new Error("database unavailable"); + return date(2_000); + }, + logger, + localClock: () => date(1_000), + }); + + await clock.refresh(); + fail = true; + await clock.refresh(); + + expect(clock.offset).toBe(1_000); + expect(logger.warn).toHaveBeenCalledWith( + "Database clock offset refresh failed; retaining previous offset", + expect.any(Error), + ); + }); + + test("refreshes periodically without overlapping samples and stops its timer", async () => { + let samples = 0; + let inFlight = 0; + let maxInFlight = 0; + const clock = new Clock({ + sampleDatabaseTime: async () => { + samples++; + inFlight++; + maxInFlight = Math.max(maxInFlight, inFlight); + await new Promise((resolve) => setTimeout(resolve, 10)); + inFlight--; + return date(1_000 + samples); + }, + logger, + localClock: () => date(0), + refreshIntervalMs: 1, + }); + + await clock.start(); + await new Promise((resolve) => setTimeout(resolve, 25)); + clock.stop(); + const samplesAfterStop = samples; + await new Promise((resolve) => setTimeout(resolve, 25)); + + expect(samples).toBeGreaterThan(1); + expect(maxInFlight).toBe(1); + expect(samples).toBe(samplesAfterStop); + }); +});