Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
24 commits
Select commit Hold shift + click to select a range
468bfde
feat(workflow-engine): the contract of the workflow engine on the ledger
rami-hatoum Oct 4, 2026
2001985
docs(global): record the decision to run workflows on the ledger
rami-hatoum Oct 4, 2026
23d8c20
feat(ledger): read a stream after a version and share the decision loop
rami-hatoum Oct 4, 2026
7cf7118
feat(workflow-engine): revise the contract after its review
rami-hatoum Oct 4, 2026
72f9fa0
docs(workflow-engine): describe the revised contract and its invariants
rami-hatoum Oct 4, 2026
c20ee5b
docs(global): weigh alternatives and costs in the workflow engine dec…
rami-hatoum Oct 4, 2026
6e4ad3d
feat(ledger): share the load-decide-append loop as decisionLoop
rami-hatoum Oct 4, 2026
49c0c19
feat(operations): keep the settlement and call-result vocabulary once
rami-hatoum Oct 4, 2026
cd9641d
feat(orchestration): bound the cache of compiled expressions
rami-hatoum Oct 4, 2026
55a9d77
test(workflow-engine): scan for host timers and tell caches from cons…
rami-hatoum Oct 4, 2026
b5187ca
refactor(workflow-engine): move the DSL from the orchestration primitive
rami-hatoum Oct 4, 2026
01afc20
refactor(workflow-engine): take the functions a workflow calls from t…
rami-hatoum Oct 4, 2026
425d9b0
refactor(orchestration): take the limits and DslError from the workfl…
rami-hatoum Oct 4, 2026
8f18fca
test(workflow-engine): scan the DSL and the jq library the machine runs
rami-hatoum Oct 4, 2026
ee84e71
feat(workflow-engine): second revision of the contract after the audit
rami-hatoum Oct 4, 2026
ff608b4
docs(workflow-engine): describe the second revision of the contract
rami-hatoum Oct 4, 2026
932c098
docs(global): correct and extend the workflow engine decision
rami-hatoum Oct 4, 2026
261a0aa
fix(workflow-engine): refuse a decision of more than one event
rami-hatoum Oct 4, 2026
6d53462
docs(workflow-engine): mark the bundle entries transitional and name …
rami-hatoum Oct 4, 2026
0d2cd65
docs(global): record the nested-workflow error kind and the transitio…
rami-hatoum Oct 4, 2026
ddcf6fa
Merge remote-tracking branch 'origin/main' into feat/workflow-engine
rami-hatoum Oct 4, 2026
da3a51b
docs(global): let a package's exports map name its entry points
rami-hatoum Oct 4, 2026
ed468d9
docs(global): describe the hosted runtime by its constraints, not its…
rami-hatoum Oct 4, 2026
323e07b
docs(global): say where the image runs without naming the host
rami-hatoum Oct 4, 2026
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
4 changes: 2 additions & 2 deletions CLAUDE.md
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,7 @@

## What this is

`auto-brain` is the runtime for business brains, a pnpm monorepo. It is source-available under the Elastic License 2.0 (see `LICENSING.md`); never call it open source. Its container image runs self-hosted and in Auto's cloud hosting on Cloudflare Containers. Auto Studio (the management plane, with governance and observability) and the Cloudflare side live in the private `on.auto` repo, not here.
`auto-brain` is the runtime for business brains, a pnpm monorepo. It is source-available under the Elastic License 2.0 (see `LICENSING.md`); never call it open source. Its container image runs self-hosted and in Auto's cloud hosting. Auto Studio (the management plane, with governance and observability) and the hosting side live in the private `on.auto` repo, not here.

