Skip to content
Open
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
48 changes: 21 additions & 27 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,15 @@ Install the package with your preferred JavaScript package manager:

The package exports TypeScript source files directly. Use it from runtimes and bundlers that can load TypeScript subpath exports, or compile it as part of your application build.

### Runtime support

This package publishes its TypeScript sources as its entry points; it does not ship a compiled `dist/`. It is therefore consumed directly by:

- **Bun**, which runs TypeScript (and the `bun:sqlite`/`bun:test` imports used by the Bun adapter) natively, and
- **bundler-based TypeScript toolchains** (Vite, esbuild, Wrangler, Webpack, etc.) that resolve and compile TypeScript subpath exports as part of your application build.

Plain Node.js consuming the package without a TypeScript-aware loader/bundler, or a standalone `tsc` emit that expects prebuilt `.d.ts`/`.js` artifacts, is not supported. Compile the SDK as part of your own build if your toolchain requires that.

## Exports

| Import path | Purpose |
Expand Down Expand Up @@ -180,27 +189,11 @@ The adapter uses Bun's `RedisClient` when available. Raw Redis commands are sent
| Mode | Use case | Behavior |
| --- | --- | --- |
| `external` | Kubernetes CronJob, systemd timer, queue worker, tests | You call `runtime.tick()` and/or `runtime.processReady()` yourself |
| `in-process` | Long-running Bun process | Registers `Bun.cron(schedule, handler)` and processes due work in the same process |
| `os` | Single-server production cron | Registers `Bun.cron(path, schedule, title)` and expects the target module to export `scheduled()` |
| `redis` | Multi-instance Bun deployment | Uses Bun cron as a wake-up mechanism and Redis for claims/leases |

For OS-level Bun cron, register cron jobs from the long-running app:

```ts
const runtime = createBunWorkflowRuntime({
registry,
adapter,
scheduler: {
mode: "os",
scriptPath: import.meta.path,
titlePrefix: "my-app-workflows",
},
});
| `in-process` | Long-running Bun process | Schedules per-cron timers from the SDK's cron parser and processes due work in the same process (plus a periodic interval tick as a safety net) |
| `redis` | Multi-instance Bun deployment | Periodic interval tick as a wake-up mechanism and Redis for claims/leases |
| `os` | — | Not supported. `start()` throws: Bun has no OS cron registration API. Use `external` with your platform scheduler instead |

runtime.start();
```

The target module must export Bun's scheduled handler:
For an external scheduler (OS cron, Kubernetes CronJob, etc.), point it at a module that exports the scheduled handler:

```ts
import {
Expand All @@ -213,12 +206,12 @@ export default createBunWorkflowScheduledHandler({
registry,
adapter: new BunSqliteWorkflowAdapter({ path: "workflows.sqlite" }),
scheduler: {
mode: "os",
mode: "external",
},
});
```

Bun cron uses standard 5-field cron expressions. Bun parses and runs in-process cron schedules in UTC. OS-level Bun cron follows the host timezone because it delegates to the platform scheduler. The SDK accounts for that in `scheduled()` by evaluating OS cron ticks in the local timezone unless a cron definition sets `timezone`.
Cron expressions use the standard 5-field form (an optional leading seconds field makes it 6). The SDK's own parser supports ranges, lists, steps, month/day names, and standard DOM-vs-DOW OR semantics. Schedules are evaluated in UTC unless a cron definition sets `timezone`; `scheduled()` invocations triggered by an OS-level scheduler (`controller.cron` set) are evaluated in the host's local timezone unless a cron definition overrides it.

### Cron Idempotency

