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
32 changes: 16 additions & 16 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,7 @@

A Corbits hub module that wakes Interchange agents on five-field UTC cron schedules, mounted as `@intx/hub-api` routes on the hub (Interchange's multi-tenant control plane) and stored in its Postgres. A ticker mails each due schedule to its agent's current run, the running instance of the agent's workflow definition.

Schedules support create, list and delete. There is no edit: delete and recreate.
Schedules support create, list and delete. There is no edit: delete and recreate. An expression must be able to fire (no Feb 30), at most 256 characters and 64 comma clauses.

## Why @corbits/cron?

Expand Down Expand Up @@ -59,15 +59,15 @@ Paths are relative to where the host mounts the sub-app.

Returns `{ start(), stop() }`.

| `opts` | Type | Purpose |
| ------------------- | -------------------------------------- | ---------------------------------------------------------------------------------- |
| `db` | `CronDb` | The hub's drizzle handle. |
| `deliver` | `DeliverCronMail` | Sends one due schedule's mail. `createRunTriggerCronDeliver(deliverer)` builds it. |
| `intervalMs` | `number` (optional) | Poll period. Defaults to 60 000. |
| `onTickError` | `(error) => void` (optional) | A tick failed before delivering, such as a lost connection. The next tick retries. |
| `onDeliveryError` | `(error, schedule) => void` (optional) | One delivery failed. |
| `onScheduleWaiting` | `(schedule) => void` (optional) | A schedule started waiting for its agent's next run. |
| `onScheduleStopped` | `(schedule) => void` (optional) | A schedule stopped because its agent was deleted. |
| `opts` | Type | Purpose |
| ------------------- | -------------------------------------- | ------------------------------------------------------------------------------------------------------- |
| `db` | `CronDb` | The hub's drizzle handle. |
| `deliver` | `DeliverCronMail` | Sends one due schedule's mail. `createRunTriggerCronDeliver(deliverer)` builds it. |
| `intervalMs` | `number` (optional) | Poll period. Defaults to 60 000. |
| `onTickError` | `(error) => void` (optional) | A tick failed before delivering, such as a lost connection. The next tick retries. |
| `onDeliveryError` | `(error, schedule) => void` (optional) | One delivery failed. |
| `onScheduleWaiting` | `(schedule) => void` (optional) | A schedule started waiting for its agent's next run. |
| `onScheduleStopped` | `(schedule) => void` (optional) | A schedule stopped: `reason` is `agent_deleted`, or `invalid_expression` for a row that can never fire. |

A deliverer throws an error carrying `RUN_GRANTS_NOT_ROUTABLE` or `RUN_MAIL_NOT_ROUTABLE` (checked with `isUnroutableRunTrigger`) when a run's address is dead. If that run's sidecar is already released, the ticker marks the run failed so the schedule waits for the agent's next run.

Expand All @@ -77,12 +77,12 @@ From `@corbits/cron/migrations`. Takes the same arguments as `@intx/db`'s `runMi

### Other exports

| Export | Use |
| ---------------------------------------------------------------------------------------------------------------------------------- | ---------------------------------------------------- |
| `isValidCronExpression(expression)` | The check `POST /` runs. |
| `nextCronFireAfter(expression, after)` | The next matching UTC minute after `after`. |
| `isUnroutableRunTrigger(error)`, `RUN_GRANTS_NOT_ROUTABLE`, `RUN_MAIL_NOT_ROUTABLE` | The dead-address error contract a deliverer follows. |
| `CronRoutesDeps`, `CreateCronTickerOpts`, `CronTicker`, `CronDb`, `DeliverCronMail`, `RunTriggerDeliverer`, `UnroutableRunTrigger` | Types for the above. |
| Export | Use |
| ---------------------------------------------------------------------------------------------------------------------------------- | -------------------------------------------------------------------- |
| `isValidCronExpression(expression)` | Syntax, range and cap check. `POST /` also requires a possible fire. |
| `nextCronFireAfter(expression, after)` | The next matching UTC minute after `after`. |
| `isUnroutableRunTrigger(error)`, `RUN_GRANTS_NOT_ROUTABLE`, `RUN_MAIL_NOT_ROUTABLE` | The dead-address error contract a deliverer follows. |
| `CronRoutesDeps`, `CreateCronTickerOpts`, `CronTicker`, `CronDb`, `DeliverCronMail`, `RunTriggerDeliverer`, `UnroutableRunTrigger` | Types for the above. |

## Using with Interchange

Expand Down
31 changes: 31 additions & 0 deletions e2e/routes.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -78,6 +78,37 @@ describeIfDb("createCronRoutes", () => {
}
});

