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
1 change: 1 addition & 0 deletions .changeset/config.json
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@
"@cleverbrush/deep",
"@cleverbrush/async",
"@cleverbrush/scheduler",
"@cleverbrush/scheduler-postgres",
"@cleverbrush/mapper",
"@cleverbrush/knex-clickhouse",
"@cleverbrush/react-form",
Expand Down
28 changes: 28 additions & 0 deletions .changeset/durable-job-execution.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,28 @@
---
"@cleverbrush/scheduler": major
"@cleverbrush/scheduler-postgres": major
---

Redesign the scheduler for versioned immediate, delayed and recurring jobs,
typed separate handlers, durable ordered progress and fenced worker leases.
Add an independently installed PostgreSQL adapter using Framework ORM and
knex-schema, explicit migrations, transactional enqueue and restart recovery.

Retries are opt-in and rerun whole handlers. Calendar triggers use explicit UTC
or IANA zones with persisted cursors and missed/overlap policies. See the
scheduler v4.x-to-v5 migration guide for the breaking API and rollout steps.
The adapter joins the fixed Framework release group.

Retain schema-driven minute/day/week/month/year definitions and their
discriminated Schedule type, exposing individual schemas and the Schemas facade.
Normalize recurrence defaults, dates and weekday order before fingerprinting;
unchanged registrations retain their cursor and start anchor. Preserve the
one-based calculator index and accept the deprecated maxOccurences spelling
while rejecting ambiguous dual spelling. Name the explicit persistence option
storageRepository. Derive PostgreSQL row and entity types from schema definitions
without parallel hand-written row types.

Use native Date/Intl calendar calculations without an additional date-time
runtime dependency. Preserve DST-gap skipping and earlier-fold selection,
including non-hour transitions and skipped calendar dates, independently of
the host time zone. Bound minute schedules to the representable Date range.
Binary file added .github/pr-evidence/durable-scheduler.png
Loading
Sorry, something went wrong. Reload?
Sorry, we cannot display this file.
Sorry, this file is invalid so it cannot be displayed.
4 changes: 4 additions & 0 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -63,6 +63,7 @@ jobs:
--health-retries 10
env:
QUERY_TEST_DATABASE_URL: postgres://framework_test:framework_test@127.0.0.1:5432/framework_queries
SCHEDULER_TEST_DATABASE_URL: postgres://framework_test:framework_test@127.0.0.1:5432/framework_queries
steps:
- uses: actions/checkout@v4
- uses: actions/setup-node@v4
Expand All @@ -72,3 +73,6 @@ jobs:
- run: npm ci
- run: npm run build
- run: npm run test:queries:integration
- run: npm run test:scheduler:integration
- run: node demos/durable-jobs/demo.ts
- run: node demos/durable-jobs/periodic.ts --fast
3 changes: 2 additions & 1 deletion AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -58,7 +58,8 @@ scripts/ ← build/release helper scripts
| `@cleverbrush/async` | Async utilities: Collector, debounce, throttle, retry |
| `@cleverbrush/mapper` | Schema-driven object mapper |
| `@cleverbrush/react-form` | React form library powered by schema PropertyDescriptors |
| `@cleverbrush/scheduler` | Cron-like job scheduler with schema-validated config |
| `@cleverbrush/scheduler` | Typed durable jobs, recurring triggers and progress |
| `@cleverbrush/scheduler-postgres` | PostgreSQL job repository and explicit migrations |
| `@cleverbrush/server` | Schema-first HTTP server: DI, auto-validation, RFC 9457 errors |
| `@cleverbrush/server-openapi` | OpenAPI 3.x generation from server endpoints |
| `@cleverbrush/client` | Type-safe HTTP client for `@cleverbrush/server` endpoints |
Expand Down
3 changes: 2 additions & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -39,7 +39,8 @@ JSON Schema, API contracts, and Standard Schema integrations.
| [`@cleverbrush/otel`](./libs/otel) | OpenTelemetry setup and instrumentation helpers for apps and clients. |
| [`@cleverbrush/async`](./libs/async) | Async utilities including collector, debounce, throttle, and retry. |
| [`@cleverbrush/deep`](./libs/deep) | Deep equality, deep extension, flattening, and object utilities. |
| [`@cleverbrush/scheduler`](./libs/scheduler) | Cron-like job scheduler with schema-validated job configuration. |
| [`@cleverbrush/scheduler`](./libs/scheduler) | Typed durable jobs, recurring triggers and ordered progress. |
| [`@cleverbrush/scheduler-postgres`](./libs/scheduler-postgres) | PostgreSQL job persistence, transactional enqueue and fenced leases. |

