Skip to content

About

Durable, backend-agnostic workflows for TypeScript — one engine over Postgres, SQLite, MySQL, MongoDB, Redis, DynamoDB, or Cloudflare Durable Objects.

Topics

Resources

Contributing

Stars

14 stars

Watchers

0 watching

Forks

Repository files navigation

iterativeflow

core license

Durable, backend-agnostic workflows for TypeScript.

Write a flow as an ordinary async function; it survives process crashes, retries failed steps, sleeps for days, and resumes deterministically by replaying memoized steps. The same engine runs behind a four-port Backend interface — Postgres, SQLite, MySQL, MongoDB, Redis, DynamoDB, Cloudflare Durable Objects, or in-memory — resident or serverless. Published under the @iterativeflow/* scope.

import { createEngine, defineFlow } from "@iterativeflow/core";
import { createPgBackend, pgPool } from "@iterativeflow/postgres";
import { Pool } from "pg";
import { z } from "zod";

const Survey = z.object({ score: z.number() });

const onboard = defineFlow({
  name: "onboard",
  version: 1,
  signals: { survey: Survey }, // any Standard Schema: typed and validated
  run: async (ctx, input: { userId: string }): Promise<{ score: number }> => {
    await ctx.step("create-account", () => createAccount(input.userId));
    await ctx.sleep(3 * 24 * 60 * 60_000); // 3 days, durable
    const survey = await ctx.signal("survey"); // typed { score: number }
    return { score: survey.score };
  },
});

const engine = createEngine(createPgBackend(pgPool(new Pool())), [onboard]);
const stop = engine.run(); // resident worker loop; returns a stop fn

const handle = await engine.submit(onboard, { userId: "u_1" });
// 3 days later, from a webhook:
await engine.signal(handle, "survey", { score: 9 });
const { output } = await engine.result(handle, { timeoutMs: 30_000, output: Survey }); // { score: 9 }
await stop();

That run lives in your backend for three days. Workers can crash, deploys can roll, the process can be killed and restarted — when the timer fires, the flow resumes from where it left off, replaying the memoized create-account step instead of re-running it.

  • Steps memoized by (runId, cursor) — ctx.step(label, fn), the unit of at-least-once execution
  • Sleeps and external signals lasting days — ctx.sleep(ms) / ctx.signal(name)
  • ctx.invoke(child, input) for child flows, and ctx.invoke([…]) for parallel fan-out + join
  • Typed contracts — signal payloads and outputs are typed by Standard Schemas that also validate them, and engine.signal checks signal names and payloads at compile time
  • At-least-once via a transactional outbox committed with each step; a reconciler re-drives anything stranded by a crash
  • Serverless or resident — engine.run() for a worker loop, or serverlessTick for one bounded cycle per Lambda/Vercel/Cron invocation
  • Structured concurrency — a run that terminates without success cancels its descendants

How it works

One engine, four ports, one datastore. The engine runs your flow, and each step commits its result and its side effects — child spawns, enqueues, timers — in a single transaction (the outbox), so a crash never leaves state and its follow-up out of sync. On restart the engine replays the flow from the start; memoized steps short-circuit, so only un-run work executes again.

flowchart TB
  flow["your flow — an async function"] --> engine
  engine["engine · replay · retry · outbox · reconcile"] --> store & queue & timer & wakeup
  subgraph backend ["Backend — four ports"]
    direction LR
    store["Store<br/>runs · steps · signals · crons"]
    queue["Queue<br/>claim · lease"]
    timer["Timer<br/>sleeps · deadlines"]
    wakeup["Wakeup<br/>completion nudge"]
  end
  backend --> db[("one datastore:<br/>Postgres · SQLite · MySQL · Mongo · Redis · DynamoDB · DO")]
Loading

Every backend implements those four ports and passes the same conformance suites — the executable definition of a correct backend — so swapping one for another changes only the create*Backend(...) call.

Packages

Package What it is
@iterativeflow/core The engine — defineFlow, ctx, createEngine, replay, outbox. Pair it with one backend.
@iterativeflow/memory In-memory backend — tests, dev, single process.
@iterativeflow/postgres Postgres backend — transactional outbox, optional LISTEN/NOTIFY push, serverless-friendly. The default.
@iterativeflow/sqlite SQLite backend (@libsql/client) — embedded, Turso, or a single-node service.
@iterativeflow/mysql MySQL/InnoDB backend — FOR UPDATE SKIP LOCKED claims.
@iterativeflow/mongodb MongoDB backend — multi-document transactional outbox (replica set required).
@iterativeflow/redis Redis backend — Lua-scripted outbox, single-node.
@iterativeflow/dynamodb DynamoDB backend — single-table, TransactWriteItems outbox, serverless AWS.
@iterativeflow/durable-objects Run the engine inside a Cloudflare Durable Object on its built-in SQLite — no external DB.
@iterativeflow/webhooks Inbound edge — verify a signed provider webhook and deliver it as a durable signal a parked flow awaits.
@iterativeflow/dashboard Dependency-free ops UI (runs list, detail, cancel/retry/signal) as a fetch handler.
@iterativeflow/conformance The shared suites every backend must pass — the executable definition of a correct backend.

Which backend?

  • memory — tests, examples, a single-process app.
  • postgres — the default: strong consistency, push completion via LISTEN/NOTIFY, works on serverless/pooled Postgres.
  • sqlite — embedded or edge, one file or Turso/libsql; zero server.
  • durable-objects — the edge, strongly consistent per object, no external database.
  • redis — low-latency, single-node (Valkey/Dragonfly work too).
  • mysql / mongodb / dynamodb — when that store is already your stack.

See patterns/backends for each backend's setup, tuning, and gotchas.

Install & quick start

npm install @iterativeflow/core @iterativeflow/memory
import { createEngine, defineFlow } from "@iterativeflow/core";
import { createMemoryBackend } from "@iterativeflow/memory";

const double = defineFlow<{ x: number }, number>({
  name: "double",
  version: 1,
  run: async (ctx, input) => ctx.step("d", () => input.x * 2),
});

const engine = createEngine(createMemoryBackend(), [double]);
const stop = engine.run();
const handle = await engine.submit(double, { x: 21 });
const res = await engine.result(handle, { timeoutMs: 5_000 });
await stop();
// res.output === 42

Swap createMemoryBackend() for createPgBackend(...), createSqliteBackend(...), etc. — nothing else changes. Each backend's guide covers its connection setup and applySchema.

The flow context

Inside run, ctx is the durable surface; every call is a checkpoint:

  • ctx.step(label, fn) — run fn once; memoize its result and replay it on every later attempt.
  • ctx.sleep(ms) — suspend and resume after a durable timer.
  • ctx.signal(name) — park until an external engine.signal(runId, name, payload) (or @iterativeflow/webhooks) delivers a typed payload.
  • ctx.invoke(flow, input) — run a child flow and await its output; pass an array to fan out in parallel and join outputs in order.

Replay assumes the body is stable: each step memo records a shape fingerprint, and a redeploy that reorders the body triggers the flow's driftPolicy (park or fail) instead of running the wrong step. Keep step order and labels stable across deploys; bump version when the body changes meaningfully.

Builder (fluent, typed)

Prefer a chain? builder compiles to the same ctx.step — each .step result is added to a typed accumulator (acc.account, acc.survey) that later steps and the output projection can read. Sleeps, signals, and invokes happen through ctx inside a step (there are no separate chain nodes for them):

import { builder } from "@iterativeflow/core";

const onboard = builder<{ userId: string }>("onboard", 1)
  .step("account", (acc) => createAccount(acc.input.userId))
  .step("survey", async (_acc, ctx) => {
    await ctx.sleep(3 * 24 * 60 * 60_000); // 3 days
    return (await ctx.signal("survey")) as { score: number };
  })
  .output((acc) => ({ score: acc.survey.score }));

Testing a flow that sleeps for days

@iterativeflow/core/testing runs the real engine on a virtual clock, so the three-day sleep above resolves in a millisecond. Time moves only when you move it, and every jump lands exactly on the next durable deadline — no polling, no fake timers, no guessing how far to advance.

import { createTestHarness } from "@iterativeflow/core/testing";
import { createMemoryBackend } from "@iterativeflow/memory";

const t = createTestHarness(createMemoryBackend(), [onboard]);

const handle = await t.engine.submit(onboard, { userId: "u_1" });
await t.advanceToNextWake(); // runs up to the sleep, then jumps the 3 days
await t.engine.signal(handle, "survey", { score: 9 });

expect(await t.settle(handle)).toMatchObject({ status: "done", output: { score: 9 } });

settle drives a run to its terminal outcome, jumping every sleep and retry backoff on the way; if the run parks on something nothing will deliver, it throws naming what it waited on rather than hanging. The harness takes any Backend, so the same test can run against a real database.

Docs

  • Backends — per-backend setup, production tuning, and gotchas.
  • Deployment — execution models, connection pooling, clocks and leases, scaling to zero, sharding, and rolling deploys.

License

MIT

About

Durable, backend-agnostic workflows for TypeScript — one engine over Postgres, SQLite, MySQL, MongoDB, Redis, DynamoDB, or Cloudflare Durable Objects.

Topics

Resources

Contributing

Stars

14 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages