Skip to content

feat(workflow-engine): the contract of a workflow engine on the ledger - #85

Merged
rami-hatoum merged 24 commits into
mainfrom
feat/workflow-engine
Oct 4, 2026
Merged

rami-hatoum merged 24 commits into
mainfrom
feat/workflow-engine

Conversation

@rami-hatoum

@rami-hatoum rami-hatoum commented Oct 4, 2026 •

Copy link
Copy Markdown
Contributor

What this adds

The contract of a workflow engine that runs on the ledger, so that workflows run the same wherever the ledger runs: a container with one SQLite file, or a hosted runtime of single-threaded, alarm-driven isolates with bounded memory. No behaviour changes for anything running today: the Temporal-based engine keeps serving workflows until the new one is proven beside it.

  • packages/workflow-engine, the contract: the machine's inputs (started, timer fired, call answered, event received, cancel requested), its events (each records the input's receipt, the change to the run's state as a strict JSON Patch, the steps taken and the outputs caused), its outputs (arm or cancel a timer, start or cancel a call, settle), the run's state and snapshots with a format version, the idempotency keys for timers, calls and external events, the admission rule (applied, stale, not started), the dispatch watermark, bounds per run, and the ports an adapter implements (run store, record store, timers, executor, serialiser, reporter). The README states 33 invariants, each marked as the engine's or an adapter's obligation. No source in the package or in what the machine imports uses a clock, randomness, a Node API or code generation; a test scans for them.
  • The workflow language moves into that package (src/dsl: expressions with the bounded jq, policy, durations, tasks), as renames. The orchestration primitive imports it from there and keeps the parsing, the summary and the Temporal runtime.
  • packages/ledger: EventStore.read(stream, after) and decisionLoop, the load-decide-append loop both the ledger and the engine run.
  • packages/operations: one home for the settlement and call-result vocabulary that was defined three times.
  • primitives/orchestration: the cache of compiled expressions is bounded; the limits and the error type come from the engine.
  • docs/decisions/0001-workflow-engine-on-the-ledger.md: why Temporal is replaced, what the engine is, what it costs, what is given up, and the evidence.
  • CLAUDE.md: the layout rule lets a package's exports map name its entry points, as the repository already does.

