API reference

This page is generated from the docstrings of the code by sphinx.ext.autodoc. The specifications give the rules behind each signature.

flowli.runtime

The runtime: what you call, and what runs your workflows.

Engine

The operator API. One engine holds the ports, the site and the registry.

class flowli.runtime.Engine(ports: Ports, site: Site, *, clock: Callable[[], Timestamp] = Timestamp.now, config: EngineConfig | None = None, registry: Registry | None = None)[source]

Bases: object

workflow(name: str, version: str) → Callable[[Callable[[...], Any]], Callable[[...], Any]][source]
provenance(actor: Actor, *, workflow: str = '-', version: str = '-', frame_kind: str = 'root', frame_name: str = 'engine', site: Site | None = None, attempt: int = 1) → Provenance[source]
async start(workflow_fn: Callable[[...], Any], /, *args: Any, key: str | None = None, queue: str | None = None, by: Actor, **kwargs: Any) → UUID[source]

Create an execution and enqueue its START task. Return the eid.

async start_result(workflow_fn: Callable[[...], Any], /, *args: Any, key: str | None = None, queue: str | None = None, by: Actor, **kwargs: Any) → StartResult[source]

start, with whether the dispatch key already had a winner.

A caller that reports back to a person needs that fact; start alone cannot tell a fresh execution from a converged one.

async start_child(execution: Execution, queue: str) → None[source]

ChildStarter. Idempotent on execution.eid.

async create_execution(execution: Execution) → bool[source]

Write the record, announce execution.created, enqueue START. Idempotent.

async enqueue(task: Task, *, repair: bool = False) → bool[source]

repair=True also puts back a task whose enqueue marker outlived it.

Repair paths only (the sweeper): it costs one listing more. See Queue.ensure.

async enqueue_resume(eid: UUID, cause: str, reason: str, by: Provenance, *, repair: bool = False) → bool[source]
async signal(eid: UUID, channel: str, payload: Any, *, by: Actor, correlation: str | None = None) → int[source]

Send on the execution-scoped channel, then enqueue a RESUME task.

async deliver(eid: UUID, channel: str, payload: Any, *, by: Actor, correlation: str | None = None, reason: str | None = None) → int[source]

Send on a channel given by its full name, then enqueue a RESUME task for eid.

async broadcast(channel: str, payload: Any, *, by: Actor) → int[source]

Send on a global channel, then enqueue one RESUME task per waiter.

async notify_parent(execution: Execution, body: dict[str, Any], by: Provenance) → None[source]

Procedure 5: tell the parent that this child reached a terminal entry.

async cancel(eid: UUID, *, by: Actor) → None[source]
async migrate(eid: UUID, version: str, *, by: Actor) → Execution[source]

Point a live execution at another workflow version. See 05-protocols.md section 10.

The record is replaced with CAS. A journal and control-log entry record who did it. A RESUME task is enqueued so a worker replays with the new code. Old memos stay valid when their frame ids and digests match.

async execution(eid: UUID) → Execution[source]
async journal(eid: UUID) → list[Sequenced[Entry]][source]

The live journal, or the archived one when retention folded it.

async archive(eid: UUID) → dict[str, Any] | None[source]
async memo(eid: UUID) → MemoTable[source]
async status(eid: UUID) → ExecutionStatus[source]

Derived from the last execution.* entry of the journal.

workflow_ref(execution: Execution) → WorkflowRef[source]
worker(queues: list[str] | None = None, worker_id: str | None = None) → Any[source]
sweeper(**kwargs: Any) → Any[source]
retention(**kwargs: Any) → Any[source]
property reviews: Any
class flowli.runtime.EngineConfig(exec_ttl: 'float' = 120.0, task_ttl: 'float' = 60.0, poll_interval: 'float' = 1.0, default_queue: 'str' = 'default', code_ref: 'str | None' = None, nack_delay: 'float' = 5.0)[source]

Bases: object

exec_ttl: float
task_ttl: float
poll_interval: float
default_queue: str
code_ref: str | None
nack_delay: float
class flowli.runtime.engine.StartResult(eid: UUID, deduplicated: bool)[source]

Bases: object

The answer of Engine.start_result.

eid: UUID
deduplicated: bool
exception flowli.runtime.UnknownExecution(eid: UUID | str)[source]

Bases: Exception

Context

The API for the author of a workflow. Every method creates one frame.

class flowli.runtime.Context(*, eid: UUID, workflow: WorkflowRef, memo: MemoTable, ports: Ports, actor: Actor, site: Site, clock: Callable[[], Timestamp], code_ref: str | None = None, starter: ChildStarter | None = None, workflow_of: Callable[[Callable[[...], Any]], WorkflowRef] | None = None, cancel_requested: Callable[[], bool] = lambda : ..., evidence: EvidenceWriter | None = None)[source]

Bases: object

property fid: str
property attempt: int
cancel_requested() → bool[source]
async run(*args: Any, **kwargs: Any) → Returned | Raised | Suspended[source]

Run the workflow function once. Return how it ended.

step(fn: Callable[[...], Any], /, *args: Any, name: str | None = None, key: str | None = None, retry: RetryPolicy | None = None, **kwargs: Any) → Awaitable[Any][source]
send(channel: str, payload: Any, *, correlation: str | None = None, scope: Literal['execution', 'global'] = 'execution', name: str = 'send', key: str | None = None) → Awaitable[int][source]
announce(kind: str, payload: dict[str, Any], *, name: str = 'announce', key: str | None = None) → Awaitable[int][source]
enqueue(queue: str, payload: Any, *, task_key: str, kind: TaskKind = TaskKind.DELEGATE, reason: str = 'delegate', name: str = 'enqueue', key: str | None = None) → Awaitable[str][source]

A step that puts a task on a queue. The memo is the task id.

task_key is the task’s dedup discriminator (see Task). key is this frame’s key.

withdraw(queue: str, task_id: str, *, name: str = 'withdraw', key: str | None = None) → Awaitable[bool][source]

A step that removes a task from a queue. The memo is whether a task was removed.

