Architecture¶
Flowli is a library with four layers and five processes. This page explains what each one owns, and where each guarantee comes from.
The shape¶
flowchart TB
subgraph app[Your application]
W[workflow functions]
R[Registry]
end
subgraph runtime[flowli.runtime]
E[Engine]
C[Context]
WK[Worker]
SW[Sweeper]
RT[Retention]
CO[Consumer]
end
subgraph domain[flowli.domain]
M[objects: Execution, Frame, Task, Message]
P[ports: Protocol definitions]
end
subgraph adapters[flowli.adapters]
CB[CairnBackend]
MB[MemoryBackend]
PR[WorkflowProjection]
end
B[(one bucket<br/>filesystem, S3, Azure, GCS)]
API[flowli.api<br/>HTTP service] --> E
WEB[web/<br/>operator interface] --> API
app --> runtime --> domain
runtime --> adapters
adapters -.implements.-> P
CB --> B
PR --> B
The layers¶
Layer |
Content |
Rule |
|---|---|---|
|
the objects, the errors, and the ports as |
pure. No I/O. It imports one thing from CairnDB: the |
|
the implementations of every port |
one module per backend: |
|
|
the protocols of the specifications, step by step. |
|
|
written against the |
|
the HTTP service |
holds no state: reads come from the projection, writes go through the engine. |
|
the |
builds the backend and the engine, then runs one job. |
Three modules sit beside the layers: flowli.codec renders objects to
JSON-compatible data and back, flowli.log configures structlog, and
flowli.evidence keeps the attempt log.
The ports¶
The domain defines each primitive as a Protocol. Each port maps to exactly
one CairnDB construct.
Port |
Purpose |
CairnDB construct |
|---|---|---|
|
the trace and the memo of one execution |
a named log |
|
the lifecycle entries of all executions |
the named log |
|
one owner per execution |
a lease |
|
an idempotent start |
a claim, which is put-if-absent |
|
the tasks |
objects plus one lease per task |
|
the messages |
one log per channel, plus wait markers |
|
the future resumes |
objects listed by prefix |
|
the |
one object, put-if-absent and CAS |
|
the folded finished executions |
one object, put-if-absent |
|
the attempt logs and the attachments |
objects listed by prefix |
CairnBackend implements all of them over one CairnDB instance.
MemoryBackend implements all of them in memory, for the tests.
backend = CairnBackend.configure({"storage": {"type": "s3", "bucket": "my-bucket"}})
engine = Engine(backend.ports, Site.local("w-1"), registry=registry)
The key layout¶
Every key of the engine starts with wf/. Every log starts with wf.
Key or log |
Content |
|---|---|
|
the control log |
|
the journal of one execution |
|
the messages of one channel |
|
the execution record |
|
the ownership lease |
|
the winner of a dispatch key |
|
the tasks, their leases and their enqueue markers |
|
a wait marker |
|
a timer |
|
a folded execution |
|
the attempt logs and the attachments |
One bucket is one tenant. There is no tenant field.
The processes¶
Process |
Command |
Needed when |
How many |
|---|---|---|---|
worker |
|
always |
as many as the throughput asks for |
sweeper |
|
always |
one is enough; two do no harm |
retention |
|
you keep the bucket small |
one, as a cron job |
HTTP service |
|
people or remote clients need access |
two or more, behind a load balancer |
consumer |
|
a workflow delegates work outside the engine |
one per pool of agents |
Each process reads the same bucket. No process talks to another process.
The worker¶
The worker repeats one procedure: dequeue a task, acquire the lease of its
execution, build the memo table, replay the function, append what happened, and
release the lease. It renews the task lease and the execution lease while it
works. On LeaseLost it stops at once and does not ack the task.
The sweeper¶
The sweeper has no lease. Each of its actions is idempotent, so two sweepers can run at the same time. It:
fires the timers that are due,
enqueues a resume for each execution whose lease expired,
enqueues the
STARTtask again for a start that was lost,repairs the control log when a worker crashed between two appends,
clears the wait markers that no frame waits for.
Nothing in a bucket fires at a time by itself
A ctx.sleep, a receive timeout and a retry delay all need the sweeper. An
installation without a sweeper never wakes a suspended execution.
The HTTP service¶
The service puts the engine on a network. It holds no state of its own: every read comes from the projection or from a port, and every write goes through the engine. The actor of a write comes from the access token, never from the body.
It serves three planes:
Plane |
For |
Content |
|---|---|---|
catalog |
the interface |
the workflows of this process and their argument schemas |
control |
operators, the interface |
executions, journals, frames, reviews, queues |
worker |
a consumer that cannot reach the bucket |
dequeue, renew, ack, nack, deliver |
The life of an execution¶
sequenceDiagram
participant O as Operator
participant E as Engine
participant B as Bucket
participant W1 as Worker 1
participant W2 as Worker 2
O->>E: start(invoice_approval, "INV-7")
E->>B: write record, announce created, enqueue START
W1->>B: dequeue START, acquire lease (epoch 1)
W1->>B: append execution.started, frame.started, frame.completed
W1->>B: append frame.suspended, execution.suspended
W1->>B: release lease
Note over W1,B: the execution holds no worker
O->>E: signal(eid, "payments", ...)
E->>B: send message, enqueue RESUME
W2->>B: dequeue RESUME, acquire lease (epoch 2)
W2->>B: read journal, build memo table
Note over W2: replay: each finished frame returns its memo
W2->>B: append frame.fulfilled, execution.completed
The worker that finishes an execution is rarely the worker that started it. Nothing in the process holds the state of an execution between two tasks.
Where each guarantee comes from¶
Guarantee |
Source |
|---|---|
A journal entry is durable and ordered when the append returns. |
the log append of CairnDB |
The memo of a frame never changes once it is written. |
the memo rule, plus the immutability of a log |
One worker owns one execution at a time. |
the lease |
A fenced worker gets |
the lease epoch, plus CAS |
One task is dequeued by one worker at a time. |
a lease on the task key |
Two starts with one dispatch key produce one execution. |
the claim |
One channel has a total order of messages. |
one log per channel |
Not guaranteed:
the order between two channels,
the order between the journal and the control log,
exactly-once external effects of a step.
The last one is the rule you design around: a step must be idempotent.
The boundary of a delegate¶
A DELEGATE task leaves the engine. The consumer is a process that you own:
an agent runner, a service, or a group of people behind a form.
flowchart LR
F["frame: delegate(ctx, 'agents', payload)"] -->|enqueue| Q[(queue 'agents')]
Q -->|dequeue| K[consumer]
K -->|deliver on the reply channel| CH[(reply channel)]
CH -->|resume| F
Rules of that boundary:
A consumer never appends to a journal, and never writes an execution record. Its only writes are the reply message, the evidence, and the state of its own task lease.
A consumer answers one time per task. The receive frame consumes the first message only.
The consumer delivers first and acks second. The reverse order can lose the answer.
A consumer that reaches the bucket uses the ports. A consumer that cannot uses the HTTP worker plane, which is the same protocol with one more hop.
flowli.runtime.Consumer implements that loop,
minus the work. The work is a Handler that you write.
What you can replace¶
You want |
You write |
|---|---|
another storage backend |
an object with the ten ports. |
another pattern |
a function that takes a |
another consumer |
a |
another interface |
an HTTP client of the control plane. The web application is one. |
another identity provider |
an |
What the engine does not do¶
It does not run a server, a scheduler daemon, or a broker.
It does not check that a step is idempotent.
It does not check what an actor is permitted to do. The HTTP service does.
It does not know about reviews, agents, boards or gates. Patterns and applications build those on the
ContextAPI.It has no way to continue an execution as a fresh one, so a journal grows for as long as the execution lives. Anything endless is a process, not a workflow.