From 648ee34757bc94dd3a1b9ede5c864942dc3058a8 Mon Sep 17 00:00:00 2001 From: psteinroe Date: Mon, 14 Sep 2026 22:24:29 +0000 Subject: [PATCH 1/5] fix(cron): align scheduling with database time --- .../pgconductor-js/src/database-client.ts | 9 +- packages/pgconductor-js/src/lib/clock-skew.ts | 82 ++++++++++ packages/pgconductor-js/src/lib/cron.ts | 5 + packages/pgconductor-js/src/orchestrator.ts | 14 +- packages/pgconductor-js/src/task-context.ts | 7 +- packages/pgconductor-js/src/worker.ts | 42 ++++-- .../integration/orchestrator-shutdown.test.ts | 34 ++++- .../tests/mocks/database-client.mock.ts | 1 + .../tests/mocks/in-memory-database-client.ts | 6 +- .../tests/unit/clock-skew.test.ts | 142 ++++++++++++++++++ 10 files changed, 315 insertions(+), 27 deletions(-) create mode 100644 packages/pgconductor-js/src/lib/clock-skew.ts create mode 100644 packages/pgconductor-js/src/lib/cron.ts create mode 100644 packages/pgconductor-js/tests/unit/clock-skew.test.ts 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-skew.ts b/packages/pgconductor-js/src/lib/clock-skew.ts new file mode 100644 index 0000000..064e7e1 --- /dev/null +++ b/packages/pgconductor-js/src/lib/clock-skew.ts @@ -0,0 +1,82 @@ +import type { Logger } from "./logger"; + +export type LocalClock = () => Date; +export type DatabaseTimeSampler = (signal?: AbortSignal) => Promise; + +export type WorkerClock = { + now(): Date; + start(signal?: AbortSignal): Promise; + stop(): void; +}; + +export const DATABASE_CLOCK_REFRESH_INTERVAL_MS = 10 * 60 * 1000; + +/** + * 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 DatabaseClockOffset implements WorkerClock { + private offsetMs = 0; + private refreshTimer: ReturnType | null = null; + private running = false; + + constructor( + private readonly sampleDatabaseTime: DatabaseTimeSampler, + private readonly logger: Logger, + private readonly localClock: LocalClock = () => new Date(), + private readonly refreshIntervalMs = DATABASE_CLOCK_REFRESH_INTERVAL_MS, + ) {} + + 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; + try { + await this.refreshOffset(signal, true); + if (this.running) this.scheduleRefresh(); + } catch (error) { + this.stop(); + throw error; + } + } + + stop(): void { + this.running = false; + if (this.refreshTimer !== null) { + clearTimeout(this.refreshTimer); + this.refreshTimer = null; + } + } + + async refresh(signal?: AbortSignal): Promise { + await this.refreshOffset(signal, false); + } + + private async refreshOffset(signal: AbortSignal | undefined, initial: boolean): 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) { + if (initial) throw error; + this.logger.warn("Database clock offset refresh failed; retaining previous offset", error); + } + } + + private scheduleRefresh(): void { + this.refreshTimer = setTimeout(async () => { + this.refreshTimer = null; + await this.refresh(); + if (this.running) this.scheduleRefresh(); + }, 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/orchestrator.ts b/packages/pgconductor-js/src/orchestrator.ts index 8f68a64..064a25e 100644 --- a/packages/pgconductor-js/src/orchestrator.ts +++ b/packages/pgconductor-js/src/orchestrator.ts @@ -189,13 +189,13 @@ export class Orchestrator { // Start heartbeat loop this.startHeartbeatLoop(); - // Kick off all workers (don't await yet!) - if (runOnce) { - // Drain mode: workers will process and stop - this.workers.forEach((w) => void w.drain(this.orchestratorId)); - } else { - // Normal mode: workers will run continuously - this.workers.forEach((w) => void w.run(this.orchestratorId)); + // Startup is observed through worker.started below; consume the run promise + // as well so the same startup error is not reported as unhandled. + for (const worker of this.workers) { + const running = runOnce + ? worker.drain(this.orchestratorId) + : worker.run(this.orchestratorId); + running.catch(noop); } // Wait for ALL workers to finish starting (register() complete) diff --git a/packages/pgconductor-js/src/task-context.ts b/packages/pgconductor-js/src/task-context.ts index 509febe..5c504d0 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 { WorkerClock } from "./lib/clock-skew"; import type { TaskDefinition, TaskName, @@ -76,6 +77,7 @@ export function createTaskSignal( export type TaskContextOptions = { abortController: TypedAbortController; db: DatabaseClient; + clock: WorkerClock; 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..d1e2f59 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 { DatabaseClockOffset, type WorkerClock } from "./lib/clock-skew"; import { createTaskSignal, isTaskAbortReason, @@ -37,6 +38,7 @@ import type { TypedAbortController } from "./lib/typed-abort-controller"; */ export type WorkerConfig = { concurrency: number; + clock?: WorkerClock; flushBatchSize: number; fetchBatchSize: number; pollIntervalMs: number; @@ -131,6 +133,7 @@ export class Worker< private readonly flushBatchSize: number; private readonly flushIntervalMs: number; private readonly pollIntervalMs: number; + private readonly clock: WorkerClock; private _startDeferred: Deferred | null = null; private _stopDeferred: Deferred | null = null; @@ -161,6 +164,9 @@ export class Worker< this.flushIntervalMs = fullConfig.flushIntervalMs; this.fetchBatchSize = fullConfig.fetchBatchSize; this.flushBatchSize = fullConfig.flushBatchSize; + this.clock = + fullConfig.clock || + new DatabaseClockOffset((signal) => this.db.getDatabaseTime({ signal }), this.logger); } /** @@ -217,15 +223,27 @@ export class Worker< } this.orchestratorId = orchestratorId; - this._startDeferred = new Deferred(); - this._stopDeferred = new Deferred(); - this._abortController = new AbortController(); + const started = (this._startDeferred = new Deferred()); + const stopped = (this._stopDeferred = new Deferred()); + const controller = (this._abortController = new AbortController()); - // Synchronous registration - await this.register(); + try { + // Sample before calculating or registering cron schedules. + await this.clock.start(controller.signal); + await this.register(); + } catch (error) { + this.clock.stop(); + controller.abort(); + started.reject(error); + stopped.resolve(); + this._startDeferred = null; + this._stopDeferred = null; + this._abortController = null; + throw error; + } // Worker is now started - this._startDeferred.resolve(); + started.resolve(); // Run pipeline in background // Build batch configs map @@ -247,6 +265,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; @@ -254,7 +273,7 @@ export class Worker< } })(); - return this._startDeferred.promise; + return started.promise; } /** @@ -317,8 +336,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 +603,7 @@ export class Worker< TaskContext.create( { db: this.db, + clock: this.clock, abortController: taskAbortController, execution: exec, logger: makeChildLogger(this.logger, { @@ -807,8 +826,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/integration/orchestrator-shutdown.test.ts b/packages/pgconductor-js/tests/integration/orchestrator-shutdown.test.ts index 832fb85..200ae1d 100644 --- a/packages/pgconductor-js/tests/integration/orchestrator-shutdown.test.ts +++ b/packages/pgconductor-js/tests/integration/orchestrator-shutdown.test.ts @@ -1,8 +1,10 @@ -import { afterAll, afterEach, beforeAll, test, expect } from "bun:test"; +import { afterAll, afterEach, beforeAll, test, expect, mock } from "bun:test"; import { TestDatabasePool } from "../fixtures/test-database"; import type { TestDatabase } from "../fixtures/test-database"; import { Conductor } from "../../src/conductor"; import { Orchestrator } from "../../src/orchestrator"; +import { defineTask } from "../../src/task-definition"; +import { TaskSchemas } from "../../src/schemas"; let pool: TestDatabasePool; const databases: TestDatabase[] = []; @@ -43,6 +45,36 @@ test("gracefully shuts down on stop()", async () => { expect(orch.isStopped).toBe(true); }); +test("reports a database clock startup failure", async () => { + const db = await pool.child(); + databases.push(db); + const taskDefinition = defineTask({ name: "clock-startup" }); + const conductor = Conductor.create({ + sql: db.sql, + tasks: TaskSchemas.fromSchema([taskDefinition]), + context: {}, + }); + await conductor.ensureInstalled(); + + const clock = { + now: () => new Date(), + start: mock(async () => { + throw new Error("clock unavailable"); + }), + stop: mock(() => {}), + }; + const task = conductor.createTask({ name: "clock-startup" }, { invocable: true }, async () => {}); + const orchestrator = Orchestrator.create({ + conductor, + tasks: [task], + defaultWorker: { clock }, + }); + + await expect(orchestrator.start()).rejects.toThrow("clock unavailable"); + await orchestrator.stopped.catch(() => {}); + expect(clock.stop).toHaveBeenCalled(); +}, 30_000); + test("stop() can be called multiple times safely", async () => { const db = await pool.child(); databases.push(db); 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-skew.test.ts b/packages/pgconductor-js/tests/unit/clock-skew.test.ts new file mode 100644 index 0000000..f08c816 --- /dev/null +++ b/packages/pgconductor-js/tests/unit/clock-skew.test.ts @@ -0,0 +1,142 @@ +import { beforeEach, describe, expect, mock, test } from "bun:test"; +import { DatabaseClockOffset } from "../../src/lib/clock-skew"; +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("DatabaseClockOffset", () => { + beforeEach(() => { + logger.warn.mockClear(); + }); + + test("corrects positive and negative local clock skew", async () => { + let now = 1_000; + const positive = new DatabaseClockOffset( + async () => date(2_000), + logger, + () => date(now), + ); + + await positive.refresh(); + now = 1_100; + expect(positive.now()).toEqual(date(2_100)); + + now = 1_000; + const negative = new DatabaseClockOffset( + async () => date(0), + logger, + () => 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 DatabaseClockOffset( + async () => date(1_100), + logger, + () => { + 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 DatabaseClockOffset( + async () => databaseNow, + logger, + () => date(databaseNow.getTime() - 5 * 60 * 1000), + ); + const negative = new DatabaseClockOffset( + async () => databaseNow, + logger, + () => 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("fails startup when the initial sample fails", async () => { + const clock = new DatabaseClockOffset( + async () => { + throw new Error("database unavailable"); + }, + logger, + () => date(1_000), + ); + + await expect(clock.start()).rejects.toThrow("database unavailable"); + }); + + test("retains the previous offset when a refresh fails", async () => { + let fail = false; + const clock = new DatabaseClockOffset( + async () => { + if (fail) throw new Error("database unavailable"); + return date(2_000); + }, + logger, + () => 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 DatabaseClockOffset( + async () => { + samples++; + inFlight++; + maxInFlight = Math.max(maxInFlight, inFlight); + await new Promise((resolve) => setTimeout(resolve, 10)); + inFlight--; + return date(1_000 + samples); + }, + logger, + () => date(0), + 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); + }); +}); From acf128d772bc8a48773ff0015dbbeead12319d0a Mon Sep 17 00:00:00 2001 From: psteinroe Date: Thu, 17 Sep 2026 19:50:23 +0000 Subject: [PATCH 2/5] refactor(clock): keep sampling failures fail-open --- packages/pgconductor-js/src/lib/clock-skew.ts | 14 +++----- packages/pgconductor-js/src/worker.ts | 27 +++++---------- .../integration/orchestrator-shutdown.test.ts | 34 +------------------ .../tests/unit/clock-skew.test.ts | 11 ++++-- 4 files changed, 22 insertions(+), 64 deletions(-) diff --git a/packages/pgconductor-js/src/lib/clock-skew.ts b/packages/pgconductor-js/src/lib/clock-skew.ts index 064e7e1..d9bcf41 100644 --- a/packages/pgconductor-js/src/lib/clock-skew.ts +++ b/packages/pgconductor-js/src/lib/clock-skew.ts @@ -38,13 +38,8 @@ export class DatabaseClockOffset implements WorkerClock { async start(signal?: AbortSignal): Promise { this.stop(); this.running = true; - try { - await this.refreshOffset(signal, true); - if (this.running) this.scheduleRefresh(); - } catch (error) { - this.stop(); - throw error; - } + await this.refreshOffset(signal); + if (this.running) this.scheduleRefresh(); } stop(): void { @@ -56,10 +51,10 @@ export class DatabaseClockOffset implements WorkerClock { } async refresh(signal?: AbortSignal): Promise { - await this.refreshOffset(signal, false); + await this.refreshOffset(signal); } - private async refreshOffset(signal: AbortSignal | undefined, initial: boolean): Promise { + private async refreshOffset(signal?: AbortSignal): Promise { try { const requestStart = this.localClock(); const databaseTime = await this.sampleDatabaseTime(signal); @@ -67,7 +62,6 @@ export class DatabaseClockOffset implements WorkerClock { const midpoint = (requestStart.getTime() + responseEnd.getTime()) / 2; this.offsetMs = databaseTime.getTime() - midpoint; } catch (error) { - if (initial) throw error; this.logger.warn("Database clock offset refresh failed; retaining previous offset", error); } } diff --git a/packages/pgconductor-js/src/worker.ts b/packages/pgconductor-js/src/worker.ts index d1e2f59..b3f2854 100644 --- a/packages/pgconductor-js/src/worker.ts +++ b/packages/pgconductor-js/src/worker.ts @@ -223,27 +223,16 @@ export class Worker< } this.orchestratorId = orchestratorId; - const started = (this._startDeferred = new Deferred()); - const stopped = (this._stopDeferred = new Deferred()); - const controller = (this._abortController = new AbortController()); + this._startDeferred = new Deferred(); + this._stopDeferred = new Deferred(); + this._abortController = new AbortController(); - try { - // Sample before calculating or registering cron schedules. - await this.clock.start(controller.signal); - await this.register(); - } catch (error) { - this.clock.stop(); - controller.abort(); - started.reject(error); - stopped.resolve(); - this._startDeferred = null; - this._stopDeferred = null; - this._abortController = null; - throw error; - } + // Sample before calculating or registering cron schedules. + await this.clock.start(this.abortController.signal); + await this.register(); // Worker is now started - started.resolve(); + this._startDeferred.resolve(); // Run pipeline in background // Build batch configs map @@ -273,7 +262,7 @@ export class Worker< } })(); - return started.promise; + return this._startDeferred.promise; } /** diff --git a/packages/pgconductor-js/tests/integration/orchestrator-shutdown.test.ts b/packages/pgconductor-js/tests/integration/orchestrator-shutdown.test.ts index 200ae1d..832fb85 100644 --- a/packages/pgconductor-js/tests/integration/orchestrator-shutdown.test.ts +++ b/packages/pgconductor-js/tests/integration/orchestrator-shutdown.test.ts @@ -1,10 +1,8 @@ -import { afterAll, afterEach, beforeAll, test, expect, mock } from "bun:test"; +import { afterAll, afterEach, beforeAll, test, expect } from "bun:test"; import { TestDatabasePool } from "../fixtures/test-database"; import type { TestDatabase } from "../fixtures/test-database"; import { Conductor } from "../../src/conductor"; import { Orchestrator } from "../../src/orchestrator"; -import { defineTask } from "../../src/task-definition"; -import { TaskSchemas } from "../../src/schemas"; let pool: TestDatabasePool; const databases: TestDatabase[] = []; @@ -45,36 +43,6 @@ test("gracefully shuts down on stop()", async () => { expect(orch.isStopped).toBe(true); }); -test("reports a database clock startup failure", async () => { - const db = await pool.child(); - databases.push(db); - const taskDefinition = defineTask({ name: "clock-startup" }); - const conductor = Conductor.create({ - sql: db.sql, - tasks: TaskSchemas.fromSchema([taskDefinition]), - context: {}, - }); - await conductor.ensureInstalled(); - - const clock = { - now: () => new Date(), - start: mock(async () => { - throw new Error("clock unavailable"); - }), - stop: mock(() => {}), - }; - const task = conductor.createTask({ name: "clock-startup" }, { invocable: true }, async () => {}); - const orchestrator = Orchestrator.create({ - conductor, - tasks: [task], - defaultWorker: { clock }, - }); - - await expect(orchestrator.start()).rejects.toThrow("clock unavailable"); - await orchestrator.stopped.catch(() => {}); - expect(clock.stop).toHaveBeenCalled(); -}, 30_000); - test("stop() can be called multiple times safely", async () => { const db = await pool.child(); databases.push(db); diff --git a/packages/pgconductor-js/tests/unit/clock-skew.test.ts b/packages/pgconductor-js/tests/unit/clock-skew.test.ts index f08c816..0e2fdf8 100644 --- a/packages/pgconductor-js/tests/unit/clock-skew.test.ts +++ b/packages/pgconductor-js/tests/unit/clock-skew.test.ts @@ -77,7 +77,7 @@ describe("DatabaseClockOffset", () => { ); }); - test("fails startup when the initial sample fails", async () => { + test("retains the local clock when the initial sample fails", async () => { const clock = new DatabaseClockOffset( async () => { throw new Error("database unavailable"); @@ -86,7 +86,14 @@ describe("DatabaseClockOffset", () => { () => date(1_000), ); - await expect(clock.start()).rejects.toThrow("database unavailable"); + 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 () => { From 84614956274ba4745dbab0dcc0169e5e0d562099 Mon Sep 17 00:00:00 2001 From: psteinroe Date: Thu, 17 Sep 2026 19:52:00 +0000 Subject: [PATCH 3/5] refactor(clock): fold offset refresh helper --- packages/pgconductor-js/src/lib/clock-skew.ts | 12 ++++-------- 1 file changed, 4 insertions(+), 8 deletions(-) diff --git a/packages/pgconductor-js/src/lib/clock-skew.ts b/packages/pgconductor-js/src/lib/clock-skew.ts index d9bcf41..6e062a6 100644 --- a/packages/pgconductor-js/src/lib/clock-skew.ts +++ b/packages/pgconductor-js/src/lib/clock-skew.ts @@ -38,8 +38,8 @@ export class DatabaseClockOffset implements WorkerClock { async start(signal?: AbortSignal): Promise { this.stop(); this.running = true; - await this.refreshOffset(signal); - if (this.running) this.scheduleRefresh(); + await this.refresh(signal); + if (this.running) this.scheduleNextRefresh(); } stop(): void { @@ -51,10 +51,6 @@ export class DatabaseClockOffset implements WorkerClock { } async refresh(signal?: AbortSignal): Promise { - await this.refreshOffset(signal); - } - - private async refreshOffset(signal?: AbortSignal): Promise { try { const requestStart = this.localClock(); const databaseTime = await this.sampleDatabaseTime(signal); @@ -66,11 +62,11 @@ export class DatabaseClockOffset implements WorkerClock { } } - private scheduleRefresh(): void { + private scheduleNextRefresh(): void { this.refreshTimer = setTimeout(async () => { this.refreshTimer = null; await this.refresh(); - if (this.running) this.scheduleRefresh(); + if (this.running) this.scheduleNextRefresh(); }, this.refreshIntervalMs); } } From 42db4531b74a590a8fc80f1c71b8346c4c247484 Mon Sep 17 00:00:00 2001 From: psteinroe Date: Thu, 17 Sep 2026 19:54:26 +0000 Subject: [PATCH 4/5] refactor(clock): use concrete worker clock --- .../src/lib/{clock-skew.ts => clock.ts} | 17 ++++------------ packages/pgconductor-js/src/orchestrator.ts | 14 ++++++------- packages/pgconductor-js/src/task-context.ts | 4 ++-- packages/pgconductor-js/src/worker.ts | 9 +++------ .../{clock-skew.test.ts => clock.test.ts} | 20 +++++++++---------- 5 files changed, 26 insertions(+), 38 deletions(-) rename packages/pgconductor-js/src/lib/{clock-skew.ts => clock.ts} (77%) rename packages/pgconductor-js/tests/unit/{clock-skew.test.ts => clock.test.ts} (87%) diff --git a/packages/pgconductor-js/src/lib/clock-skew.ts b/packages/pgconductor-js/src/lib/clock.ts similarity index 77% rename from packages/pgconductor-js/src/lib/clock-skew.ts rename to packages/pgconductor-js/src/lib/clock.ts index 6e062a6..5be0ac0 100644 --- a/packages/pgconductor-js/src/lib/clock-skew.ts +++ b/packages/pgconductor-js/src/lib/clock.ts @@ -1,29 +1,20 @@ import type { Logger } from "./logger"; -export type LocalClock = () => Date; -export type DatabaseTimeSampler = (signal?: AbortSignal) => Promise; - -export type WorkerClock = { - now(): Date; - start(signal?: AbortSignal): Promise; - stop(): void; -}; - -export const DATABASE_CLOCK_REFRESH_INTERVAL_MS = 10 * 60 * 1000; +const DATABASE_CLOCK_REFRESH_INTERVAL_MS = 10 * 60 * 1000; /** * 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 DatabaseClockOffset implements WorkerClock { +export class Clock { private offsetMs = 0; private refreshTimer: ReturnType | null = null; private running = false; constructor( - private readonly sampleDatabaseTime: DatabaseTimeSampler, + private readonly sampleDatabaseTime: (signal?: AbortSignal) => Promise, private readonly logger: Logger, - private readonly localClock: LocalClock = () => new Date(), + private readonly localClock: () => Date = () => new Date(), private readonly refreshIntervalMs = DATABASE_CLOCK_REFRESH_INTERVAL_MS, ) {} diff --git a/packages/pgconductor-js/src/orchestrator.ts b/packages/pgconductor-js/src/orchestrator.ts index 064a25e..8f68a64 100644 --- a/packages/pgconductor-js/src/orchestrator.ts +++ b/packages/pgconductor-js/src/orchestrator.ts @@ -189,13 +189,13 @@ export class Orchestrator { // Start heartbeat loop this.startHeartbeatLoop(); - // Startup is observed through worker.started below; consume the run promise - // as well so the same startup error is not reported as unhandled. - for (const worker of this.workers) { - const running = runOnce - ? worker.drain(this.orchestratorId) - : worker.run(this.orchestratorId); - running.catch(noop); + // Kick off all workers (don't await yet!) + if (runOnce) { + // Drain mode: workers will process and stop + this.workers.forEach((w) => void w.drain(this.orchestratorId)); + } else { + // Normal mode: workers will run continuously + this.workers.forEach((w) => void w.run(this.orchestratorId)); } // Wait for ALL workers to finish starting (register() complete) diff --git a/packages/pgconductor-js/src/task-context.ts b/packages/pgconductor-js/src/task-context.ts index 5c504d0..184940c 100644 --- a/packages/pgconductor-js/src/task-context.ts +++ b/packages/pgconductor-js/src/task-context.ts @@ -1,6 +1,6 @@ import type { DatabaseClient, JsonValue, Execution, Payload } from "./database-client"; import { nextCronOccurrence } from "./lib/cron"; -import type { WorkerClock } from "./lib/clock-skew"; +import type { Clock } from "./lib/clock"; import type { TaskDefinition, TaskName, @@ -77,7 +77,7 @@ export function createTaskSignal( export type TaskContextOptions = { abortController: TypedAbortController; db: DatabaseClient; - clock: WorkerClock; + clock: Clock; execution: Execution; logger: Logger; window?: [string, string]; diff --git a/packages/pgconductor-js/src/worker.ts b/packages/pgconductor-js/src/worker.ts index b3f2854..61a68c7 100644 --- a/packages/pgconductor-js/src/worker.ts +++ b/packages/pgconductor-js/src/worker.ts @@ -18,7 +18,7 @@ import { Deferred } from "./lib/deferred"; import { type PollableAsyncIterable } from "./lib/async-queue"; import { BatchingAsyncQueue, type BatchGroup } from "./lib/batching-async-queue"; import { nextCronOccurrence } from "./lib/cron"; -import { DatabaseClockOffset, type WorkerClock } from "./lib/clock-skew"; +import { Clock } from "./lib/clock"; import { createTaskSignal, isTaskAbortReason, @@ -38,7 +38,6 @@ import type { TypedAbortController } from "./lib/typed-abort-controller"; */ export type WorkerConfig = { concurrency: number; - clock?: WorkerClock; flushBatchSize: number; fetchBatchSize: number; pollIntervalMs: number; @@ -133,7 +132,7 @@ export class Worker< private readonly flushBatchSize: number; private readonly flushIntervalMs: number; private readonly pollIntervalMs: number; - private readonly clock: WorkerClock; + private readonly clock: Clock; private _startDeferred: Deferred | null = null; private _stopDeferred: Deferred | null = null; @@ -164,9 +163,7 @@ export class Worker< this.flushIntervalMs = fullConfig.flushIntervalMs; this.fetchBatchSize = fullConfig.fetchBatchSize; this.flushBatchSize = fullConfig.flushBatchSize; - this.clock = - fullConfig.clock || - new DatabaseClockOffset((signal) => this.db.getDatabaseTime({ signal }), this.logger); + this.clock = new Clock((signal) => this.db.getDatabaseTime({ signal }), this.logger); } /** diff --git a/packages/pgconductor-js/tests/unit/clock-skew.test.ts b/packages/pgconductor-js/tests/unit/clock.test.ts similarity index 87% rename from packages/pgconductor-js/tests/unit/clock-skew.test.ts rename to packages/pgconductor-js/tests/unit/clock.test.ts index 0e2fdf8..faf03f1 100644 --- a/packages/pgconductor-js/tests/unit/clock-skew.test.ts +++ b/packages/pgconductor-js/tests/unit/clock.test.ts @@ -1,5 +1,5 @@ import { beforeEach, describe, expect, mock, test } from "bun:test"; -import { DatabaseClockOffset } from "../../src/lib/clock-skew"; +import { Clock } from "../../src/lib/clock"; import { nextCronOccurrence } from "../../src/lib/cron"; const logger = { @@ -11,14 +11,14 @@ const logger = { const date = (ms: number): Date => new Date(ms); -describe("DatabaseClockOffset", () => { +describe("Clock", () => { beforeEach(() => { logger.warn.mockClear(); }); test("corrects positive and negative local clock skew", async () => { let now = 1_000; - const positive = new DatabaseClockOffset( + const positive = new Clock( async () => date(2_000), logger, () => date(now), @@ -29,7 +29,7 @@ describe("DatabaseClockOffset", () => { expect(positive.now()).toEqual(date(2_100)); now = 1_000; - const negative = new DatabaseClockOffset( + const negative = new Clock( async () => date(0), logger, () => date(now), @@ -42,7 +42,7 @@ describe("DatabaseClockOffset", () => { test("uses the request midpoint to account for latency", async () => { const localTimes = [date(1_000), date(1_300)]; - const clock = new DatabaseClockOffset( + const clock = new Clock( async () => date(1_100), logger, () => { @@ -58,12 +58,12 @@ describe("DatabaseClockOffset", () => { test("gives skewed workers the same cron slot", async () => { const databaseNow = date(Date.UTC(2025, 0, 1, 12, 0, 0)); - const positive = new DatabaseClockOffset( + const positive = new Clock( async () => databaseNow, logger, () => date(databaseNow.getTime() - 5 * 60 * 1000), ); - const negative = new DatabaseClockOffset( + const negative = new Clock( async () => databaseNow, logger, () => date(databaseNow.getTime() + 5 * 60 * 1000), @@ -78,7 +78,7 @@ describe("DatabaseClockOffset", () => { }); test("retains the local clock when the initial sample fails", async () => { - const clock = new DatabaseClockOffset( + const clock = new Clock( async () => { throw new Error("database unavailable"); }, @@ -98,7 +98,7 @@ describe("DatabaseClockOffset", () => { test("retains the previous offset when a refresh fails", async () => { let fail = false; - const clock = new DatabaseClockOffset( + const clock = new Clock( async () => { if (fail) throw new Error("database unavailable"); return date(2_000); @@ -122,7 +122,7 @@ describe("DatabaseClockOffset", () => { let samples = 0; let inFlight = 0; let maxInFlight = 0; - const clock = new DatabaseClockOffset( + const clock = new Clock( async () => { samples++; inFlight++; From 2fb94125d677861440e4e843f3f82ffeca9c57ff Mon Sep 17 00:00:00 2001 From: psteinroe Date: Thu, 17 Sep 2026 19:56:44 +0000 Subject: [PATCH 5/5] refactor(clock): use options object --- packages/pgconductor-js/src/lib/clock.ts | 29 ++++++-- packages/pgconductor-js/src/worker.ts | 5 +- .../pgconductor-js/tests/unit/clock.test.ts | 66 +++++++++---------- 3 files changed, 60 insertions(+), 40 deletions(-) diff --git a/packages/pgconductor-js/src/lib/clock.ts b/packages/pgconductor-js/src/lib/clock.ts index 5be0ac0..1c63402 100644 --- a/packages/pgconductor-js/src/lib/clock.ts +++ b/packages/pgconductor-js/src/lib/clock.ts @@ -2,6 +2,13 @@ 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. @@ -11,12 +18,22 @@ export class Clock { private refreshTimer: ReturnType | null = null; private running = false; - constructor( - private readonly sampleDatabaseTime: (signal?: AbortSignal) => Promise, - private readonly logger: Logger, - private readonly localClock: () => Date = () => new Date(), - private readonly refreshIntervalMs = DATABASE_CLOCK_REFRESH_INTERVAL_MS, - ) {} + 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); diff --git a/packages/pgconductor-js/src/worker.ts b/packages/pgconductor-js/src/worker.ts index 61a68c7..6337d45 100644 --- a/packages/pgconductor-js/src/worker.ts +++ b/packages/pgconductor-js/src/worker.ts @@ -163,7 +163,10 @@ export class Worker< this.flushIntervalMs = fullConfig.flushIntervalMs; this.fetchBatchSize = fullConfig.fetchBatchSize; this.flushBatchSize = fullConfig.flushBatchSize; - this.clock = new Clock((signal) => this.db.getDatabaseTime({ signal }), this.logger); + this.clock = new Clock({ + sampleDatabaseTime: (signal) => this.db.getDatabaseTime({ signal }), + logger: this.logger, + }); } /** diff --git a/packages/pgconductor-js/tests/unit/clock.test.ts b/packages/pgconductor-js/tests/unit/clock.test.ts index faf03f1..03979d1 100644 --- a/packages/pgconductor-js/tests/unit/clock.test.ts +++ b/packages/pgconductor-js/tests/unit/clock.test.ts @@ -18,22 +18,22 @@ describe("Clock", () => { test("corrects positive and negative local clock skew", async () => { let now = 1_000; - const positive = new Clock( - async () => date(2_000), + const positive = new Clock({ + sampleDatabaseTime: async () => date(2_000), logger, - () => date(now), - ); + localClock: () => date(now), + }); await positive.refresh(); now = 1_100; expect(positive.now()).toEqual(date(2_100)); now = 1_000; - const negative = new Clock( - async () => date(0), + const negative = new Clock({ + sampleDatabaseTime: async () => date(0), logger, - () => date(now), - ); + localClock: () => date(now), + }); await negative.refresh(); now = 1_100; @@ -42,15 +42,15 @@ describe("Clock", () => { test("uses the request midpoint to account for latency", async () => { const localTimes = [date(1_000), date(1_300)]; - const clock = new Clock( - async () => date(1_100), + 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); @@ -58,16 +58,16 @@ describe("Clock", () => { test("gives skewed workers the same cron slot", async () => { const databaseNow = date(Date.UTC(2025, 0, 1, 12, 0, 0)); - const positive = new Clock( - async () => databaseNow, + const positive = new Clock({ + sampleDatabaseTime: async () => databaseNow, logger, - () => date(databaseNow.getTime() - 5 * 60 * 1000), - ); - const negative = new Clock( - async () => databaseNow, + localClock: () => date(databaseNow.getTime() - 5 * 60 * 1000), + }); + const negative = new Clock({ + sampleDatabaseTime: async () => databaseNow, logger, - () => date(databaseNow.getTime() + 5 * 60 * 1000), - ); + localClock: () => date(databaseNow.getTime() + 5 * 60 * 1000), + }); await positive.refresh(); await negative.refresh(); @@ -78,13 +78,13 @@ describe("Clock", () => { }); test("retains the local clock when the initial sample fails", async () => { - const clock = new Clock( - async () => { + const clock = new Clock({ + sampleDatabaseTime: async () => { throw new Error("database unavailable"); }, logger, - () => date(1_000), - ); + localClock: () => date(1_000), + }); await clock.start(); clock.stop(); @@ -98,14 +98,14 @@ describe("Clock", () => { test("retains the previous offset when a refresh fails", async () => { let fail = false; - const clock = new Clock( - async () => { + const clock = new Clock({ + sampleDatabaseTime: async () => { if (fail) throw new Error("database unavailable"); return date(2_000); }, logger, - () => date(1_000), - ); + localClock: () => date(1_000), + }); await clock.refresh(); fail = true; @@ -122,8 +122,8 @@ describe("Clock", () => { let samples = 0; let inFlight = 0; let maxInFlight = 0; - const clock = new Clock( - async () => { + const clock = new Clock({ + sampleDatabaseTime: async () => { samples++; inFlight++; maxInFlight = Math.max(maxInFlight, inFlight); @@ -132,9 +132,9 @@ describe("Clock", () => { return date(1_000 + samples); }, logger, - () => date(0), - 1, - ); + localClock: () => date(0), + refreshIntervalMs: 1, + }); await clock.start(); await new Promise((resolve) => setTimeout(resolve, 25));