receive(channel: str, *, timeout: timedelta | None = None, scope: Literal['execution', 'global'] = 'execution', name: str = 'receive', key: str | None = None) → Awaitable[Message | None][source]
sleep(duration: timedelta, *, name: str = 'sleep', key: str | None = None) → Awaitable[None][source]
sleep_until(instant: Timestamp, *, name: str = 'sleep_until', key: str | None = None) → Awaitable[None][source]
child(workflow_fn: Callable[[...], Any], /, *args: Any, name: str | None = None, key: str | None = None, queue: str = 'default', **kwargs: Any) → Awaitable[Any][source]
gather(*aws: Awaitable[Any]) → Awaitable[list[Any]][source]
now() → Awaitable[Timestamp][source]
random() → Awaitable[float][source]
uuid() → Awaitable[str][source]
class flowli.runtime.Wait(fid: str, on: str, deadline: Timestamp | None)[source]

Bases: object

A frame that waits. The worker registers it. See 05-protocols.md section 3.

fid: str
on: str
deadline: Timestamp | None
class flowli.runtime.Returned(value: 'Any')[source]

Bases: object

value: Any
class flowli.runtime.Raised(error: 'BaseException')[source]

Bases: object

error: BaseException
class flowli.runtime.Suspended(waits: 'tuple[Wait, ...]')[source]

Bases: object

waits: tuple[Wait, ...]

Registry

class flowli.runtime.Registry[source]

Bases: object

workflow(name: str, version: str) → Callable[[Callable[[...], Any]], Callable[[...], Any]][source]
register(ref: WorkflowRef) → None[source]
entries() → list[WorkflowRef][source]

Every registered workflow, ordered by name and version.

get(name: str, version: str) → WorkflowRef[source]
of(fn: Callable[[...], Any]) → WorkflowRef[source]
class flowli.runtime.WorkflowRef(name: 'str', version: 'str', fn: 'Callable[..., Any]')[source]

Bases: object

name: str
version: str
fn: Callable[[...], Any]

Worker, sweeper and retention

class flowli.runtime.Worker(engine: Engine, queues: list[str], worker_id: str)[source]

Bases: object

property ports: Ports
async run_once() → bool[source]

Dequeue at most one task and process it. Return True when a task was processed.

async run_forever(stop: Event | None = None) → None[source]

Loop until stop. One failed iteration is logged, not fatal: the bucket arbitrates, so a lost race or a storage hiccup is retried on the next pass.

async run_execution(task: Task) → Done | Retry[source]
async run_detached_step(task: Task) → Done | Retry[source]

payload: {“fn”: “module:qualname”, “args”: […], “kwargs”: {…}, “name”: str}.

class flowli.runtime.Sweeper(engine: Engine, *, repair_window: timedelta = timedelta(seconds=60), sweeper_id: str = 'sweeper', source: ControlSource | None = None)[source]

Bases: object

async run_once() → SweepReport[source]
async run_forever(stop: Event | None = None, interval: float = 60.0) → None[source]
class flowli.runtime.SweepReport(timers_fired: 'list[str]' = <factory>, recovered: 'list[Eid]' = <factory>, restarted: 'list[Eid]' = <factory>, repaired: 'list[Eid]' = <factory>, waits_cleared: 'list[tuple[str, Eid]]'=<factory>)[source]

Bases: object

timers_fired: list[str]
recovered: list[UUID]
restarted: list[UUID]
repaired: list[UUID]
waits_cleared: list[tuple[str, UUID]]
property total: int
class flowli.runtime.Retention(engine: Engine, *, delay: timedelta = timedelta(days=30), source: ControlSource | None = None, retention_id: str = 'retention')[source]

Bases: object

async run_once() → RetentionReport[source]
async run_forever(stop: Event | None = None, interval: float = 3600.0) → None[source]
async fold(eid: UUID, report: RetentionReport | None = None) → bool[source]

Archive one execution, then delete its live state. Return True when archived now.

Order: archive (put-if-absent), deletes, then the announcement. A rerun after a crash finds the archive, finishes the deletes and announces. The announcement can repeat after a crash between the last delete and the announcement. That is harmless.

class flowli.runtime.RetentionReport(archived: 'list[Eid]' = <factory>, cleaned: 'list[Eid]' = <factory>)[source]

Bases: object

archived: list[UUID]
cleaned: list[UUID]
class flowli.runtime.ControlSource(*args, **kwargs)[source]

Bases: Protocol

Where the sweeper learns execution statuses: the in-memory fold or a projection.

async snapshot() → dict[UUID, KnownExecution][source]

Return the known executions, terminal ones optional.

async terminal_before(before: Timestamp) → dict[UUID, KnownExecution][source]

Return terminal executions whose last entry is older than before. For retention.

class flowli.runtime.ControlView(engine: Engine | None = None)[source]

Bases: object

A fold of the control log: eid -> last known lifecycle entry. Keeps a cursor.

executions: dict[UUID, KnownExecution]
async snapshot() → dict[UUID, KnownExecution][source]
async terminal_before(before: Timestamp) → dict[UUID, KnownExecution][source]
apply(seq: int, entry: Entry) → None[source]
class flowli.runtime.KnownExecution(status: 'ExecutionStatus', last_type: 'str', last_at: 'Timestamp', queue: 'str' = 'default', archived: 'bool' = False)[source]

Bases: object

status: ExecutionStatus
last_type: str
last_at: Timestamp
queue: str = 'default'
archived: bool = False

Consumer

The loop of a consumer of delegate tasks, minus the work.

class flowli.runtime.Consumer(engine: Any, handler: Handler, config: ConsumerConfig | None = None, *, actor: Actor | None = None)[source]

Bases: object

Takes delegate tasks from queues and answers them.

One task at a time, per the rule that one session holds one task (section 2). Run several consumers for more.

property ports: Any
async run_once() → bool[source]

Take one task and see it through. False when every queue was empty.

async run_forever(stop: Event | None = None) → None[source]
async recover() → list[str][source]

Take back the tasks this holder still owns, after a restart.

attach writes nothing and does not bump the epoch: this is the same ownership period continuing, after an interruption of the process that held it (section 4.3).

class flowli.runtime.ConsumerConfig(queues: tuple[str, ...] = ('agents',), holder: str = 'consumer', ttl: float = 300.0, poll_interval: float = 1.0, grace: float = 120.0, nack_delay: timedelta = datetime.timedelta(seconds=30), conflict_delay: timedelta = datetime.timedelta(seconds=10))[source]

Bases: object

See specs/10-agent-runner.md section 11.

queues: tuple[str, ...]
holder: str
ttl: float
poll_interval: float
grace: float
nack_delay: timedelta
conflict_delay: timedelta
class flowli.runtime.ConsumerReport(delivered: 'list[str]' = <factory>, refused: 'list[str]' = <factory>, recovered: 'list[str]' = <factory>, terminated: 'list[str]' = <factory>)[source]

Bases: object

delivered: list[str]
refused: list[str]
recovered: list[str]
terminated: list[str]
class flowli.runtime.Handler(*args, **kwargs)[source]

Bases: Protocol

The work. Every method takes the Held task it belongs to.

start must be write-ahead: whatever it starts is addressable by the handle it returns, and the consumer writes that handle before the next heartbeat (section 4.2).

async prepare(held: Held) → None[source]

Acquire what the work needs. Raise Refused on a conflict.

async start(held: Held) → Any[source]

Start the work. Return a JSON-compatible handle.

async reattach(held: Held, handle: Any) → Any | None[source]

Re-attach to live work. None when it is gone. Never starts anything.

async poll(held: Held, handle: Any) → Any | None[source]

The reply payload when the work is done, None while it runs.

async interrupt(held: Held, handle: Any, *, grace: float) → Any[source]

Stop the work within grace, and return the reply payload to deliver.

async terminate(held: Held, handle: Any) → None[source]

Stop work this consumer will not wait for. Never raises.

async release(held: Held) → None[source]

Give back what prepare acquired. Never raises.

class flowli.runtime.Held(claimed: ClaimedTask, task: DelegateTask, adopted: bool = False)[source]

Bases: object

One task this consumer owns, and the state it keeps in the lease.

The lease document is the whole ownership model of a long attempt: it fences, it survives a restart, and its state carries whatever the consumer must find again (section 4.1).

claimed: ClaimedTask
task: DelegateTask
adopted: bool = False
property lease: Any
property task_id: str
property queue: str
property state: dict[str, Any]
property handle: Any
property cancel_requested: bool
async write_handle(handle: Any) → None[source]

Record the address of the work, before the work can outlive us.

exception flowli.runtime.Refused(reason: str, delay: timedelta | None = None)[source]

Bases: Exception

The consumer cannot take this task now. The task goes back on the queue.

flowli.patterns

Helpers written against the Context API only. The engine does not know them.

async flowli.patterns.delegate(ctx: Context, queue: str, payload: Any, *, timeout: timedelta | None = None, name: str = 'delegate', key: str | None = None, task_key: str | None = None) → Message | None[source]

Enqueue once, then receive the reply. None when timeout expires first.

Frames: {name}-enqueue, {name}-receive, and {name}-withdraw on timeout, all under the caller’s frame with key.

async flowli.patterns.review(ctx: Context, queue: str, payload: Any, *, timeout: timedelta | None = None, name: str = 'review', key: str | None = None) → Decision | None[source]

Ask the humans on queue. Return their Decision, or None when timeout expires.

name names this review’s frames. Two reviews under one frame with one key need two names, or their frames collide.

class flowli.patterns.Decision(verdict: 'str', data: 'Any', by: 'Actor', at: 'Timestamp')[source]

Bases: object

verdict: str
data: Any
by: Actor
at: Timestamp
class flowli.patterns.Reviews(engine: Engine, *, take_ttl: float = 30.0)[source]

Bases: object

Operator side. engine.reviews.decide(…).

async decide(rid: str, *, eid: UUID, queue: str, verdict: str, by: Actor, data: Any = None) → Decision[source]

Record exactly one decision for rid, deliver it to the workflow, remove the task.

eid and queue come from the inbox row (review.requested carries both).

Idempotent: a second call returns the first decision. A call after a crash between the claim and the delivery completes the delivery.

The delivery does not wait for the task. A task that another consumer holds, or that a worker pushed into the future with a nack, would otherwise swallow a decision that is already recorded. A repeat sends the same payload again, and the receive frame consumes the first message only.

async flowli.patterns.saga(ctx: Context, actions: list[tuple[Callable[[...], Any], Callable[[...], Any]]], *, name: str = 'saga', undo_retry: RetryPolicy | None = None) → list[Any][source]

Each (act, undo) pair runs as steps {name}-act:{i} and {name}-undo:{i}.

undo receives the value act returned. On failure the compensations run in reverse order, then the error is raised again. A replay never runs an action or a compensation twice.

async flowli.patterns.fan_out(ctx: Context, child_fn: Callable[[...], Any], items: list[Any], *, name: str | None = None, queue: str = 'default') → list[Any][source]

Start a child per item, keyed by position. Return the values in item order.

async flowli.patterns.on_tick(engine: Engine, workflow_fn: Callable[[...], Any], schedule_name: str, tick: Timestamp, /, *args: Any, queue: str | None = None, **kwargs: Any) → UUID[source]

Start workflow_fn for tick. Two schedulers that fire the same tick get the same eid.

flowli.domain

The objects, the errors and the ports. This layer has no I/O.

Execution and frames

class flowli.domain.Execution(eid: 'Eid', workflow: 'str', version: 'str', args: 'Any', created_by: 'Provenance', parent: 'FrameRef | None' = None, dispatch_key: 'str | None' = None, queue: 'str' = 'default')[source]

Bases: object

eid: UUID
workflow: str
version: str
args: Any
created_by: Provenance
parent: FrameRef | None
dispatch_key: str | None
queue: str
class flowli.domain.ExecutionStatus(*values)[source]

Bases: StrEnum

PENDING = 'pending'
RUNNING = 'running'
SUSPENDED = 'suspended'
COMPLETED = 'completed'
FAILED = 'failed'
CANCELLED = 'cancelled'
property is_terminal: bool
class flowli.domain.FrameRef(eid: UUID, fid: str)[source]

Bases: object

One frame of one execution. Everything derived from that pair is derived here.

eid: UUID
fid: str
property child_eid: UUID

a UUIDv5, so it is deterministic.

Type:

The id of the child execution this frame starts

property child_channel: str

Where the child execution reports its terminal outcome to this frame.

property step_channel: str

Where a detached run of this step frame reports its outcome.

reply_channel(frame: str) → str[source]

Where a delegate sub-frame frame of this frame receives its answer.

class flowli.domain.Frame(ref: 'FrameRef', kind: 'FrameKind', name: 'str', args_digest: 'str', retry: 'RetryPolicy' = <factory>)[source]

Bases: object

ref: FrameRef
kind: FrameKind
name: str
args_digest: str
retry: RetryPolicy
class flowli.domain.FrameKind(*values)[source]

Bases: StrEnum