test("an expression that can never fire or is over the caps is rejected", async () => {
const { db, close } = createDB(requireDatabase().config);
try {
const tenantId = `tnt_cron_expr_${randomUUID().slice(0, 8)}`;
await seedTenant(db, tenantId);
await seedDeployment(db, tenantId, "agent-expr-source");
const app = cronRoutesApp(db, tenantId, allowAll);

const huge = Array.from({ length: 5000 }, () => "1").join(",");
for (const expression of ["0 0 31 2 *", `${huge} * * * *`]) {
const response = await post(app, {
expression,
definitionName: "agent-expr-source",
subject: "s",
body: "b",
});
expect(response.status).toBe(400);
expect(await response.json()).toEqual({ error: "invalid_expression" });
}
const leap = await post(app, {
expression: "0 0 29 2 *",
definitionName: "agent-expr-source",
subject: "s",
body: "b",
});
expect(leap.status).toBe(201);
} finally {
await close();
}
});

test("each route is gated by the host's requireGrant", async () => {
const { db, close } = createDB(requireDatabase().config);
try {
Expand Down
67 changes: 67 additions & 0 deletions e2e/ticker.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -228,6 +228,73 @@ describeIfDb("createCronTicker", () => {
}
});

test("a long-gap schedule fires; one that can never fire stops, fast", async () => {
const { db, close } = createDB(requireDatabase().config);
try {
const tenantId = `tnt_cron_expr_${randomUUID().slice(0, 8)}`;
await seedTenant(db, tenantId);
await seedDeployment(db, tenantId, "agent-expr-source");

const suffix = randomUUID().slice(0, 8);
const leapId = `sched_leap_${suffix}`;
const feb31Id = `sched_feb31_${suffix}`;
const hugeId = `sched_huge_${suffix}`;
const huge = Array.from({ length: 5000 }, () => "1").join(",");
const row = (id: string, expression: string, createdAt: Date) => ({
id,
tenantId,
expression,
definitionName: "agent-expr-source",
subject: id,
body: "b",
createdAt,
});
// Feb 29 2024 is more than a year after this row was saved.
await db
.insert(cronScheduleTable)
.values([
row(leapId, "0 0 29 2 *", new Date("2021-03-01T00:00:00Z")),
row(feb31Id, "0 0 31 2 *", new Date(Date.now() - 2 * 60_000)),
row(hugeId, `${huge} 0 31 2 *`, new Date(Date.now() - 2 * 60_000)),
]);

const delivered: string[] = [];
const stopped: Array<{ id: string; reason: string }> = [];
const ticker = createCronTicker({
db,
intervalMs: 20,
deliver: (message) => {
delivered.push(message.subject);
},
onScheduleStopped: (schedule) =>
stopped.push({ id: schedule.id, reason: schedule.reason }),
});
const started = performance.now();
ticker.start();
// Polling, not a fixed sleep: tick latency follows DB load.
for (
let i = 0;
i < 250 && (delivered.length === 0 || stopped.length < 2);
i++
) {
await new Promise((resolve) => setTimeout(resolve, 20));
}
const elapsed = performance.now() - started;
ticker.stop();
// Let an in-flight tick settle before close() ends the client.
await new Promise((resolve) => setTimeout(resolve, 300));

expect(delivered).toEqual([leapId]);
expect(stopped.sort((a, b) => a.id.localeCompare(b.id))).toEqual([
{ id: feb31Id, reason: "invalid_expression" },
{ id: hugeId, reason: "invalid_expression" },
]);
expect(elapsed).toBeLessThan(2_000);
} finally {
await close();
}
});

test("two concurrent tickers deliver a due row exactly once", async () => {
const { db: dbA, close: closeA } = createDB(requireDatabase().config);
const { db: dbB, close: closeB } = createDB(requireDatabase().config);
Expand Down
105 changes: 105 additions & 0 deletions src/cron.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,11 @@ import { describe, expect, test } from "bun:test";

import {
cronExpressionCanFire,
cronIsDue,
cronMatchesMinute,
MAX_CRON_CLAUSES,
MAX_CRON_EXPRESSION_LENGTH,
parseCronExpression,
isValidCronExpression,
isValidTimeZone,
nextCronFireAfter,
Expand Down Expand Up @@ -268,3 +272,104 @@ describe("timezone matching and DST", () => {
expect(zonedParts(next, "America/Los_Angeles").hour).toBe(9);
});
});

describe("step without a range end runs to the field maximum", () => {
test("5/2 on minutes is 5,7,…,59", () => {
const at = (minute: number) => new Date(Date.UTC(2026, 0, 1, 0, minute));
expect(cronMatchesMinute("5/2 * * * *", at(5))).toBe(true);
expect(cronMatchesMinute("5/2 * * * *", at(7))).toBe(true);
expect(cronMatchesMinute("5/2 * * * *", at(59))).toBe(true);
expect(cronMatchesMinute("5/2 * * * *", at(6))).toBe(false);
expect(cronMatchesMinute("5/2 * * * *", at(4))).toBe(false);
});
});

describe("limits", () => {
test("rejects an expression over the length cap", () => {
const long = `${"0,".repeat(MAX_CRON_EXPRESSION_LENGTH)}0 * * * *`;
expect(isValidCronExpression(long)).toBe(false);
});

test("rejects an expression over the clause cap", () => {
const clauses = Array.from({ length: MAX_CRON_CLAUSES }, () => "1").join(
",",
);
expect(isValidCronExpression(`${clauses} * * * *`)).toBe(false);
});

test("never-firing expressions cannot fire; Feb 29 can", () => {
expect(cronExpressionCanFire("0 0 31 2 *")).toBe(false);
expect(cronExpressionCanFire("0 0 30 2 *")).toBe(false);
expect(cronExpressionCanFire("0 0 31 4,6,9,11 *")).toBe(false);
expect(cronExpressionCanFire("0 0 29 2 *")).toBe(true);
expect(cronExpressionCanFire("0 0 31 2 5")).toBe(true);
});
});

describe("cronIsDue", () => {
const parsed = (expression: string) => {
const result = parseCronExpression(expression);
if (result === undefined) throw new Error(`unparseable: ${expression}`);
return result;
};

test("a Feb 29 schedule created years earlier fires on Feb 29", () => {
const after = new Date("2025-03-01T00:00:00Z");
expect(
cronIsDue(parsed("0 0 29 2 *"), after, new Date("2028-02-28T23:59:00Z")),
).toBe(false);
expect(
cronIsDue(parsed("0 0 29 2 *"), after, new Date("2028-02-29T00:00:30Z")),
).toBe(true);
});

test("Feb 29 across the 2100 non-leap gap", () => {
expect(
nextCronFireAfter("0 0 29 2 *", new Date("2096-03-01T00:00:00Z")),
).toEqual(new Date("2104-02-29T00:00:00Z"));
});

test("keeps Vixie OR when both day fields are restricted", () => {
// 2026-01-16 is a Friday, not the 13th.
expect(
cronIsDue(
parsed("0 0 13 * 5"),
new Date("2026-01-15T00:00:00Z"),
new Date("2026-01-16T00:00:00Z"),
),
).toBe(true);
});

test("a row last fired decades ago costs a bounded scan", () => {
const start = performance.now();
expect(
cronIsDue(
parsed("0 0 29 2 *"),
new Date("1970-01-01T00:00:00Z"),
new Date("2027-12-31T00:00:00Z"),
),
).toBe(true);
expect(performance.now() - start).toBeLessThan(100);
});
});

describe("cost is bounded", () => {
test("a 5,000-clause expression is rejected without scanning", () => {
const big = Array.from({ length: 5000 }, () => "1").join(",");
const start = performance.now();
expect(parseCronExpression(`${big} 0 31 2 *`)).toBeUndefined();
expect(cronExpressionCanFire(`${big} 0 31 2 *`)).toBe(false);
expect(performance.now() - start).toBeLessThan(50);
});

test("the costliest accepted expression finds its fire quickly", () => {
const minutes = Array.from({ length: MAX_CRON_CLAUSES - 4 }, (_, i) =>
String(59 - i),
).join(",");
const expression = `${minutes} 23 29 2 *`;
expect(isValidCronExpression(expression)).toBe(true);
const start = performance.now();
nextCronFireAfter(expression, new Date("2096-03-01T00:00:00Z"));
expect(performance.now() - start).toBeLessThan(100);
});
});
Loading
Loading