Expand Down Expand Up @@ -305,6 +298,8 @@ export default createCloudflareDispatchHandler({

The handler calls the Cloudflare binding with `Workflow.create({ id, params })`. Scheduled envelopes are passed as params and delayed by the Workflow entrypoint helper.

> **Rate limiting is best-effort.** The optional `rateLimit` option uses an in-memory fixed-window counter scoped to a single isolate. Cloudflare Workers run many isolates across PoPs, each with its own window, so the effective accepted rate is `max × <isolate count>` and a client spreading requests can bypass it. For true global enforcement, back it with a [Cloudflare Rate Limiting binding](https://developers.cloudflare.com/workers/runtime-apis/bindings/rate-limit/) or a Durable Object counter.

### Workflow Entrypoint Helper

Use `createCloudflareWorkflowEntrypoint()` when you want Cloudflare Workflows to execute SDK workflow definitions directly.
Expand Down Expand Up @@ -414,16 +409,15 @@ The root export includes:
- Keep workflow instance IDs under Cloudflare's current instance ID limit when using the REST adapter.
- Keep workflow names under Cloudflare's current workflow name limit when using the REST adapter.
- Use unique, stable step names. Step results are keyed by instance ID and step name.
- Do not rely on Bun's fallback cron parser for production semantics. Production Bun scheduling should use Bun's native cron support.
- In-process Bun cron uses UTC. OS-level Bun cron uses the host timezone.
- `Bun.cron(path, schedule, title)` re-registers the same title in place, so keep `titlePrefix`, workflow name, and cron name stable.
- Cron schedules are parsed by the SDK's own parser on every runtime (full 5/6-field support: ranges, lists, steps, names, DOM/DOW OR semantics).
- In-process cron timers evaluate in UTC unless a cron definition sets `timezone`; OS-triggered `scheduled()` ticks use the host timezone by default.
- This package currently ships TypeScript source via exports; compile before publishing to runtimes that cannot load TypeScript directly.

## Verification

Useful package-level checks:

```bash
bun test packages/workflows-sdk/src
bunx tsc -p packages/workflows-sdk/tsconfig.json --noEmit
bun test
bunx tsc --noEmit
```
1 change: 1 addition & 0 deletions package.json
Original file line number Diff line number Diff line change
Expand Up @@ -54,6 +54,7 @@
},
"files": [
"src",
"!src/**/*.test.ts",
"README.md",
"LICENSE"
],
Expand Down
116 changes: 111 additions & 5 deletions src/bun/redis-adapter.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@ class FakeRedis {
readonly strings = new Map<string, string>();
readonly hashes = new Map<string, Record<string, string>>();
readonly sortedSets = new Map<string, Map<string, number>>();
readonly expirations = new Map<string, number>();

async send(command: string, args: unknown[]): Promise<unknown> {
switch (command.toUpperCase()) {
Expand All @@ -19,8 +20,19 @@ class FakeRedis {
}
case "GET":
return this.strings.get(String(args[0])) ?? null;
case "DEL":
return this.strings.delete(String(args[0])) ? 1 : 0;
case "DEL": {
const key = String(args[0]);
const removed = this.strings.delete(key);
this.hashes.delete(key);
return removed ? 1 : 0;
}
case "EXPIRE": {
const key = String(args[0]);
this.expirations.set(key, Number(args[1]));
return this.strings.has(key) || this.hashes.has(key) || this.sortedSets.has(key)
? 1
: 0;
}
case "HSET": {
const key = String(args[0]);
const fields = args.slice(1).map(String);
Expand Down Expand Up @@ -62,8 +74,31 @@ class FakeRedis {
case "EVAL": {
const script = String(args[0]);
if (script.includes("ZRANGEBYSCORE")) {
const [, , queueKey, leasePrefix, processingKey, maxScore, ttlMs, now, token] =
args as [string, number, string, string, string, string | number, number, number, string];
const [
,
,
queueKey,
leasePrefix,
processingKey,
instancePrefix,
maxScore,
,
now,
token,
updatedAt,
] = args as [
string,
number,
string,
string,
string,
string,
string | number,
number,
number,
string,
string,
];
const max = maxScore === "+inf" ? Infinity : Number(maxScore);
const due = [...(this.sortedSets.get(queueKey)?.entries() ?? [])]
.filter(([, score]) => score <= max)
Expand All @@ -76,8 +111,14 @@ class FakeRedis {
this.strings.set(leaseKey, String(token));
this.sortedSets.get(queueKey)?.delete(due);
const processing = this.sortedSets.get(processingKey) ?? new Map<string, number>();
processing.set(due, Number(now) + Number(ttlMs));
// Score is the claim time (mirrors the real Lua script).
processing.set(due, Number(now));
this.sortedSets.set(processingKey, processing);
const instanceHash = this.hashes.get(`${instancePrefix}:${due}`);
if (instanceHash) {
instanceHash.status = "running";
instanceHash.updatedAt = String(updatedAt);
}
return due;
}

Expand Down Expand Up @@ -242,4 +283,69 @@ describe("BunRedisWorkflowAdapter", () => {
await expect(adapter.hasStepResult("job-1", "side-effect")).resolves.toBe(true);
await expect(adapter.getStepResult("job-1", "side-effect")).resolves.toBeUndefined();
});

test("recovers a stalled job whose status never advanced to running", async () => {
const redis = new FakeRedis();
const adapter = new BunRedisWorkflowAdapter({ client: redis });

await adapter.dispatch(envelope("job-x"));

// Simulate a crash after the atomic claim moved the id into processing but
// before the running-status write: the id sits in processing, the instance
// JSON still says "queued", and there is no dedicated status field.
redis.sortedSets.get("workflows:ready")?.delete("job-x");
const processing =
redis.sortedSets.get("workflows:processing") ?? new Map<string, number>();
processing.set("job-x", Date.now() - 10_000);
redis.sortedSets.set("workflows:processing", processing);
redis.strings.set("workflows:lease:job-x", "token");

const recovered = await adapter.recoverStalled(new Date(Date.now() - 5_000));
expect(recovered).toBe(1);
expect((await adapter.getInstance("job-x"))?.status).toBe("queued");
expect(
await adapter.claimNext(new Date("2026-05-24T09:00:00.000Z")),
).toMatchObject({ id: "job-x" });
});

test("does not re-enqueue a terminal job found in processing", async () => {
const redis = new FakeRedis();
const adapter = new BunRedisWorkflowAdapter({ client: redis });

await adapter.dispatch(envelope("job-done"));
await adapter.updateInstance("job-done", "complete", { output: 1 });

const processing =
redis.sortedSets.get("workflows:processing") ?? new Map<string, number>();
processing.set("job-done", Date.now() - 10_000);
redis.sortedSets.set("workflows:processing", processing);

const recovered = await adapter.recoverStalled(new Date(Date.now() - 5_000));
expect(recovered).toBe(0);
expect((await adapter.getInstance("job-done"))?.status).toBe("complete");
expect(redis.sortedSets.get("workflows:processing")?.has("job-done")).toBe(false);
});

test("expires terminal instance, step, and dead-letter hashes", async () => {
const redis = new FakeRedis();
const adapter = new BunRedisWorkflowAdapter({
client: redis,
retention: { completedTtlSeconds: 100, deadLetterTtlSeconds: 200 },
});

await adapter.dispatch(envelope("job-r"));
await adapter.saveStepResult("job-r", "step", 1);
await adapter.updateInstance("job-r", "complete", { output: 1 });
await adapter.recordDeadLetter(envelope("job-d"), new Error("boom"));

expect(redis.expirations.get("workflows:instance:job-r")).toBe(100);
expect(redis.expirations.get("workflows:step:job-r:step")).toBe(200);
expect(redis.expirations.get("workflows:dead:job-d")).toBe(200);
});

test("throws at construction for an unusable Redis client", () => {
expect(() => new BunRedisWorkflowAdapter({ client: {} })).toThrow(
/Unusable Redis client/,
);
});
});
Loading