ROOT = 'root'
STEP = 'step'
CHILD = 'child'
RECEIVE = 'receive'
SLEEP = 'sleep'
class flowli.domain.RetryPolicy(max_attempts: 'int' = 1, backoff: 'timedelta' = datetime.timedelta(0), backoff_factor: 'float' = 1.0, max_backoff: 'timedelta | None' = None)[source]

Bases: object

max_attempts: int
backoff: timedelta
backoff_factor: float
max_backoff: timedelta | None
allows_retry(failed_attempt: int, retryable: bool) → bool[source]

True when attempt failed_attempt failed and another attempt is allowed.

delay_after(failed_attempt: int) → timedelta[source]

Delay before attempt failed_attempt + 1. See 05-protocols.md section 6.

class flowli.domain.Attempt(frame: 'FrameRef', number: 'int', provenance: 'Provenance')[source]

Bases: object

frame: FrameRef
number: int
provenance: Provenance
class flowli.domain.Completed(value: 'Any')[source]

Bases: object

value: Any
class flowli.domain.Failed(error_type: 'str', message: 'str', retryable: 'bool' = True)[source]

Bases: object

error_type: str
message: str
retryable: bool
classmethod from_exception(exc: BaseException, retryable: bool = True) → Failed[source]

The journal

class flowli.domain.Entry(type: str, fid: str, payload: dict[str, Any], provenance: Provenance)[source]

Bases: object

One event in the journal or in the control log.

type is an EntryType value or an ‘announce.*’ string.

type: str
fid: str
payload: dict[str, Any]
provenance: Provenance
property is_terminal: bool
classmethod execution_started(prov: Provenance, args: Any) → Entry[source]
classmethod frame_started(prov: Provenance, fid: str, kind: str, name: str, args_digest: str, attempt: int) → Entry[source]
classmethod frame_completed(prov: Provenance, fid: str, attempt: int, value: Any) → Entry[source]
classmethod frame_failed(prov: Provenance, fid: str, attempt: int, failed: Failed, retry_at: str | None = None) → Entry[source]
classmethod frame_suspended(prov: Provenance, fid: str, attempt: int, on: str, deadline: str | None = None) → Entry[source]
classmethod frame_fulfilled(prov: Provenance, fid: str, attempt: int, on: str, message_seq: int | None = None) → Entry[source]
classmethod execution_suspended(prov: Provenance, on: list[str]) → Entry[source]
classmethod execution_resumed(prov: Provenance, epoch: int, reason: str) → Entry[source]
classmethod execution_completed(prov: Provenance, value: Any) → Entry[source]
classmethod execution_failed(prov: Provenance, failed: Failed) → Entry[source]
classmethod execution_cancelled(prov: Provenance, by: dict[str, Any]) → Entry[source]
classmethod execution_migrated(prov: Provenance, from_version: str, version: str) → Entry[source]
classmethod announce(prov: Provenance, fid: str, kind: str, payload: dict[str, Any]) → Entry[source]
classmethod task_enqueued(task: Any) → Entry[source]

Control-log twin of a Task. task is a domain Task (kept untyped to avoid a cycle).

class flowli.domain.EntryType(*values)[source]

Bases: StrEnum

EXECUTION_STARTED = 'execution.started'
EXECUTION_SUSPENDED = 'execution.suspended'
EXECUTION_RESUMED = 'execution.resumed'
EXECUTION_COMPLETED = 'execution.completed'
EXECUTION_FAILED = 'execution.failed'
EXECUTION_CANCELLED = 'execution.cancelled'
FRAME_STARTED = 'frame.started'
FRAME_COMPLETED = 'frame.completed'
FRAME_FAILED = 'frame.failed'
FRAME_SUSPENDED = 'frame.suspended'
FRAME_FULFILLED = 'frame.fulfilled'
EXECUTION_MIGRATED = 'execution.migrated'
EXECUTION_CREATED = 'execution.created'
EXECUTION_ARCHIVED = 'execution.archived'
TASK_ENQUEUED = 'task.enqueued'
class flowli.domain.MemoTable(memos: dict[str, ~flowli.domain.frames.Completed]=<factory>, failures: dict[str, list[~flowli.domain.frames.Failed]]=<factory>, suspended: dict[str, str]=<factory>, fulfilled: dict[str, int | None]=<factory>, deadlines: dict[str, str]=<factory>, retry_at: dict[str, str]=<factory>, digests: dict[str, str]=<factory>, consumed: dict[str, int]=<factory>, terminal: Completed | Failed | None = None, cancelled: bool = False, tail: int = 0)[source]

Bases: object

Derived state of one journal. See specs/02-journal.md section 4.

memos: dict[str, Completed]
failures: dict[str, list[Failed]]
suspended: dict[str, str]
fulfilled: dict[str, int | None]
deadlines: dict[str, str]
retry_at: dict[str, str]
digests: dict[str, str]
consumed: dict[str, int]
terminal: Completed | Failed | None = None
cancelled: bool = False
tail: int = 0
classmethod build(entries: list[Sequenced[Entry]]) → MemoTable[source]
apply(s: Sequenced[Entry]) → None[source]
property is_terminal: bool
next_attempt(fid: str) → int[source]
last_consumed(channel: str) → int[source]
check_digest(fid: str, computed: str) → str | None[source]

Return the recorded digest when it differs from computed, else None.

class flowli.domain.Sequenced(seq: 'int', item: 'T')[source]

Bases: Generic

seq: int
item: T
class flowli.domain.Condition[source]

Bases: object

The on field of frame.suspended / frame.fulfilled. A prefixed string.

CHANNEL = 'channel'
TIMER = 'timer'
CHILD = 'child'
OPERATOR = 'operator'
static channel(name: str) → str[source]
static timer(timer_id: str) → str[source]
static child(eid: str) → str[source]
static operator() → str[source]
static parse(on: str) → tuple[str, str | None][source]

Tasks, messages and timers

class flowli.domain.Task(queue: str, kind: TaskKind, target: FrameRef, reason: str, enqueued_by: Provenance, key: str = '', not_before: Timestamp | None = None, payload: Any = None)[source]

Bases: object

A unit of work on a queue. Its id is derived: {kind}:{eid} plus :{key} when set.

key is the dedup discriminator inside one kind and target: the message that caused a resume, the frame of a detached step, the name of a delegate. Two tasks with the same id are one task, so the enqueue of a duplicate is a no-op.