## How The Pieces Fit

Expand Down
22 changes: 22 additions & 0 deletions demos/durable-jobs/README.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,22 @@
# Durable job contracts and separated handlers

After npm ci and npm run build, run `node demos/durable-jobs/demo.ts`
with Node 24. It prints queued, running, progress and succeeded events, then
the typed result. This small example deliberately uses process-local memory.

The core scheduler README shows PostgreSQL setup for real durability.
The sibling crash-worker.mjs and thread-handler.mjs are process-level fixtures
exercised by the scheduler integration and packaged-thread tests, respectively.

## Periodic execution

Run `node demos/durable-jobs/periodic.ts` to execute two reports one minute
apart. The example validates a minute schedule, registers it with
`upsertSchedule`, starts both the dispatcher and worker, waits for both
completed runs, and stops dispatch before draining the worker.

For a fast smoke test, run `node demos/durable-jobs/periodic.ts --fast`.
Only this example's in-memory repository clock advances by one minute after
the first completion; the periodic dispatcher and worker still execute normally.
CI runs this mode. This is a test/demo clock, not a production scheduling option.
Choose PostgreSQL via `storageRepository` for restart-safe recurring jobs.
9 changes: 9 additions & 0 deletions demos/durable-jobs/contracts.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,9 @@
import { object, string, number } from '@cleverbrush/schema';
import { defineJob } from '@cleverbrush/scheduler';

export const Report = defineJob({
name: 'report', version: 1,
input: object({ reportId: string() }),
progress: object({ percent: number() }),
output: object({ downloadUrl: string() })
});
14 changes: 14 additions & 0 deletions demos/durable-jobs/crash-worker.mjs
Original file line number Diff line number Diff line change
@@ -0,0 +1,14 @@
// Integration fixture: deliberately never finishes after committing progress.
import knex from 'knex';
import { object, string, number } from '@cleverbrush/schema';
import { defineJob, JobScheduler } from '@cleverbrush/scheduler';
import { PostgresJobRepository } from '@cleverbrush/scheduler-postgres';
const db = knex({ client: 'pg', connection: process.env.SCHEDULER_TEST_DATABASE_URL });
const job = defineJob({ name: 'report', version: 1, input: object({ id: string() }), progress: object({ percent: number() }), output: object({ url: string() }), retry: { maxAttempts: 2, initialDelayMs: 1 } });
const scheduler = new JobScheduler({ storageRepository: new PostgresJobRepository(db, { tablePrefix: process.env.JOB_TABLE_PREFIX }), namespace: process.env.JOB_NAMESPACE });
const worker = scheduler.createWorker({ jobs: [job.handle(async (_, context) => {
await context.report({ percent: 50 });
process.send?.({ ready: true });
await new Promise(() => {});
})], pollIntervalMs: 10, leaseMs: 500, heartbeatMs: 100 });
await worker.start();
16 changes: 16 additions & 0 deletions demos/durable-jobs/demo.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,16 @@
import { JobScheduler, InMemoryJobRepository } from '@cleverbrush/scheduler';
import { Report } from './contracts.ts';
import { handleReport } from './handler.ts';

