Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 6 additions & 3 deletions packages/pgconductor-js/src/database-client.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<Date> {
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;
}
Expand Down
80 changes: 80 additions & 0 deletions packages/pgconductor-js/src/lib/clock.ts
Original file line number Diff line number Diff line change
@@ -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<Date>;
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<typeof setTimeout> | 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<Date>;
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<void> {
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<void> {
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);
}
}
5 changes: 5 additions & 0 deletions packages/pgconductor-js/src/lib/cron.ts
Original file line number Diff line number Diff line change
@@ -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();
}
7 changes: 4 additions & 3 deletions packages/pgconductor-js/src/task-context.ts
Original file line number Diff line number Diff line change
@@ -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,
Expand Down Expand Up @@ -76,6 +77,7 @@ export function createTaskSignal(
export type TaskContextOptions = {
abortController: TypedAbortController<TaskAbortReasons>;
db: DatabaseClient;
clock: Clock;
execution: Execution;
logger: Logger;
window?: [string, string];
Expand Down Expand Up @@ -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(
Expand Down
19 changes: 13 additions & 6 deletions packages/pgconductor-js/src/worker.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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<void> | null = null;
private _stopDeferred: Deferred<void> | null = null;
Expand Down Expand Up @@ -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,
});
}

/**
Expand Down Expand Up @@ -221,7 +227,8 @@ export class Worker<
this._stopDeferred = new Deferred<void>();
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
Expand All @@ -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;
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -585,6 +592,7 @@ export class Worker<
TaskContext.create<Tasks, Events, typeof extraContext>(
{
db: this.db,
clock: this.clock,
abortController: taskAbortController,
execution: exec,
logger: makeChildLogger(this.logger, {
Expand Down Expand Up @@ -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}`;

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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");
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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<Date> {
return this.getInternalTime();
}

async getDatabaseTime(): Promise<Date> {
return this.getInternalTime();
}

// ============================================================================
// Orchestrator Management
// ============================================================================
Expand Down
149 changes: 149 additions & 0 deletions packages/pgconductor-js/tests/unit/clock.test.ts
Original file line number Diff line number Diff line change
@@ -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);
});
});
Loading