queue: str
kind: TaskKind
target: FrameRef
reason: str
enqueued_by: Provenance
key: str
not_before: Timestamp | None
payload: Any
static id_for(kind: TaskKind, eid: UUID, key: str = '') → str[source]
property task_id: str
visible_at(default: Timestamp) → Timestamp[source]

The instant from which a worker may dequeue this task.

class flowli.domain.TaskKind(*values)[source]

Bases: StrEnum

START = 'start'
RESUME = 'resume'
RUN_STEP = 'run_step'
DELEGATE = 'delegate'
class flowli.domain.DelegateTask(eid: UUID, fid: str, reply_channel: str, payload: Any, reply_fid: str | None = None)[source]

Bases: object

The payload of a DELEGATE task, as its consumer sees it.

The contract between the pattern that enqueues (06-patterns.md, section 1) and the consumer that answers (10-agent-runner.md, section 2), so it belongs to neither of them.

eid: UUID
fid: str
reply_channel: str
payload: Any
reply_fid: str | None
classmethod from_task_payload(payload: dict[str, Any]) → DelegateTask[source]
property waiting_fid: str

The frame a consumer files its evidence under.

class flowli.domain.Message(channel: 'str', seq: 'int', payload: 'Any', sent_by: 'Provenance', correlation: 'str | None' = None)[source]

Bases: object

channel: str
seq: int
payload: Any
sent_by: Provenance
correlation: str | None
class flowli.domain.Timer(due_at: Timestamp, target: FrameRef)[source]

Bases: object

A future instant at which the sweeper resumes one frame.

The id is derived, not stored: {due_at_iso}-{eid}-{fid_digest}. It starts with the instant, so ids sort by due time, and two timers for one frame and instant are one.

due_at: Timestamp
target: FrameRef
property timer_id: str
is_due(now: Timestamp | None = None) → bool[source]
flowli.domain.execution_channel(eid: UUID, name: str) → str[source]

A channel scoped to one execution: ‘{eid}.{name}’.

Provenance

class flowli.domain.Actor(kind: 'ActorKind', id: 'str', on_behalf_of: 'str | None' = None)[source]

Bases: object

kind: Literal['worker', 'human', 'system', 'schedule']
id: str
on_behalf_of: str | None
classmethod worker(worker_id: str) → Actor[source]
classmethod human(email: str, on_behalf_of: str | None = None) → Actor[source]
classmethod system(name: str, on_behalf_of: str | None = None) → Actor[source]
classmethod schedule(name: str) → Actor[source]
class flowli.domain.Site(host: 'str', pid: 'int', worker_id: 'str', instance: 'str | None' = None, region: 'str | None' = None, epoch: 'int | None' = None)[source]

Bases: object

host: str
pid: int
worker_id: str
instance: str | None
region: str | None
epoch: int | None
classmethod local(worker_id: str, instance: str | None = None, region: str | None = None) → Site[source]
with_epoch(epoch: int | None) → Site[source]
class flowli.domain.Code(workflow: 'str', version: 'str', frame_kind: 'str', frame_name: 'str', code_ref: 'str | None' = None)[source]

Bases: object

workflow: str
version: str
frame_kind: str
frame_name: str
code_ref: str | None
class flowli.domain.Provenance(actor: 'Actor', site: 'Site', code: 'Code', attempt: 'int', at: 'Timestamp')[source]

Bases: object

actor: Actor
site: Site
code: Code
attempt: int
at: Timestamp
classmethod now(actor: Actor, site: Site, code: Code, attempt: int = 1) → Provenance[source]

Errors

Domain errors. See specs/04-api.md section 4.

exception flowli.domain.errors.FlowliError[source]

Bases: Exception

Base of all domain errors.

exception flowli.domain.errors.InvalidName[source]

Bases: FlowliError, ValueError

An identifier does not match its required pattern.

exception flowli.domain.errors.NondeterminismError(fid: str, recorded: str, computed: str)[source]

Bases: FlowliError

Replay found a frame id with a different args digest.

exception flowli.domain.errors.DuplicateFrameError(fid: str)[source]

Bases: FlowliError

A frame key repeats under one parent.

exception flowli.domain.errors.ChildFailed(child_eid: Any, status: str, error: str | None = None)[source]

Bases: FlowliError

A child execution failed or was cancelled.

exception flowli.domain.errors.LeaseLost[source]

Bases: FlowliError

The execution lease was stolen. Stop all work on the execution.

exception flowli.domain.errors.Cancelled[source]

Bases: FlowliError

Raised inside the workflow function at a frame boundary after a cancel request.

exception flowli.domain.errors.WorkflowNotRegistered(workflow: str, version: str)[source]

Bases: FlowliError

exception flowli.domain.errors.NonRetryableError[source]

Bases: Exception

A step raises this, or a subclass, to stop retries. The frame fails at once.

exception flowli.domain.errors.StepFailed(fid: str, failed: Any)[source]

Bases: FlowliError

A step frame has no attempt left. Raised in the workflow function on every replay.

Ports

Each port is a Protocol. An adapter implements all of them.

Ports the domain requires. See specs/03-ports.md.

One adapter implements all of them. The domain never imports the adapter.

class flowli.domain.ports.Lease(*args, **kwargs)[source]

Bases: Protocol

Fenced ownership. renew() raises LeaseLost when fenced.

property epoch: int
property state: Any
property deadline_at: Timestamp | None

When the lease expires unless it is renewed.

async renew() → None[source]
async release(state: Any = None) → None[source]
async refresh_state() → Any[source]

Absorb cooperative writes. Return the fresh state.

async update_state(fn: Callable[[Any], Any]) → Any[source]

Apply fn to the freshest state and write the result. fn must be pure.

The state of a task lease belongs to its holder. A consumer of a long task keeps the address of its work there (spec 10, section 4.2).

class flowli.domain.ports.Journal(*args, **kwargs)[source]

Bases: Protocol

Trace and memo of one execution. Spec 03 section 3.

async append(eid: UUID, entry: Entry) → int[source]
async read(eid: UUID, after: int = 0) → list[Sequenced[Entry]][source]
async tail(eid: UUID) → int[source]
async delete(eid: UUID) → None[source]

Delete the whole journal. Retention only, after the archive is written.

class flowli.domain.ports.ControlLog(*args, **kwargs)[source]

Bases: Protocol

Lifecycle entries of all executions. Spec 03 section 4.