How it was checked

  • Spikes measured the ground before the design: replay cost of the current interpreter, what a step-function rewrite touches, executor idempotency under duplicate delivery, timers and the single-writer question, portability of the packages to an isolate runtime, one alarm driving many timers, a cold fold of 40,000 events, settlement across stores.
  • The contract was reviewed independently twice (a design review, then a full audit against the decision, the alternatives, correctness traced through a run's lifecycle, DRY, domain-agnosticism, patterns, portability, the plan and the record), each followed by a revision; the re-check accepted it for the machine to be built on.
  • pnpm check passes with the engine at 100% coverage and the orchestration suite and workflow bundle unchanged in behaviour after the moves.

What follows

The machine itself (the interpreter rewritten as the step function, with every existing test and the replay corpus carried over), then the Node adapter and the removal of Temporal.

🤖 Generated with Claude Code

rami-hatoum and others added 24 commits October 4, 2026 18:13
@beonauto/workflow-engine holds the contract between the machine that
runs a workflow and the adapters that store, time and execute for it,
grouped by the engine's layers: machine, run log, timers, inbox,
executor, dispatch, serialisation, settlement and engine. It knows
workflows, the inputs of a run and an executor port, and nothing of
brains, specs or models.

- Inputs: started, timer fired, call answered, event received and cancel
  requested, each with the execution id and the time it arrived.
- Events: one kind, input_applied, with the input's receipt, the change
  to the run's state as a JSON Patch, so evolve evaluates nothing, and
  the outputs it caused: arm or cancel a timer, start or cancel a call,
  settle.
- State: plain JSON with a schema, so a snapshot is the state as it is;
  snapshots are due after 1,000 inputs or 1 MiB of events and stored in
  chunks of at most 65,536 UTF-16 code units, cut between characters.
- Keys: timer ids, call keys of execution, reference and run, event ids
  and the execution id, with isStale saying which inputs a run has
  already taken.
- Ports: run store, timers, executor, record store, dispatch watermark
  and run serialiser, each answering with an Effect; outputsAbove gives
  what a wake dispatches.

The README states the invariants as sentences a reviewer can check, the
limits (a run holds 4 MiB, an event 1.5 MiB, an input runs 100 tasks)
and the open design points. A test checks that no source of the package
uses a Node-only API, generates code or imports Temporal. The machine is
the next step; the orchestration primitive still runs on Temporal.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
docs/decisions/0001-workflow-engine-on-the-ledger.md is the first
decision record: why workflows move from Temporal to an engine on the
ledger, the shape of that engine, what we give up and must keep
correct, and the measurements of the two spikes it rests on, cited by
branch and file.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
The workflow engine keeps one stream per run and loads it from its
latest snapshot, so it must read only the events after a version, and
it must decide and append the way the ledger does rather than through a
copy. EventStore.read takes an optional version to read after; Emmett
answers a read past the end of a stream with version 0, so read
answers with that version instead. The retries after a version conflict
move into retriedOnVersionConflict, which the ledger's own execute now
uses, and the index exports it with VersionConflict, the event store,
its appender and its codec, and sqliteEventStore, which opens the store
on any of Emmett's SQLite drivers without the layer.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
- Values a run holds live once in a value table and frames refer to
  them by id, so a patch adds an answer once however many places it
  lands; events are measured in UTF-8 bytes against 1.5 MiB.
- Admission names a stale reason, answers not_started for an input to
  a run that has not started, and dies on another execution or a
  second start with another document.
- Every event and snapshot names its state format; patches apply
  strictly and the fold decodes the state once, from the snapshot and
  the tail the run store returns.
- Dispatch runs in stream order and the watermark stops before the
  first event whose output failed; cancels of unseen keys tombstone.
- Limits gain 100,000 inputs, 512 MiB of history and the expression
  work budget; snapshots are due at max(1 MiB, the last snapshot's
  size) and chunked at 1 MiB of UTF-8.
- Events record the receipt with an answer's status or an event's
  type, and the steps an input ran; `at` is clamped to never go back.
- The run store's append fails with the ledger's VersionConflict, and
  the engine's submit with the ledger's Conflict.
- Purity and no-cache scans cover the machine and the run log.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
The README describes the run's log as a state-transition log, the value
table and how it bounds a patch, submission outcomes, receipts and
tombstones, the dispatch order, the sweep over live runs, state formats,
the new bounds and snapshot rules, Node's serialisation, and invariants
22 to 31, each marked as the engine's or an adapter's.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…ision

The record now names the run's stream a state-transition log and says
why, weighs Cloudflare Workflows, Restate and Inngest, names the
interpreter rewrite as the main cost, and states the snapshot rules,
bounds, Node serialisation and timers, PostgreSQL assumptions,
migration, retention and what replaces Temporal's UI. The 30 s CPU
figure is Cloudflare's documentation: local workerd ran 35 s of CPU
without a limit (spikes/cloudflare/results/fold.json, cpuLimits).

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
decisionLoop(load, append, decider) is the loop Ledger.execute runs,
now exported so a caller with its own load, such as the workflow
engine's load from a snapshot and its tail, runs the same loop and the
same retry instead of a copy. It answers with what the load gave, the
decided events and the folded state.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
SettlementSchema and CallResultSchema, with invalidArguments, now live
in @beonauto/operations beside Outcome. The workflow engine's run
outputs, inputs and record store use them instead of copies; the
record store's settlement in @beonauto/specs and the interpreter's
RunSettlement are derived from them. The interpreter's call result and
the Temporal worker's spec result keep their own shapes until Temporal
goes.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
The module-level map of compiled expressions grew with every expression
a process ever compiled. Under Temporal each workflow had an isolate of
its own; on the workflow engine one process serves every tenant. The
cache now holds at most 262,144 characters of expression source and
lets go of the expression used longest ago. A compiled expression
measured 22 to 34 bytes of heap per source character, so the cache
holds at most about 9 MiB, and compiling one again took 10 to 150 us
(1,000 expressions of 1 to 100 terms, parsed and validated).

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…tants

The portability scan also refuses setTimeout and setInterval and the
JavaScript Temporal global. The cache scan now flags a Map or Set built
empty at module level, a cache, and lets through a constant collection
of literals and a WeakMap, whose entries go with the values they
describe, ahead of the DSL joining the package.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
The DSL the machine needs moves to packages/workflow-engine/src/dsl
with its tests, as renames: JSON helpers, expressions with the jq work
budget and the bounded cache, durations and their limits, tasks, and
the policy with its checks and nesting rules. Eighteen of the files are
unchanged. retained-size, the interpreter's estimate of held data under
Temporal, moves beside the interpreter.

The orchestration primitive imports the DSL through the engine's new
dsl/* entry, so the bundled workflow code takes neither Effect nor the
ledger from the engine's main entry; the replay of the recorded
histories against a freshly built bundle passes. The jq dependency and
its pnpm patch, registered by package and version, follow the code. The
two nesting and fork tests that start a workflow through the
interpreter stay with it, in refused-documents.test.ts. The workflow
bundle's inputs, and the image's bundle and final stages, now include
packages/workflow-engine.

The pre-commit lint refuses a commit whose moved files do not resolve,
so the renames and the import changes are one commit.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…he caller

The DSL policy knew execute_spec, its arguments, the orchestration
primitive and the specs of a brain. policyOf(functions) now takes the
functions a workflow may call, each with the checks of its arguments,
and the words that explain where a workflow reaches the world and how
it starts. The orchestration primitive gives execute_spec and its rules
in src/document/workflow-functions.ts, with their tests, so the engine
package knows no brains, specs or primitives. The messages a document
gets are the same as before.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…ow engine

The limits the interpreter repeated verbatim, tasks without waiting,
expression work, work in one activation, tasks in one activation and
the four limits on events, now live once in the engine's limits, which
the package exports as @beonauto/workflow-engine/limits, a module of
plain numbers that the bundled workflow code can import. DslError is
the engine's DslErrorSchema type; the interpreter imports it as a type.
The limits that differ under Temporal, the 16 MiB a workflow holds and
its history, stay with the interpreter.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
The purity scan now covers the DSL beside the machine and the run log,
and a scan of jq-ts's distribution, the one library the machine runs
expressions with, finds no Node-only API, code generation, host timer,
clock read or random source. jq reaches the host's time zone only
through localtime and strflocaltime, the two builtins it builds with
local time, and the DSL refuses both; the test checks both facts.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
- History bytes live in the state. withHistoryBytes adds the event's
  own replace of /historyBytes, solved for the size of the event that
  carries it, so decide can bound the history as it bounds inputs, and
  a load dies when the state's count is not the bytes of its events.
- The whole-state decode refuses members the format does not know.
- Each event is folded under its own state format, the state is
  upcast where the format changes, formats never go back within a
  stream, and corpus/format-1.json, a committed stream and snapshot,
  must load as a test.
- A call frame holds the opaque function, its arguments as a value id
  and a label; the timer purpose call_deadline pairs every started
  call with a deadline; start receipts tell a call started again from
  one still running or answered again.
- Values are held while reachable: withReachableValuesOnly sweeps the
  value table after a decision, and a property test over 400 states
  checks heldBytes against the bytes the frames, context and input
  reach.
- The record notes each live run's next due time and whether its
  dispatch fell behind, so a sweep touches only overdue runs; troubling
  settle receipts go to a RunReporter.
- A second start with the same document and another input dies.
- A fired timer's input takes at least its due time.
- Snapshots are due by bytes alone and chunked with encodeInto.
- runLoopOf runs the ledger's decisionLoop with a load from a snapshot
  and its tail.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
The README now covers the entries for the DSL and the limits, the DSL
taking its call functions from the caller, the bounded expression
cache, history bytes and how they are counted, the format rule with
its corpus in place of upcasters for patches, held values by
reachability, the opaque call frame, the call deadline and the start
receipts, where tombstones matter, due runs with the sweep's cadence
and grace, troubling receipts, the API's answers for a repeated or late
event, and invariants 32 to 35. Invariant 14 no longer offers one
transaction, which the ledger's port does not.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Two figures said more than was measured: the 15.6 s late alarm was one
due while wrangler dev was stopped, firing when it restarted, and local
workerd enforced no CPU limit at all, not even cpu_ms = 50
(spikes/cloudflare/results/timers.json and fold.json). Settlement is
two writes, never one transaction, since the ledger's port appends to
one stream. The 4 MiB a run holds, and repeated events no longer
counting toward its life, move under Consequences as changes tenants
see. The record now states the format rule with its corpus, the call
deadline, the sweep of overdue runs, the DSL living in the engine
package, the test count of the rewrite (129 interpreter, 79 DSL), a
plan section, and the Cloudflare executor for long calls as an open
decision due before the Cloudflare adapters.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
runLoopOf appended each decided event as a separate write with
expectedVersion + index, so a decider that returned two events would
have lost the atomicity of an input without a sound. A decision is one
event (invariant 1); the loop now dies with SplitDecision on more, and
appends nothing.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…the first test

The dsl/* and limits entries are there only to keep Effect and the
ledger out of the Temporal workflow bundle, and go when Temporal goes.
Invariant 5 names SplitDecision. The nested-workflow error kind leaves
the open points for the decision record, and a section on the machine
names invariant 33, every open call answered, as its first test.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
…nal entries

A workflow that executes a workflow through a computed primitive name
fails with a validation error once the executor rejects it, where the
interpreter raises a configuration error today. The written case is
still refused at create_spec as forbidden; both errors have status 400
and settle as invalid_input; the visible difference is the error type a
catch.errors.with filter matches. The engine's dsl/* and limits entries
are transitional and go when Temporal goes.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
The layout rule named two files as the only entry points, while several
packages already export subpaths that their READMEs explain. The rule
now says what the repository does: the exports map names the few entry
points, and a subpath needs a reason in the README.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
… vendor

Where Auto's cloud hosting runs the image is internal. The decision
record now states once what that runtime allows, a single-threaded
isolate per run woken by at-least-once alarms that are dropped after
bounded retries, about 128 MB of memory, bounded CPU per wake-up, no
long-lived process, no code generation, its own SQLite with 2 MB rows,
and a shared database without interactive transactions that binds at
most 100 parameters, and reuses it everywhere the vendor was named. The
layout is one isolate per run, one for the brain's record and one for
the org's registry. Every number stays; the hosted-runtime measurements
are cited as kept in the private repository. The engine and ledger
READMEs and two test names follow.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Where Auto's cloud hosting runs is internal to the private repo. This
repository says only that the image runs self-hosted and in Auto's
cloud hosting.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
@rami-hatoum
rami-hatoum added this pull request to the merge queue Oct 4, 2026
Merged via the queue into main with commit a862223 Oct 4, 2026
10 checks passed
@rami-hatoum
rami-hatoum deleted the feat/workflow-engine branch October 4, 2026 21:00
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant