feat(workflow-engine): the contract of a workflow engine on the ledger - #85
Merged
Merged
Conversation
@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>
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
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.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)anddecisionLoop, 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'sexportsmap name its entry points, as the repository already does.How it was checked
pnpm checkpasses 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