async announce(entry: Entry) → int[source]
async read(after: int = 0) → list[Sequenced[Entry]][source]
class flowli.domain.ports.LeaseInfo(epoch: int, holder: str | None, deadline_at: Timestamp | None, released: bool, state: Any)[source]

Bases: object

A read of a lease document, without acquisition.

epoch: int
holder: str | None
deadline_at: Timestamp | None
released: bool
state: Any
is_expired(now: Timestamp | None = None) → bool[source]

True when nobody holds the lease: released, or past its deadline.

class flowli.domain.ports.Ownership(*args, **kwargs)[source]

Bases: Protocol

One owner per execution. Spec 03 section 5.

async acquire(eid: UUID, holder: str, ttl: float) → Lease | None[source]
async request_cancel(eid: UUID, by: Provenance) → None[source]
async inspect(eid: UUID) → LeaseInfo | None[source]

Read the lease document. None when no lease was ever acquired.

async delete(eid: UUID) → None[source]

Delete the lease document. Retention only.

class flowli.domain.ports.Dispatch(*args, **kwargs)[source]

Bases: Protocol

Idempotent start, and generic exactly-one claims. Spec 03 section 6.

async claim_start(key: str, eid: UUID) → tuple[bool, UUID][source]

Return (won, eid). The loser gets the winner’s eid.

async claim(key: str, value: Any) → tuple[bool, Any][source]

Put-if-absent of a JSON value under wf/claims/{key}. Return (won, winner’s value).

class flowli.domain.ports.QueueDepth(total: int, visible: int, claimed: int)[source]

Bases: object

What a queue holds now. See specs/03-ports.md section 7.

total: int
visible: int
claimed: int
class flowli.domain.ports.ClaimedTask(task: 'Task', key: 'str', lease: 'Lease')[source]

Bases: object

task: Task
key: str
lease: Lease
class flowli.domain.ports.Queue(*args, **kwargs)[source]

Bases: Protocol

Tasks. Spec 03 section 7.

async enqueue(task: Task) → bool[source]

Return True when this call created the task. False is a no-op.

async ensure(task: Task) → bool[source]

enqueue, and repair an enqueue marker that outlived its task.

An adapter that dedups with a separate marker writes the marker first, so a crash between the two writes leaves a marker that refuses every later enqueue of that task id. Repair paths only: it costs one listing, where enqueue costs none.

async dequeue(queue: str, worker: str, ttl: float) → ClaimedTask | None[source]
async take(queue: str, task_id: str, worker: str, ttl: float) → ClaimedTask | None[source]

Lease one specific task. None when absent, hidden, or held by someone else.

async ack(claimed: ClaimedTask) → None[source]
async nack(claimed: ClaimedTask, delay: timedelta) → None[source]
async depth(queue: str) → QueueDepth[source]

Counts only: one list of the tasks and one of the leases.

async pending(queue: str, limit: int = 100) → list[Task][source]

The tasks on the queue, oldest first.

async peek(queue: str, task_id: str) → Task | None[source]

Read one task. A read takes no lease: ownership is for removal.

async attach(queue: str, task_id: str, holder: str, ttl: float) → ClaimedTask | None[source]

Re-attach to a task that holder holds now. Writes nothing.

None when the task is gone, when the holder does not match, or when the lease expired. See spec 10, sections 4.3 and 4.4.

async request_cancel(queue: str, task_id: str, by: Provenance) → None[source]

Ask the holder of a task to stop. It fences nobody.

The holder reads it on its next renew. This is Ownership.request_cancel one level down.

class flowli.domain.ports.Channel(*args, **kwargs)[source]

Bases: Protocol

Messages. Spec 03 section 8.

async send(message: Message) → int[source]
async read(channel: str, after: int = 0) → list[Message][source]
async register_wait(channel: str, ref: FrameRef) → None[source]
async clear_wait(channel: str, ref: FrameRef) → None[source]
async waiters(channel: str) → list[FrameRef][source]
async all_waits() → list[tuple[str, FrameRef]][source]

Every wait marker, as (channel, ref). For the sweeper.

async scoped(eid: UUID) → list[str][source]

Names of the channels scoped to eid: those with the prefix ‘{eid}.’.

async delete_channel(channel: str) → None[source]

Delete a channel’s log. Retention only.

class flowli.domain.ports.Timers(*args, **kwargs)[source]

Bases: Protocol

Future resumes. Spec 03 section 9.

async schedule(timer: Timer) → None[source]
async due(now: Timestamp) → list[Timer][source]
async remove(timer: Timer) → None[source]
async remove_for(eid: UUID) → int[source]

Remove every timer of one execution. Return how many. Retention only.

class flowli.domain.ports.ExecutionStore(*args, **kwargs)[source]

Bases: Protocol

The wf/exec/{eid}/meta object. Put-if-absent, then CAS for migrate.

async create(execution: Execution) → bool[source]
async read(eid: UUID) → Execution | None[source]
async replace(eid: UUID, fn: Callable[[Execution], Execution]) → Execution | None[source]

Read-modify-write with CAS. fn must be pure. None when the record is absent.

async delete(eid: UUID) → None[source]
class flowli.domain.ports.Archive(*args, **kwargs)[source]

Bases: Protocol

Folded journals of finished executions. wf/archive/{eid}. Put-if-absent.

async write(eid: UUID, data: dict[str, Any]) → bool[source]
async read(eid: UUID) → dict[str, Any] | None[source]
class flowli.domain.ports.EvidenceRef(eid: UUID, fid: str, attempt: int)[source]

Bases: object

One attempt of one frame. The key of evidence is the attempt, because a retry writes a second log.

eid: UUID
fid: str
attempt: int
class flowli.domain.ports.EvidenceItem(ref: 'EvidenceRef', name: 'str', media_type: 'str', size: 'int')[source]

Bases: object

ref: EvidenceRef
name: str
media_type: str
size: int
class flowli.domain.ports.Evidence(*args, **kwargs)[source]

Bases: Protocol

What a step or an agent wrote while it ran. Spec 03 section 12.

The journal says what happened. Evidence explains it. The engine never reads evidence, and no decision depends on it.

async append_log(ref: EvidenceRef, part: bytes) → None[source]
async read_log(ref: EvidenceRef) → bytes[source]
async put(ref: EvidenceRef, name: str, data: bytes, media_type: str) → str[source]
async get(ref: EvidenceRef, name: str) → bytes | None[source]
async list(eid: UUID) → list[EvidenceItem][source]
async delete_for(eid: UUID) → int[source]