const jobs = new JobScheduler({ storageRepository: new InMemoryJobRepository(), pollIntervalMs: 10 });
const worker = jobs.createWorker({ jobs: [Report.handle(handleReport)], pollIntervalMs: 10 });
const run = await jobs.enqueue(Report, { reportId: 'quarterly' });
await worker.start();
try {
for await (const event of jobs.events(Report, run.id)) {
console.log(event.sequence, event.type, event.data);
}
const snapshot = await jobs.getRun(Report, run.id);
console.log(snapshot?.status, snapshot?.output);
if (snapshot?.status !== 'succeeded') process.exitCode = 1;
} finally { await worker.stop(); }
9 changes: 9 additions & 0 deletions demos/durable-jobs/handler.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,9 @@
import type { JobHandler } from '@cleverbrush/scheduler';
import type { Report } from './contracts.ts';

/** Handler in a separate module retains all contract-inferred types. */
export const handleReport: JobHandler<typeof Report> = async (input, context) => {
context.signal.throwIfAborted();
await context.report({ percent: 50 });
return { downloadUrl: '/reports/' + input.reportId };
};
42 changes: 42 additions & 0 deletions demos/durable-jobs/periodic.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,42 @@
import { setTimeout as delay } from 'node:timers/promises';
import { InMemoryJobRepository, JobScheduler, ScheduleSchema } from '@cleverbrush/scheduler';
import { Report } from './contracts.ts';
import { handleReport } from './handler.ts';

// --fast advances only this demo's in-memory clock after the first completion.
const fast = process.argv.includes('--fast');
let clock = Date.now();
const jobs = new JobScheduler({
storageRepository: new InMemoryJobRepository({ now: () => fast ? clock : Date.now() }),
pollIntervalMs: 10
});
const worker = jobs.createWorker({ jobs: [Report.handle(handleReport)], pollIntervalMs: 10 });
await jobs.upsertSchedule('minute-report', Report, { reportId: 'periodic' }, {
schedule: ScheduleSchema.parse({ every: 'minute', interval: 1, maxOccurrences: 2 }),
missed: 'coalesce', overlap: 'skip'
});

try {
await worker.start(); // Executes jobs accepted by the dispatcher.
await jobs.start(); // Materializes due occurrences, not handler execution.
const deadline = Date.now() + (fast ? 5000 : 65000);
let advanced = false;
for (;;) {
const completed = (await jobs.health()).counts.succeeded ?? 0;
if (completed === 2) {
console.log('Completed both periodic reports.');
break;
}
if (jobs.lastError) throw jobs.lastError;
if (worker.lastError) throw worker.lastError;
if (Date.now() > deadline) throw new Error('Periodic demo timed out');
if (fast && completed === 1 && !advanced) {
clock += 60000;
advanced = true;
}
await delay(10);
}
} finally {
await jobs.stop(); // Stop producing before draining the worker.
await worker.stop();
}
12 changes: 12 additions & 0 deletions demos/durable-jobs/thread-handler.mjs
Original file line number Diff line number Diff line change
@@ -0,0 +1,12 @@
// Trusted default-export handler for a worker-thread execution.
export default async function handler(input, context) {
if (input.mode === 'exit') process.exit(7);
if (input.mode === 'hang') await new Promise(() => {});
if (input.mode === 'invalid') { await context.report({ percent: 'wrong' }); }
if (input.mode === 'accessor') {
await context.report({ get percent() { throw new Error('must not execute'); } });
}
if (input.mode === 'unawaited-invalid') void context.report({ percent: 'wrong' });
await context.report({ percent: 50 });
return { url: '/reports/' + input.id };
}
1 change: 1 addition & 0 deletions demos/todo-backend/Dockerfile
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@ COPY libs/react-form/package.json ./libs/react-form/
COPY libs/schema/package.json ./libs/schema/
COPY libs/schema-json/package.json ./libs/schema-json/
COPY libs/scheduler/package.json ./libs/scheduler/
COPY libs/scheduler-postgres/package.json ./libs/scheduler-postgres/
COPY libs/server/package.json ./libs/server/
COPY libs/server-openapi/package.json ./libs/server-openapi/
COPY libs/otel/package.json ./libs/otel/
Expand Down
1 change: 1 addition & 0 deletions demos/todo-frontend/Dockerfile
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@ COPY libs/react-form/package.json ./libs/react-form/
COPY libs/schema/package.json ./libs/schema/
COPY libs/schema-json/package.json ./libs/schema-json/
COPY libs/scheduler/package.json ./libs/scheduler/
COPY libs/scheduler-postgres/package.json ./libs/scheduler-postgres/
COPY libs/server/package.json ./libs/server/
COPY libs/server-openapi/package.json ./libs/server-openapi/
COPY libs/client/package.json ./libs/client/
Expand Down
91 changes: 91 additions & 0 deletions libs/scheduler-postgres/README.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,91 @@
# @cleverbrush/scheduler-postgres

PostgreSQL persistence for [@cleverbrush/scheduler](../scheduler/README.md).
Uses @cleverbrush/orm and @cleverbrush/knex-schema for schemas, migrations and
routine queries. Isolated native queries provide database-time leases,
SKIP LOCKED row claims and aggregate health checks.

## Install and migrate

```sh
npm install @cleverbrush/scheduler @cleverbrush/scheduler-postgres knex pg
```

```ts
// An explicit migration, invoked once by your existing migration runner.
import { createSchedulerTables, dropSchedulerTables } from '@cleverbrush/scheduler-postgres';
import type { Knex } from 'knex';

export const up = (knex: Knex) => createSchedulerTables(knex);
export const down = (knex: Knex) => dropSchedulerTables(knex);
```

Down is destructive: it deletes runs, schedules, progress and attempts.
No migration runs implicitly when constructing a repository or starting workers.
A tablePrefix option (default cb_jobs) supports separate storage ownership.
Use exactly the same prefix for migrations and repositories.

```ts
import knex from 'knex';
import { JobScheduler } from '@cleverbrush/scheduler';
import { PostgresJobRepository } from '@cleverbrush/scheduler-postgres';

const database = knex({
client: 'pg', connection: process.env.DATABASE_URL,
pool: { min: 0, max: 10 }, acquireConnectionTimeout: 5000
});
const jobs = new JobScheduler({
storageRepository: new PostgresJobRepository(database),
namespace: 'reporting'
});
```

The adapter does not close the caller's pool. Stop workers and dispatchers first,
then database.destroy(). Use PostgreSQL's default READ COMMITTED isolation.
Transitions set local lock and statement timeouts (5 and 15 seconds). Configure
connection acquisition and network timeouts for your deployment as well.

## Transactional enqueue

```ts
await database.transaction(async transaction => {
// Write application data using this same transaction.
const producer = new JobScheduler({
storageRepository: new PostgresJobRepository(transaction),
namespace: 'reporting'
});
await producer.enqueue(Report, { reportId }, { idempotencyKey: reportId });
});
```

The enqueue resolves within a savepoint; durable acceptance occurs only when the
outer transaction commits. Rollback removes both the application write and job.
Do not create/start workers on a transaction-bound repository. Local timeout
settings also apply to the enclosing transaction after savepoint release.

## Storage and guarantees

Four library-owned tables store runs, attempts, events and recurring triggers.
Indexed scalar columns support claims, expiry, overlap and health queries.
Opaque strict-JSON record snapshots are stored as text so arbitrary validated
payloads do not need application-specific database schemas or lossy driver
decoding. Large artifacts belong in object storage, referenced by payload keys.

Unique namespace/dedupe keys protect concurrent producers. Locked schedule
cursors and occurrence keys protect competing dispatchers. State transitions,
attempt history and ordered events commit together. Cleanup cascades dependent
events and attempts, but never removes active jobs or triggers.

The repository uses PostgreSQL clock_timestamp for lease decisions, not worker
wall clocks. Expired owners cannot resurrect a lease. Physical job side effects
remain at-least-once: see the core guide's retry/idempotency rules.

Integration tests require a disposable PostgreSQL database. They create unique
table prefixes and drop only those tables:

```sh
SCHEDULER_TEST_DATABASE_URL=postgres://... npm run test:scheduler:integration
```

Tests exercise competing producers/workers/dispatchers, rollback, SKIP LOCKED,
lease expiry, retention and SIGKILL/restart recovery. CI runs them on PostgreSQL 16.
Loading
Loading