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- 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 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 journal(eid: UUID) list[Sequenced[Entry]][source]¶
The live journal, or the archived one when retention folded it.
- async status(eid: UUID) ExecutionStatus[source]¶
Derived from the last execution.* entry of the journal.
- workflow_ref(execution: Execution) WorkflowRef[source]¶
- 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
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- 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_until(instant: Timestamp, *, name: str = 'sleep_until', key: str | None = None) Awaitable[None][source]¶
- class flowli.runtime.Wait(fid: str, on: str, deadline: Timestamp | None)[source]¶
Bases:
objectA frame that waits. The worker registers it. See 05-protocols.md section 3.
- class flowli.runtime.Raised(error: 'BaseException')[source]¶
Bases:
object- error: BaseException¶
Registry¶
Worker, sweeper and retention¶
- class flowli.runtime.Worker(engine: Engine, queues: list[str], worker_id: str)[source]¶
Bases:
object- async run_once() bool[source]¶
Dequeue at most one task and process it. Return True when a task was processed.
- 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]¶
- 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
- 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 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
- class flowli.runtime.ControlSource(*args, **kwargs)[source]¶
Bases:
ProtocolWhere 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:
objectA 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]¶
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:
objectTakes 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.
- 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:
objectSee specs/10-agent-runner.md section 11.
- class flowli.runtime.ConsumerReport(delivered: 'list[str]' = <factory>, refused: 'list[str]' = <factory>, recovered: 'list[str]' = <factory>, terminated: 'list[str]' = <factory>)[source]¶
Bases:
object
- class flowli.runtime.Handler(*args, **kwargs)[source]¶
Bases:
ProtocolThe 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 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.
- class flowli.runtime.Held(claimed: ClaimedTask, task: DelegateTask, adopted: bool = False)[source]¶
Bases:
objectOne 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¶
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- at: Timestamp¶
- class flowli.patterns.Reviews(engine: Engine, *, take_ttl: float = 30.0)[source]¶
Bases:
objectOperator 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.
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- created_by: Provenance¶
- class flowli.domain.ExecutionStatus(*values)[source]¶
Bases:
StrEnum- PENDING = 'pending'¶
- RUNNING = 'running'¶
- SUSPENDED = 'suspended'¶
- COMPLETED = 'completed'¶
- FAILED = 'failed'¶
- CANCELLED = 'cancelled'¶
- class flowli.domain.FrameRef(eid: UUID, fid: str)[source]¶
Bases:
objectOne frame of one execution. Everything derived from that pair is derived here.
- class flowli.domain.Frame(ref: 'FrameRef', kind: 'FrameKind', name: 'str', args_digest: 'str', retry: 'RetryPolicy' = <factory>)[source]¶
Bases:
object- 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
- class flowli.domain.Attempt(frame: 'FrameRef', number: 'int', provenance: 'Provenance')[source]¶
Bases:
object- provenance: Provenance¶
The journal¶
- class flowli.domain.Entry(type: str, fid: str, payload: dict[str, Any], provenance: Provenance)[source]¶
Bases:
objectOne event in the journal or in the control log.
type is an EntryType value or an ‘announce.*’ string.
- provenance: Provenance¶
- 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_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_completed(prov: Provenance, value: Any) Entry[source]¶
- classmethod execution_failed(prov: Provenance, failed: Failed) Entry[source]¶
- 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:
objectDerived state of one journal. See specs/02-journal.md section 4.
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:
objectA 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.
- enqueued_by: Provenance¶
- 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:
objectThe 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.
- class flowli.domain.Message(channel: 'str', seq: 'int', payload: 'Any', sent_by: 'Provenance', correlation: 'str | None' = None)[source]¶
Bases:
object- sent_by: Provenance¶
- class flowli.domain.Timer(due_at: Timestamp, target: FrameRef)[source]¶
Bases:
objectA 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¶
Provenance¶
- class flowli.domain.Actor(kind: 'ActorKind', id: 'str', on_behalf_of: 'str | None' = None)[source]¶
Bases:
object
- 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
Errors¶
Domain errors. See specs/04-api.md section 4.
- exception flowli.domain.errors.InvalidName[source]¶
Bases:
FlowliError,ValueErrorAn identifier does not match its required pattern.
- exception flowli.domain.errors.NondeterminismError(fid: str, recorded: str, computed: str)[source]¶
Bases:
FlowliErrorReplay found a frame id with a different args digest.
- exception flowli.domain.errors.DuplicateFrameError(fid: str)[source]¶
Bases:
FlowliErrorA frame key repeats under one parent.
- exception flowli.domain.errors.ChildFailed(child_eid: Any, status: str, error: str | None = None)[source]¶
Bases:
FlowliErrorA child execution failed or was cancelled.
- exception flowli.domain.errors.LeaseLost[source]¶
Bases:
FlowliErrorThe execution lease was stolen. Stop all work on the execution.
- exception flowli.domain.errors.Cancelled[source]¶
Bases:
FlowliErrorRaised 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:
ExceptionA 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:
FlowliErrorA 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:
ProtocolFenced ownership. renew() raises LeaseLost when fenced.
- class flowli.domain.ports.Journal(*args, **kwargs)[source]¶
Bases:
ProtocolTrace and memo of one execution. Spec 03 section 3.
- class flowli.domain.ports.ControlLog(*args, **kwargs)[source]¶
Bases:
ProtocolLifecycle entries of all executions. Spec 03 section 4.
- class flowli.domain.ports.LeaseInfo(epoch: int, holder: str | None, deadline_at: Timestamp | None, released: bool, state: Any)[source]¶
Bases:
objectA read of a lease document, without acquisition.
- class flowli.domain.ports.Ownership(*args, **kwargs)[source]¶
Bases:
ProtocolOne owner per execution. Spec 03 section 5.
- async request_cancel(eid: UUID, by: Provenance) None[source]¶
- class flowli.domain.ports.Dispatch(*args, **kwargs)[source]¶
Bases:
ProtocolIdempotent start, and generic exactly-one claims. Spec 03 section 6.
- class flowli.domain.ports.QueueDepth(total: int, visible: int, claimed: int)[source]¶
Bases:
objectWhat a queue holds now. See specs/03-ports.md section 7.
- class flowli.domain.ports.ClaimedTask(task: 'Task', key: 'str', lease: 'Lease')[source]¶
Bases:
object
- class flowli.domain.ports.Queue(*args, **kwargs)[source]¶
Bases:
ProtocolTasks. 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 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.
- class flowli.domain.ports.Channel(*args, **kwargs)[source]¶
Bases:
ProtocolMessages. Spec 03 section 8.
- async all_waits() list[tuple[str, FrameRef]][source]¶
Every wait marker, as (channel, ref). For the sweeper.
- class flowli.domain.ports.Timers(*args, **kwargs)[source]¶
Bases:
ProtocolFuture resumes. Spec 03 section 9.
- class flowli.domain.ports.ExecutionStore(*args, **kwargs)[source]¶
Bases:
ProtocolThe wf/exec/{eid}/meta object. Put-if-absent, then CAS for migrate.
- class flowli.domain.ports.Archive(*args, **kwargs)[source]¶
Bases:
ProtocolFolded journals of finished executions. wf/archive/{eid}. Put-if-absent.
- class flowli.domain.ports.EvidenceRef(eid: UUID, fid: str, attempt: int)[source]¶
Bases:
objectOne attempt of one frame. The key of evidence is the attempt, because a retry writes a second log.
- class flowli.domain.ports.EvidenceItem(ref: 'EvidenceRef', name: 'str', media_type: 'str', size: 'int')[source]¶
Bases:
object- ref: EvidenceRef¶
- class flowli.domain.ports.Evidence(*args, **kwargs)[source]¶
Bases:
ProtocolWhat 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 list(eid: UUID) list[EvidenceItem][source]¶
- 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:
objectEverything the engine needs, bundled.
- control: ControlLog¶
- executions: ExecutionStore¶
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.adapters¶
CairnDB¶
- class flowli.adapters.cairndb.CairnBackend(db: CairnDB, *, now: Callable[[], Timestamp] = Timestamp.now, log_cache_size: int = 128)[source]¶
Bases:
objectAll ports over one CairnDB engine.
- 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'¶
- 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]¶
- 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.
- 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.
- 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:
objectAll ports over one shared clock. backend.ports is what the engine takes.
- clock: ManualClock¶
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:
objectSee specs/09-http-api.md section 3.
- class flowli.api.Catalog(registry: Registry, *, default_queue: str = 'default')[source]¶
Bases:
objectEvery registered workflow, with its argument schema.
Built once, at startup: a schema never changes while the process runs.
- class flowli.api.Principal(actor: Actor, capabilities: frozenset[str] = <factory>)[source]¶
Bases:
objectWho is calling, and what they may do.
- class flowli.api.StaticAuthenticator(principals: Mapping[str, Principal])[source]¶
Bases:
objectA 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.
- class flowli.api.OIDCAuthenticator(config: OIDCConfig, *, jwk_client: Any | None = None)[source]¶
Bases:
objectBearer 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.
- 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
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.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:
objectThe 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¶
- 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:
objectBuffers 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:
objectThe 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 list(eid: UUID) list[EvidenceItem][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.