- `packages/server`: the Node.js server (`@beonauto/server`), plus everything that packages it into a container (`Dockerfile`, `Dockerfile.dockerignore`)
- `packages/api`: the API (`@beonauto/api`), a Hono app that answers every request, with problem documents, the `Origin` and `Host` checks, authentication and the operation routes
Expand Down Expand Up @@ -43,7 +43,7 @@ Run one package's gate with `pnpm turbo run lint typecheck test --filter @beonau
- 100% coverage per file is the gate. Never add coverage-ignore comments or coverage excludes; write the test.
- Do not write comments. Make the code read like English through names and ordering.
- Tests live next to the code as `*.test.ts`, test behaviour through the public interface, and prefer injected fakes over mocks.
- A package with more than about 12 source files groups them one level deep, in folders named after concepts that hold fewer than about 15 files each, with tests beside the code they test. Nothing new goes directly under `src`, and there are no barrel files: the entry points are `src/index.ts`, `src/testing/index.ts` and a package's commands (the server's `src/main.ts` and `dev.ts`, the identity package's `src/key-command.ts`).
- A package with more than about 12 source files groups them one level deep, in folders named after concepts that hold fewer than about 15 files each, with tests beside the code they test. Nothing new goes directly under `src`, and there are no barrel files: a package's entry points are the few paths its `exports` map names (`src/index.ts`, `src/testing/index.ts` and, where its README says why, a subpath) and its commands (the server's `src/main.ts` and `dev.ts`, the identity package's `src/key-command.ts`).
- Commits are conventional with a scope named after a package or primitive folder (`feat(server): ...`, `docs(inference): ...`); `global`, `deps`, `ci` and `release` are the other scopes.
- `pnpm check` must pass before you finish. Fixing a problem is a change, so rerun it.
- When something fails, assume your change broke it. What is on `main` passed the same gate.
85 changes: 85 additions & 0 deletions docs/decisions/0001-workflow-engine-on-the-ledger.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,85 @@
# 1. Run workflows on an engine on the ledger, not on Temporal

- Status: accepted
- Date: 2026-10-04

## Context

auto-brain's image runs self-hosted, as one Node process with one SQLite file, and in Auto's cloud hosting. There the engine must run in a single-threaded isolate for each run, woken by alarms that fire at least once and are dropped after a bounded number of failed retries, with bounded memory (about 128 MB) and bounded CPU for each wake-up, no long-lived process and no code generation; each isolate keeps its state in its own SQLite, whose rows take at most 2 MB, and the database the isolates share has no interactive transactions and binds at most 100 parameters in one statement. This record calls that the hosted runtime.

Workflow specs run on Temporal today, the wrong place for them:

- Temporal cannot run in the hosted runtime: a Temporal worker is a long-lived process, and the hosted runtime has none.
- Nobody asked for Temporal. People ask for workflows that wait, retry and survive a restart, and a self-hosted server must run a Temporal service beside it to get them.
- Tenant data is stored twice: Temporal's history holds each workflow's document, input, the outputs of its calls and its events, unencrypted, beside the brain's ledger.
- We already have the store: the ledger keeps every brain's events with Emmett on SQLite, through the sqlite3 driver self-hosted and through the drivers for the isolate's own and the shared SQLite in the hosted runtime.

We weighed three other engines. The hosted runtime's own workflow service runs only there, so self-hosted servers would need a second engine. Restate and Inngest are services of their own: a self-hosted server would run one beside it, as with Temporal, and each keeps step results in its own store, so tenant data would still be stored twice.

## Decision

A run is a decider in Emmett's workflow shape, on the ledger's own load-decide-append loop and conflict retry, `decisionLoop`, with a load that folds a snapshot and its tail:

- `decide(input, state)` says what an input changes and `evolve(state, event)` applies it. The inputs are a start, a timer fired, a call answered, an event received and a cancel request, each with the time it arrived, never earlier than the input before it.
- One stream per run. An applied input appends one event, with the version it was decided on as the expected version.
- The stream is a state-transition log, not classic event sourcing. Each event holds the change to the state as a JSON Patch, a receipt naming the input, the steps it ran and the outputs: arm or cancel a timer, start or cancel a call, settle. Replay applies patches and evaluates nothing, so a run started under one version of the interpreter loads under the next, and the receipt and steps keep the log readable.
- Every event and snapshot names its state format. Each event folds under its own format, the state is upcast where the format changes, formats never go back within a stream, and a committed corpus of every format must load as a test.
- Outputs are dispatched after the append, in stream order behind a watermark per run, and again on wake until they all succeed. Each is idempotent by its key: timer id, call key (execution, task reference, run) or execution id. Every started call has a deadline timer, so every call is answered.
- Deduplication lives in the run's state. Emmett is the store, never the engine: we use neither its workflow handler, which folds the whole stream for every input, nor its processors.
- A snapshot follows once the events since the last one take as many bytes as it did, and at least 1 MiB; only the latest is kept, in chunks of at most 1 MiB under the 2 MB row limit of the hosted runtime's SQLite. A run may take 100,000 inputs and write 512 MiB of history, both checked in `decide`.
- Four adapters sit behind small ports: run store, timers, executor and record store, with the watermark and per-run serialisation beside them.
- In the hosted runtime: one isolate per run, its stream in the isolate's SQLite and its timers on the isolate's alarm; one for the brain's record; one for the org's registry; a scheduled sweep wakes the runs the record says are overdue.
- Self-hosted, one server keeps the ledger and every run in one SQLite file, in one process; a second process on that file is unsupported. Timers go through the ledger's own SQLite driver or a separate file.
- A PostgreSQL adapter, later, assumes no 2 MB row limit and takes a lease per run for serialisation.
- The package `@beonauto/workflow-engine` is the workflow machine's contract, since the state is shaped by the workflow DSL, and the DSL lives in it. The orchestration primitive imports the DSL from it and keeps parsing, the primitive, the event and cancel operations and, for now, the Temporal runtime. The engine knows no brains, specs or primitives: the caller names the functions a workflow may call. Its `./dsl/*` and `./limits` entries are transitional: they keep Effect and the ledger out of the Temporal workflow bundle and go when Temporal goes.

Still open, due before the hosted adapters: how the hosted runtime executes a call that runs longer than the CPU one wake-up allows, as a step of a durable workflow service outside the isolate or through a queue to a container.

## Consequences

The main cost is rewriting the interpreter, 129 tests of it beside the DSL's 79, as a machine that steps from state to state instead of an async function Temporal replays. Its DSL, expressions and policy stay.

Tenants see three changes. A run may hold 4 MiB of data instead of 16. A repeated event no longer counts toward the events a run takes over its life. And a workflow that executes a workflow through a primitive name it computes fails with a `validation` error, once the executor rejects the call as `invalid_arguments`, where the interpreter raises a `configuration` error today. The written case, `primitive: orchestration`, is still refused at `create_spec` as forbidden; only a computed name reaches the executor. Both errors have status 400 and settle the execution as `invalid_input`; the visible difference is the error `type`, which a `catch.errors.with` filter matches.

We give up Temporal's durable timers, deduplicated delivery, replay, web UI and operator tools. We must build and keep correct:

- Timers that fire at least once; a fire of a timer no longer armed changes nothing.
- Deduplication in state, every key bounded, or snapshots grow with the run.
- The watermark: a crash between append and dispatch loses nothing.
- The sweep, for alarms that fire late or give up after their retries.
- Serialisation per run: an in-process lock in Node, the isolate's single thread in the hosted runtime.
- Settlement in two stores, the run's stream and the brain's record, each idempotent by execution id, the second retried; the ledger's port appends to one stream, so they are never one transaction.

Reads over the ledger and Studio replace Temporal's UI.

Nothing running on Temporal is migrated: nothing is in production, and the switch happens before a release.

Ended streams are kept; a deletion policy is a later decision.

Tenant data is stored once, and a workflow needs no service beyond the server.

## Plan

- The machine, on the contract and the DSL in `@beonauto/workflow-engine`.
- A conformance suite: the same workflow probes through a fake driver, the Node adapter and the hosted adapter, seeded from the server's `src/workflow-executions` tests.
- The adapters, a `cancel_execution` operation, and the workflow SDK's validators precompiled, since the hosted runtime allows no code generation.
- The cutover: `pnpm dev` without Temporal's dev server, and the README.
- Measurements before and after: the 15 recorded histories through the driver; inputs per second, and bytes per input against Temporal's 9.26 MiB for 40,000 inputs; timer lateness at p99; heap per live run; snapshot bytes per run.

## Evidence

Branch `spike/engine-node`:

- `spikes/node/results/replay.json`: the interpreter as it runs on Temporal replays 40,000 inputs in 4.6 s and retains up to 103 MiB; its log takes 9.26 MiB.
- `spikes/node/results/message-id-probe.json`: Emmett appends a message with an id it has seen as a new message.
- `spikes/node/results/executor-virtual.json`, `executor-real.json`: a result delivered four times settles once, a result after its timeout is ignored, a crash after the append is recovered on wake.
- `spikes/node/results/timers-precision.json`, `timers-recovery.json`: timers fire 3.7 ms late at p99 when idle; after a killed scheduler all 200 fire, none twice.
- `spikes/node/results/timers-two-processes.json`: two processes double-fire 227 of 300 timers unless each claims a timer first.
- `spikes/node/results/lost-write-repeat.json`: timers written through a second SQLite library to the ledger's file lost committed cancels in three runs of three.

The measurements in the hosted runtime are kept in the private repository:

- Folding 40,000 events cold takes 239 ms and holds 66 MiB; from a snapshot every 1,000 events, 9 ms and 1.6 MiB. The hosted runtime documents 30 s of CPU for each wake-up by default; its local runtime enforced no CPU limit at all: 35 s of CPU finished with no limit set, and 2 s with a limit of 50 ms.
- Alarms fire 5 ms late at p99; an alarm due while the local runtime was stopped fired 15.6 s late, when it was restarted; a sweep re-armed one that had given up.
- The record was written exactly once, or given up as intended, under every injected fault of the shared database and every crash; the shared database refuses eleven events in one append.
- The interpreter, the DSL policy, jq and both hosted ledger drivers run in the isolate; the workflow SDK's validators run once precompiled.
16 changes: 13 additions & 3 deletions packages/ledger/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,16 @@ This package is the event store behind the `Ledger` port of `@beonauto/operation
- A decision of more than eight events is a defect.
- Any other failure of the database is a defect, never a `Conflict` or a rejection.

## Building on the ledger's loop

Code that keeps its own streams, such as `@beonauto/workflow-engine`, uses the same pieces as the ledger itself rather than a copy of them:

- `sqliteEventStore(optionsOf)` opens the event store on any of Emmett's SQLite drivers, without the layer.
- `EventStore.read(stream, after)` gives the events after version `after` and the version of the whole stream, so a reader that holds a snapshot at version `after` reads only the tail. Emmett answers a read past the end of a stream with version 0; `read` answers with `after` instead.
- `eventAppenderOf(store)` encodes and appends events with an expected version, at most eight in one append, and fails with `VersionConflict` when another writer appended first.
- `retriedOnVersionConflict(attempt)` runs a load-decide-append attempt again after a version conflict, up to three more times, and then fails with `Conflict`.
- `decisionLoop(load, append, decider)` is the load-decide-append loop itself, the one `Ledger.execute` runs: it loads, decides, appends the decided events with the loaded version expected, retries with `retriedOnVersionConflict`, and answers with what the load gave, the events and the folded state. The ledger's load folds the whole stream; a caller with snapshots passes a load that folds a snapshot and its tail.

## Creating the layer

```ts
Expand All @@ -34,9 +44,9 @@ Each SQLite connection may cache up to 8 MiB of pages and maps none of the file

## Portability

The same ledger will run on Cloudflare D1, so it follows these rules:
The same ledger must also run on hosted SQLite databases that bind at most 100 parameters in one statement and offer no interactive transactions, so it follows these rules:

- It uses only the event store's own operations: read a stream, append with an expected version, migrate, close. It writes no SQL and registers no projections or consumers.
- It uses only the event store's own operations: read a stream, or its tail after a version, append with an expected version, migrate, close. It writes no SQL and registers no projections or consumers.
- It does not rely on transactions or rollback: each command makes at most one append.
- An append carries at most eight events. Emmett binds ten parameters for each event it inserts, and D1 accepts at most 100 bound parameters in one query.
- An append carries at most eight events. Emmett binds ten parameters for each event it inserts, and such a database binds at most 100 in one statement.
- Only `src/open-event-store.ts` knows which SQLite driver is in use and that the database is a file.
84 changes: 84 additions & 0 deletions packages/ledger/src/decision-loop.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,84 @@
import { Conflict } from '@beonauto/operations';
import { Effect, Result } from 'effect';
import { describe, expect, it } from 'vitest';