Retention only. Returns how many objects were removed.

class flowli.domain.ports.Ports(journal: Journal, control: ControlLog, ownership: Ownership, dispatch: Dispatch, queue: Queue, channel: Channel, timers: Timers, executions: ExecutionStore, archive: Archive, evidence: Evidence)[source]

Bases: object

Everything the engine needs, bundled.

journal: Journal
control: ControlLog
ownership: Ownership
dispatch: Dispatch
queue: Queue
channel: Channel
timers: Timers
executions: ExecutionStore
archive: Archive
evidence: Evidence

Names and ids

Identifier rules and generators. See specs/01-domain-model.md.

flowli.domain.names.Eid

An execution id: UUIDv7 for a started execution, UUIDv5 for a derived child.

flowli.domain.names.parse_eid(value: UUID | str) → UUID[source]

Accept a UUID or its canonical string. Raise InvalidName otherwise.

flowli.domain.names.check_frame_name(name: str) → str[source]
flowli.domain.names.check_channel(channel: str) → str[source]
flowli.domain.names.check_queue(queue: str) → str[source]
flowli.domain.names.digest_text(text: str) → str[source]

sha256 hex of a string. Used to derive names from frame ids, which hold ‘/’ and ‘#’.

flowli.domain.names.new_eid() → UUID[source]

A UUIDv7: time-ordered, carries its creation instant (cairndb Timestamp.from_uuid7).

flowli.adapters

CairnDB

class flowli.adapters.cairndb.CairnBackend(db: CairnDB, *, now: Callable[[], Timestamp] = Timestamp.now, log_cache_size: int = 128)[source]

Bases: object

All ports over one CairnDB engine.

classmethod configure(config: dict[str, Any], **kwargs: Any) → CairnBackend[source]
property ports: Ports
projection(**kwargs: Any) → Any[source]

The SQLite projection of the control log. See cairndb_projection.py.

async close() → None[source]
class flowli.adapters.cairndb_projection.WorkflowProjection(db: CairnDB, *, db_path: str | None = None, version: str = '1', poll_interval: float = 5.0)[source]

Bases: object

NAME = 'wf_view'
async refresh() → int | None[source]
async start() → None[source]
async stop() → None[source]
async wait_for(seq: int, timeout: float = 30.0) → bool[source]
property path: str
property ready: bool

False until the first refresh wrote the file.

A control log with no commit yet leaves no file, and a read-only open of a file that is not there fails. A projection with no file has seen nothing, so every query answers empty.

connect() → Connection[source]
execution(eid: UUID | str) → ExecutionRow | None[source]
executions(*, status: ExecutionStatus | None = None, workflow: str | None = None, parent_eid: UUID | str | None = None, limit: int = 100, after: tuple[str, str] | None = None) → list[ExecutionRow][source]

Newest first. after is the (updated_at, eid) of the last row of the previous page: a keyset cursor over the same order.

children(eid: UUID | str) → list[ExecutionRow][source]
counts_by_status() → dict[ExecutionStatus, int][source]
tasks(*, eid: UUID | str | None = None, queue: str | None = None, kind: str | None = None, limit: int = 100) → list[TaskRow][source]

Enqueued tasks, newest first. See the note on TaskRow.

pending_reviews(queue: str | None = None) → list[ReviewRow][source]
review(rid: str) → ReviewRow | None[source]
announcements(*, kind: str | None = None, eid: UUID | str | None = None, limit: int = 100) → list[AnnouncementRow][source]
async snapshot() → dict[UUID, KnownExecution][source]

ControlSource for the Sweeper: refresh, then the non-terminal executions.

async terminal_before(before: Timestamp) → dict[UUID, KnownExecution][source]

ControlSource for Retention: terminal, not archived, last entry older than before.

In memory

For the tests. No bucket, no disk.

class flowli.adapters.memory.MemoryBackend(clock: ManualClock = <factory>)[source]

Bases: object

All ports over one shared clock. backend.ports is what the engine takes.

clock: ManualClock
property ports: Ports
class flowli.adapters.memory.ManualClock(start: Timestamp | datetime | None = None)[source]

Bases: object

A clock that moves only when told to.

advance(delta: timedelta) → Timestamp[source]
set(at: Timestamp | datetime) → None[source]

flowli.api

The HTTP service over one engine. Install the api extra.

flowli.api.create_app(engine: Any, *, authenticator: Authenticator, projection: Any, catalog: Catalog | None = None, config: ApiConfig | None = None) → FastAPI[source]

Wire one engine, one projection and one catalog into a FastAPI app.

The app owns the refresh loop of the projection: it catches up once at startup and polls while it runs (08-projection.md, section 4).

class flowli.api.ApiConfig(queues: tuple[str, ...] = ('default',), default_queue: str = 'default', wait_ms: int = 2000, default_limit: int = 50, max_limit: int = 500, refresh_after_write: bool = True, task_ttl: float = 300.0, max_task_ttl: float = 3600.0, nack_delay: float = 5.0, max_wait_seconds: float = 20.0, poll_interval: float = 0.5, title: str = 'flowli')[source]

Bases: object

See specs/09-http-api.md section 3.

queues: tuple[str, ...]
default_queue: str
wait_ms: int
default_limit: int
max_limit: int
refresh_after_write: bool
task_ttl: float
max_task_ttl: float
nack_delay: float
max_wait_seconds: float
poll_interval: float
title: str
class flowli.api.Catalog(registry: Registry, *, default_queue: str = 'default')[source]

Bases: object

Every registered workflow, with its argument schema.

Built once, at startup: a schema never changes while the process runs.

entries() → list[CatalogEntry][source]
entry(name: str, version: str) → CatalogEntry[source]
validate(name: str, version: str, args: dict[str, Any]) → dict[str, Any][source]

The arguments, checked against the schema. Unchecked when there is none.

class flowli.api.Principal(actor: Actor, capabilities: frozenset[str] = <factory>)[source]

Bases: object

Who is calling, and what they may do.

actor: Actor
capabilities: frozenset[str]
allows(capability: str, scope: str | None = None) → bool[source]
require(capability: str, scope: str | None = None) → None[source]
class flowli.api.StaticAuthenticator(principals: Mapping[str, Principal])[source]

Bases: object

A fixed map from token to principal, for tests and for a local run.

It is not a development back door: an unknown token is refused, and there is no way to name an actor from outside the map.

principals: Mapping[str, Principal]
authenticate(token: str | None, on_behalf_of: str | None) → Principal[source]
class flowli.api.OIDCAuthenticator(config: OIDCConfig, *, jwk_client: Any | None = None)[source]

Bases: object

Bearer JWT against the JWKS of one issuer.

The service starts no flow of its own. The web application does the authorization-code flow with PKCE and sends the access token.

authenticate(token: str | None, on_behalf_of: str | None) → Principal[source]
class flowli.api.OIDCConfig(issuer: 'str', audience: 'str', jwks_url: 'str', actor_claim: 'str' = 'email', groups_claim: 'str' = 'groups', scope_claim: 'str' = 'scope', client_claim: 'str' = 'client_id', roles: 'Mapping[str, Sequence[str]]'=<factory>, machine_actor_kind: 'str' = 'worker', leeway: 'float' = 30.0)[source]

Bases: object

issuer: str
audience: str
jwks_url: str
actor_claim: str = 'email'
groups_claim: str = 'groups'
scope_claim: str = 'scope'
client_claim: str = 'client_id'
roles: Mapping[str, Sequence[str]]
machine_actor_kind: str = 'worker'
leeway: float = 30.0
class flowli.api.CachingAuthenticator(inner: Authenticator, ttl: float = 60.0)[source]

Bases: object

Verifies each token once per ttl. A JWT verification is CPU work.

authenticate(token: str | None, on_behalf_of: str | None) → Principal[source]

flowli.log

Structured logging for flowli, on structlog (which CairnDB uses too).

Every runtime component logs events with key-value fields, never formatted strings:

log.info(“execution_completed”, eid=eid, epoch=2, worker_id=”w-1”)

Context variables carry worker_id, task_id, eid, epoch, fid, attempt through every event emitted while that work is in flight, so a line can be correlated without repeating the fields at each call site.

Call configure_logging() once per process. The CLI does it. A library user who does not call it gets structlog’s defaults, which print to stdout at every level.

The chain holds one processor of its own, evidence.capture: it keeps each event as evidence of the attempt that is running, and does nothing when none is (see flowli.evidence).

flowli.log.bound(**fields: Any) → Iterator[None][source]

Bind fields to every event emitted inside the block, in this task’s context.

flowli.log.configure_logging(level: str = 'INFO', fmt: Literal['console', 'json'] = 'console') → None[source]

Configure structlog for this process. Shared with CairnDB’s loggers.

flowli.log.get_logger(name: str) → Any[source]

flowli.evidence

Attempt logs: what a step wrote while it ran. See specs/03-ports.md section 12.

An author writes log.info(“fetched”, n=12) and nothing more. log.py binds eid, fid and attempt to every event of a worker already, and the capture processor below puts the event in the buffer of the running attempt. The buffer flushes as one object per flush, never one object per event.

Evidence is never authority. A write of it is best effort: an error is logged and is not raised into the frame, because a frame must not fail because a log did not write.

class flowli.evidence.AttemptBuffer(ref: EvidenceRef, limit: int = 1048576, pending: list[bytes] = None, written: int = 0, dropped: int = 0, tail: deque[bytes] = None)[source]

Bases: object

The lines of one attempt, between two flushes.

Past limit the buffer keeps a bounded tail only and counts what it dropped, so a runaway step costs a known amount of memory and of storage.

ref: EvidenceRef
limit: int = 1048576
pending: list[bytes] = None
written: int = 0
dropped: int = 0
tail: deque[bytes] = None
add(event: dict[str, Any]) → None[source]
drain() → bytes | None[source]

The next part to write, or None when there is nothing new.

final() → bytes | None[source]

The last part: what is pending, then the truncation marker and the tail.

class flowli.evidence.EvidenceWriter(port: Evidence, *, limit: int = DEFAULT_LIMIT, flush_interval: float = DEFAULT_FLUSH_INTERVAL, on_error: Callable[[BaseException], None] | None = None)[source]

Bases: object

Buffers the log of an attempt and writes it to the Evidence port.

collecting(ref: EvidenceRef) → AsyncIterator[AttemptBuffer][source]

Capture what this attempt logs, and flush it when the attempt ends.

class flowli.evidence.NullEvidence[source]

Bases: object

The default: accept every write, store nothing.

A process that wants attempt logs configures a real adapter. One that does not pays nothing for the machinery.

async append_log(ref: EvidenceRef, part: bytes) → None[source]
async read_log(ref: EvidenceRef) → bytes[source]
async put(ref: EvidenceRef, name: str, data: bytes, media_type: str) → str[source]
async get(ref: EvidenceRef, name: str) → bytes | None[source]
async list(eid: UUID) → list[EvidenceItem][source]
async delete_for(eid: UUID) → int[source]
flowli.evidence.capture(_logger: Any, _method: str, event: MutableMapping[str, Any]) → MutableMapping[str, Any][source]

A structlog processor. It adds the event to the buffer of this attempt.

It never changes the event and never raises: logging must work whether or not an attempt is running.

flowli.codec

Serialization of domain objects, in one place, outside the domain. Also the content digest of values, since hashing a value means rendering it first.

The domain dataclasses describe data and check intrinsic invariants only. Turning them into JSON-compatible dicts and back is this module’s job, through a cattrs converter with hooks for the three value types that are not JSON-native:

Timestamp <-> RFC 3339 string (Timestamp.to_iso / from_iso) UUID <-> canonical string timedelta <-> seconds (float) datetime -> RFC 3339 string, as a Timestamp (user data in step arguments)

unstructure(obj) gives the dict form. structure(data, cls) parses it back, and is where inbound values are validated and coerced. Adapters, the engine and the patterns use these two functions. Nothing in flowli.domain imports this module.

flowli.codec.canonical_json(data: Any) → str[source]

Deterministic JSON of JSON-native data: sorted keys, no whitespace, no NaN, UTF-8.

Render domain values with unstructure first. Anything else raises TypeError.

flowli.codec.digest(value: Any) → str[source]

sha256 hex of the canonical JSON of unstructure(value).

The frame argument digest of the memo rule (02-journal.md section 5). One rendering rule set, the codec’s, decides how a Timestamp, a UUID or a dataclass looks here.

flowli.codec.structure(data: Any, cls: type[T]) → T[source]

Parse the JSON-compatible form back into cls. Validates and coerces on the way.

flowli.codec.unstructure(obj: Any) → Any[source]

The JSON-compatible form of a domain object (dict, list, or scalar).