import { decisionLoop, VersionConflict, type DecisionLoop } from './index.ts';
import { outcomeOf } from './testing/open-ledger.ts';
import { tally, type Amounts } from './testing/tally.ts';

interface Loaded {
readonly state: number;
readonly version: number;
readonly loadedFrom: string;
}

interface Appended {
readonly stream: string;
readonly events: readonly unknown[];
readonly expectedVersion: number;
}

interface Counted {
readonly type: 'counted';
readonly by: number;
}

interface Loop {
readonly loop: DecisionLoop<Loaded, number, Amounts, Counted, 'conflict'>;
readonly appended: readonly Appended[];
}

function loopOver(loaded: Loaded, conflicts: number): Loop {
const appended: Appended[] = [];
let remaining = conflicts;
const loop = decisionLoop(
() => Effect.succeed(loaded),
(stream, events, expectedVersion) =>
Effect.suspend(() => {
remaining -= 1;
appended.push({ stream, events, expectedVersion });
return remaining >= 0 ? Effect.fail(new VersionConflict()) : Effect.void;
}),
tally,
);
return { loop, appended };
}

const snapshotted: Loaded = { state: 40, version: 7, loadedFrom: 'a snapshot and two events' };

describe('the decision loop over any store', () => {
it('decides on what its load gave, appends with that version expected, and gives the load back with the events', async () => {
const { loop, appended } = loopOver(snapshotted, 0);

expect(await Effect.runPromise(loop('run/1', [1, 1]))).toEqual({
loaded: snapshotted,
events: [
{ type: 'counted', by: 1 },
{ type: 'counted', by: 1 },
],
state: 42,
version: 9,
});
expect(appended).toEqual([
{
stream: 'run/1',
events: [
{ type: 'counted', by: 1 },
{ type: 'counted', by: 1 },
],
expectedVersion: 7,
},
]);
});

it('loads and decides again after a version conflict, and fails with Conflict after three more attempts', async () => {
const { loop, appended } = loopOver(snapshotted, 4);

expect(await outcomeOf(loop('run/1', [1]))).toEqual(
Result.fail(
new Conflict({ detail: 'The state changed while the command was decided', kind: 'concurrent_change' }),
),
);
expect(appended).toHaveLength(4);
});
});
42 changes: 42 additions & 0 deletions packages/ledger/src/event-store.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,42 @@
import { sqlite3EventStoreDriver } from '@event-driven-io/emmett-sqlite/sqlite3';
import { afterEach, describe, expect, it } from 'vitest';

import { sqliteEventStore, type EventStore } from './index.ts';

const opened: EventStore[] = [];

async function aStore(): Promise<EventStore> {
const store = sqliteEventStore(() => ({ driver: sqlite3EventStoreDriver, fileName: ':memory:' }));
opened.push(store);
await store.migrate();
return store;
}

afterEach(async () => {
await Promise.all(opened.splice(0).map((store) => store.close()));
});

const run = 'run/0199a3c4-7d2e-7c1a-9b3f-2f1e0d9c8b7a';

function numbered(from: number, count: number) {
return Array.from({ length: count }, (_, index) => ({ type: 'counted', data: { n: from + index } }));
}

describe('reading a stream after a version', () => {
it('gives the events after that version, and the version of the whole stream', async () => {
const store = await aStore();
await store.append(run, numbered(1, 3), 0);
await store.append(run, numbered(4, 2), 3);

expect(await store.read(run, 3)).toEqual({ version: 5, events: [{ n: 4 }, { n: 5 }] });
expect(await store.read(run)).toEqual({ version: 5, events: [1, 2, 3, 4, 5].map((n) => ({ n })) });
});

it('gives no events and the version it was asked after when nothing follows it', async () => {
const store = await aStore();
await store.append(run, numbered(1, 3), 0);

expect(await store.read(run, 3)).toEqual({ version: 3, events: [] });
expect(await store.read('run/nobody-wrote', 0)).toEqual({ version: 0, events: [] });
});
});
Loading
Loading