Skip to content

reflexr.workspace

Workspaces: tenant-scoped event logs, the storage protocol, and in-memory storage.

Open a Workspace handle with Workspaces.open, and publish, read and operate through it. Every surface hands commands to Workspaces.execute, which carries each out through a handle, once per id. Rules are evaluated and runs executed by the Reactor.

Workspaces

Tenant-scoped handles. See Workspaces and the log.

Workspaces

Workspaces(
    storage: Storage,
    *,
    events: Iterable[type[Event]] | None = None,
    emitted: Iterable[type[Event]] = (),
    registry: EventRegistry = DEFAULT_REGISTRY,
    rules: Iterable[Rule] = (),
    predicates: Predicates | None = None,
    stored_rules: StoredRules | None = None,
    schedules: Iterable[Schedule] = (),
    clock: Clock = utc_now,
    max_depth: int = 8,
    tracer_provider: TracerProvider | None = None,
    meter_provider: MeterProvider | None = None,
    results: CommandResults | None = None,
)

Opens scoped Workspace handles over one storage, and holds the rules.

Parameters:

Name Type Description Default
storage Storage

Where workspaces are kept.

required
events Iterable[type[Event]] | None

The event types clients may publish. Others are rejected even if they are registered, so clients cannot publish arbitrary types. None accepts every type in registry. reflexr's own events can never be published.

None
emitted Iterable[type[Event]]

Event types only runs may publish, such as an incident a triage agent opens. Clients, over REST, the WebSocket or MCP, are refused them.

()
registry EventRegistry

The namespaces whose types these workspaces accept. events and emitted must be in it. Applications in one process each use their own, so neither accepts the other's types.

DEFAULT_REGISTRY
rules Iterable[Rule]

The rules every workspace evaluates. Each is checked against the event types and predicates when the workspaces are created, so a mistake fails at startup. A disabled rule is registered and checked too, but the reactor leaves it be.

()
predicates Predicates | None

The Python predicates rules refer to, by name. They must be pure.

None
stored_rules StoredRules | None

What rules installed in a workspace at runtime may do. Without it, every change to a stored rule is forbidden, and no workspace reads any.

None
schedules Iterable[Schedule]

The schedules that publish ticks into the workspaces.

()
clock Clock

Returns the current time for run transitions. Defaults to the system clock.

utc_now
max_depth int

The deepest causal chain an event may extend: a run's events beyond it are rejected and firings beyond it refused, so workflows that trigger themselves stop.

8
tracer_provider TracerProvider | None

Where spans go. Defaults to OpenTelemetry's global provider.

None
meter_provider MeterProvider | None

Where metrics go. Defaults to OpenTelemetry's global provider.

None
results CommandResults | None

Where execute remembers commands' results. Defaults to the 10,000 most recent, in this process.

None

Raises:

Type Description
InvalidRule

If a rule refers to an event type, field or predicate that does not exist.

ValueError

If two rules, or two schedules, have the same name, a type in events or emitted is not in registry, or a rule is in a namespace of stored rules.

storage property

storage: Storage

Where the workspaces are kept.

rules property

The code rules, which every workspace evaluates, by name.

stored_rules property

stored_rules: StoredRules | None

What stored rules may do, or None if they are off.

schedules property

schedules: Mapping[str, Schedule]

The schedules that publish ticks, by name.

predicates property

predicates: Predicates

The predicates rules refer to, by name.

clock property

clock: Clock

Returns the current time.

max_depth property

max_depth: int

The deepest causal chain an event may extend.

telemetry property

telemetry: Telemetry

The tracer and instruments reflexr records with.

rules_in async

Return a workspace's rules: the code rules, then its active stored rules, by name.

Everything that needs a workspace's rules finds them here, and nothing is cached. Without stored_rules there are only the code rules, and nothing to read; with it, only stored rules in its namespaces, so a namespace a deploy dropped is hidden.

Parameters:

Name Type Description Default
at WorkspaceRef | Transaction

A transaction on the workspace, to read its stored rules as the transaction sees them, or the workspace, to read them with a storage call of their own.

required

open async

open(
    tenant_id: TenantId,
    workspace_id: WorkspaceId,
    *,
    actor: Actor,
    authorize: Authorize | None = None,
) -> Workspace

Return a handle on a tenant's workspace, acting as actor.

This is the only place a tenant id enters; nothing on the handle can reach another tenant. The surfaces pass their authorize hook, so each refuses a workspace alike.

Raises:

Type Description
Forbidden

If authorize is given and refuses the actor this workspace.

execute async

execute(
    workspace: Workspace,
    command: Command,
    *,
    command_id: str,
) -> CommandResult

Carry out a command the way every surface should, once per command_id.

The command is carried out through the handle, as its actor. A rejection is the result, not raised.

A command is known by its tenant, workspace, sender (the handle's actor, as a participant) and command_id. The first time, it is carried out and its result is remembered. A repeated id returns the remembered result and carries nothing out, whatever command it comes with, so a retry is safe on every surface.

schedules_for

schedules_for(tenant_id: TenantId) -> list[Schedule]

Return the schedules that tick in a tenant's workspaces, as the tenant may see them.

A schedule that targets particular workspaces lists only the tenant's own, so no tenant sees another's ids; one that targets none of them is left out.

Workspace

Workspace(
    context: _Context,
    ref: WorkspaceRef,
    actor: Actor,
    cause: _Cause | None = None,
)

A handle on one tenant's workspace, bound to the actor it acts as.

Create handles with Workspaces.open, derive handles for other actors with as_actor, and handles whose events a run caused with caused_by.

tenant_id property

tenant_id: TenantId

The tenant that owns the workspace.

workspace_id property

workspace_id: WorkspaceId

The workspace's id.

actor property

actor: Actor

Who this handle acts as.

telemetry property

telemetry: Telemetry

The tracer and instruments reflexr records with, for actions' own spans.

causation property

causation: Causation | None

What caused the events this handle publishes, if a run did.

as_actor

as_actor(actor: Actor) -> Workspace

Return a handle on the same workspace that acts as actor.

caused_by

caused_by(
    causation: Causation, *, correlation_id: str
) -> Workspace

Return a handle whose events a run caused, in that run's causal chain.

The reactor gives each run such a handle, so the events its action emits record the firing and run behind them and count towards the chain's depth.

publish async

publish(
    event: Event,
    *,
    id: EventId | None = None,
    correlation_id: str | None = None,
) -> Published

Publish an event into the workspace's log.

Parameters:

Name Type Description Default
event Event

The event.

required
id EventId | None

Its id. Publishing an id that is already in the log appends nothing and returns the logged envelope, so producers can retry safely. Defaults to a new id.

None
correlation_id str | None

The causal chain to join: the id of the chain's first event. By default the event starts a new chain, unless this handle belongs to a run.

None

Raises:

Type Description
NotFound

If the event's type is not accepted, or the chain does not exist.

ValidationFailed

If correlation_id names an event that did not start its chain. The message names the chain it belongs to.

Forbidden

If the event is one reflexr records itself.

DepthExceeded

If a run's event would extend its chain beyond the limit.

publish_many async

publish_many(
    events: Sequence[Event],
    *,
    ids: Sequence[EventId | None] | None = None,
    correlation_id: str | None = None,
) -> list[Published]

Publish events atomically and in order: all of them are logged, or none.

Parameters:

Name Type Description Default
events Sequence[Event]

The events.

required
ids Sequence[EventId | None] | None

Their ids, in the same order; None for a new id.

None
correlation_id str | None

The causal chain every event joins, as in publish.

None

Raises:

Type Description
ValueError

If ids is not as long as events.

give_feedback async

give_feedback(
    feedback: Feedback, *, on: FeedbackTarget
) -> Envelope

Record feedback on a run, a firing or a whole causal chain.

The feedback joins the causal chain of what it is about, so it appears in that chain's session.

Raises:

Type Description
ValidationFailed

If the feedback's type cannot be given on this kind of target, or a chain target names an event that did not start its chain.

NotFound

If the target does not exist.

retry_run async

retry_run(run_id: RunId) -> Run

Make a run runnable now.

A dead-lettered, cancelled or skipped run gets a fresh retry budget; a run waiting to retry stops waiting.

Raises:

Type Description
NotFound

If the run does not exist.

InvalidState

If the run succeeded or is running.

skip_run async

skip_run(
    run_id: RunId, *, reason: str | None = None
) -> Run

Give up on a waiting or dead-lettered run, unblocking later runs of its scope.

Raises:

Type Description
NotFound

If the run does not exist.

InvalidState

If the run is running or has finished otherwise.

cancel_run async

cancel_run(
    run_id: RunId, *, reason: str | None = None
) -> Run

Cancel a run that has not finished. A running attempt's result is discarded.

Raises:

Type Description
NotFound

If the run does not exist.

InvalidState

If the run has finished.

checkpoint_run async

checkpoint_run(
    run_id: RunId,
    *,
    attempt: int,
    step: str,
    state: JsonValue,
) -> Run

Save a running run's progress after a step, and append reflexr:run_progressed.

Actions call this through Reaction.checkpoint, so a retry resumes after the last step.

Parameters:

Name Type Description Default
run_id RunId

The run.

required
attempt int

The attempt saving it, which must still be the run's current one.

required
step str

The step that completed.

required
state JsonValue

What resuming needs, as JSON.

required

Raises:

Type Description
NotFound

If the run does not exist.

InvalidState

If the run is no longer running this attempt: it was cancelled, or another executor took it over.

replay_rule async

replay_rule(
    rule: RuleName,
    *,
    from_seq: int = 0,
    mode: Literal["rebuild", "refire"] = "rebuild",
) -> RuleProgress

Reset a rule to evaluate the log again from after from_seq.

Parameters:

Name Type Description Default
rule RuleName

The rule's name.

required
from_seq int

Where to start again; 0 is the beginning of the log.

0
mode Literal['rebuild', 'refire']

"rebuild" recomputes the rule's state up to the head of the log without firing, then carries on as normal. "refire" fires again for everything it finds, creating new runs with new ids.

'rebuild'

Returns:

Type Description
RuleProgress

The rule's new progress. The reactor does the evaluating.

Raises:

Type Description
NotFound

If the workspace has no such rule.

ValidationFailed

If from_seq is not between 0 and the head of the log.

install_rule async

install_rule(
    rule: Rule,
    *,
    provenance: Mapping[str, JsonValue] | None = None,
) -> RuleVersion

Install a stored rule: its first version, or the next version of an archived one.

The rule starts at its reflexr:rule_installed fact, so it acts on what is logged after it. A rule installed again is reset there, continuing its generation, so no firing id repeats.

Parameters:

Name Type Description Default
rule Rule

The rule, in one of the namespaces StoredRules names.

required
provenance Mapping[str, JsonValue] | None

Where it came from, such as the artifact a person accepted: JSON that reflexr keeps on the rule and its fact, and does not read.

None

Returns:

Type Description
RuleVersion

Its new version, or, if it is installed as given already, that version as a

RuleVersion

duplicate.

Raises:

Type Description
Forbidden

If stored rules are off, the name is a code rule's, or allow refuses the change.

ValidationFailed

If check, check_stored or check_provenance finds problems, listing them all, or the workspace has as many active stored rules as it may.

InvalidState

If the rule is installed already: update it instead.

update_rule async

update_rule(
    rule: Rule,
    *,
    expected_version: int | None = None,
    provenance: Mapping[str, JsonValue] | None = None,
) -> RuleVersion

Replace an active stored rule with a new version.

A new condition or scope resets the rule at its reflexr:rule_installed fact, as a changed code rule is reset; any other change keeps its state.

Parameters:

Name Type Description Default
rule Rule

The new version, with the name of the rule it replaces.

required
expected_version int | None

The version it replaces, so a change someone else made first is not overwritten.

None
provenance Mapping[str, JsonValue] | None

Where it came from, as for install_rule.

None

Returns:

Type Description
RuleVersion

Its new version, or, if it is as given already, that version as a duplicate, whatever

RuleVersion

expected_version says, so a retry of an update that succeeded is no conflict.

Raises:

Type Description
Forbidden

As for install_rule.

ValidationFailed

If the checks install_rule makes of a rule find problems.

NotFound

If no stored rule has the name.

InvalidState

If the rule is not at expected_version, which the message names, or is archived: install it again.

archive_rule async

archive_rule(
    rule: RuleName,
    *,
    expected_version: int | None = None,
    reason: str | None = None,
) -> RuleVersion

Archive a stored rule, and cancel its unfinished runs.

Its scope states are cleared and its progress kept, so a rule installed again under its name continues its generation. The runs' reflexr:run_cancelled facts follow the reflexr:rule_archived fact, in the same transaction.

Parameters:

Name Type Description Default
rule RuleName

The stored rule's name.

required
expected_version int | None

The version it archives, as for update_rule.

None
reason str | None

Why, for the fact.

None

Returns:

Type Description
RuleVersion

The archived version, or, if the rule is archived already, that version as a

RuleVersion

duplicate.

Raises:

Type Description
Forbidden

As for install_rule.

NotFound

If no stored rule has the name.

InvalidState

If the rule is not at expected_version, which the message names.

head_seq async

head_seq() -> int

Return the seq of the latest envelope, or 0 if the log is empty.

read async

read(
    *,
    after_seq: int = 0,
    before_seq: int | None = None,
    types: Collection[str] | None = None,
    limit: int | None = None,
    last: int | None = None,
) -> list[Envelope]

Return envelopes in the window after_seq < seq < before_seq, in order.

Parameters:

Name Type Description Default
after_seq int

Only envelopes after this seq.

0
before_seq int | None

Only envelopes before this seq; None reads to the head.

None
types Collection[str] | None

Only envelopes of these event types; None reads every type.

None
limit int | None

At most this many: the first ones in the window that match.

None
last int | None

At most this many: the last ones in the window that match, still in order. To page backwards, read last=n, then last=n before the oldest seq returned.

None

Raises:

Type Description
ValidationFailed

If both limit and last are given, or a number is negative.

subscribe async

subscribe(
    *, after_seq: int = 0
) -> AsyncGenerator[Envelope]

Yield envelopes after after_seq, then each new one as it is logged.

Reading the log is untraced, so a subscription that polls storage makes no trace per poll (ADR-0040); what the subscriber does with each envelope is traced as usual.

cursor async

cursor(name: str) -> int

Return how far a named consumer of the log has got: the seq it saved, or 0.

A consumer, such as a feedback mirror, saves its cursor with save_cursor and carries on after it when it starts again.

save_cursor async

save_cursor(name: str, seq: int) -> None

Save how far a named consumer of the log has got: it is done with seq.

A cursor only moves forward: saving a seq below the saved one leaves it, so a consumer that runs in several processes, each at its own pace, cannot move it back.

run async

run(run_id: RunId) -> Run

Return a run.

Raises:

Type Description
NotFound

If the run does not exist.

runs async

runs(
    *,
    rule: RuleName | None = None,
    status: RunStatus | None = None,
    scope_key: ScopeKey | None = None,
    limit: int | None = None,
) -> list[Run]

Return runs, newest first, optionally of one rule, status or scope.

dead_letters async

dead_letters(
    *, rule: RuleName | None = None
) -> list[EvaluationError]

Return the envelopes rules could not evaluate, oldest first.

schedule_ticks async

schedule_ticks() -> dict[str, datetime]

Return when each schedule last ticked in this workspace.

rule_progress async

rule_progress() -> dict[RuleName, RuleProgress]

Return each rule's progress: its cursor, generation and pending deadlines.

get_rule async

get_rule(rule: RuleName) -> WorkspaceRule

Return one of the workspace's rules, with a stored rule's version and provenance.

Raises:

Type Description
NotFound

If the workspace has no such rule, as for a stored rule that is archived.

rule_statuses async

rule_statuses() -> list[RuleStatus]

Return the status of each of the workspace's rules, code rules first.

Every surface reports this, so they agree. A rule that has not evaluated the workspace yet is at cursor 0, and the progress of a rule the workspace no longer has, such as a code rule no longer registered or an archived stored rule, is left out.

schedule_statuses async

schedule_statuses() -> list[ScheduleStatus]

Return when each schedule that targets this workspace last ticked and ticks next.

Every surface reports this, so they agree. A schedule that has not checked the workspace yet has neither.

Published dataclass

Published(envelope: Envelope, duplicate: bool = False)

The outcome of publishing an event.

envelope instance-attribute

envelope: Envelope

The event's envelope: the new one, or the one already logged with its id.

duplicate class-attribute instance-attribute

duplicate: bool = False

Whether the id was already in the log, so nothing was appended.

RuleVersion dataclass

RuleVersion(
    stored: StoredRule, seq: int, duplicate: bool = False
)

The outcome of a change to a stored rule.

stored instance-attribute

stored: StoredRule

The stored rule as the change left it, or as it already was for a duplicate.

seq instance-attribute

seq: int

The seq of the change's fact, reflexr:rule_installed or reflexr:rule_archived.

For a duplicate, the head of the log: nothing has changed the rule since its version's fact, so reading on from here misses nothing about it.

duplicate class-attribute instance-attribute

duplicate: bool = False

Whether the rule was already as the change would leave it, so nothing changed.

WorkspaceRule pydantic-model

Bases: BaseModel

One of a workspace's rules: a code rule, or the current version of an active stored one.

Config:

  • frozen: True
  • extra: forbid

Fields:

origin pydantic-field

origin: Literal['code', 'stored']

code: registered in code, for every workspace. stored: installed in this one.

version pydantic-field

version: int | None = None

A stored rule's current version. None for a code rule.

provenance pydantic-field

provenance: dict[str, JsonValue] | None = None

Where a stored rule's current version came from, as its installer said. None for a code rule.

RuleStatus pydantic-model

Bases: BaseModel

A rule's progress in a workspace.

Config:

  • frozen: True
  • extra: forbid

Fields:

version pydantic-field

version: int | None = None

A stored rule's current version. None for a code rule.

enabled pydantic-field

enabled: bool

Whether the reactor evaluates the rule. A disabled rule's cursor holds.

lag pydantic-field

lag: int

How many envelopes the rule is behind the head of the log.

ScheduleStatus pydantic-model

Bases: BaseModel

When a schedule last ticked in a workspace, and when it ticks next.

Config:

  • frozen: True
  • extra: forbid

Fields:

Authorize module-attribute

Decides whether an actor may use a workspace of its tenant.

The surfaces that serve workspaces, the FastAPI router and the MCP server, take one as their authorize hook, and pass it to Workspaces.open for every request that names a workspace.

The reactor

Evaluating rules and executing runs. See The reactor.

Reactor

Reactor(
    workspaces: Workspaces,
    *,
    actions: Mapping[str, Action[None]] = ...,
    holder: str | None = None,
    batch_size: int = 500,
    lease_ttl: timedelta = ...,
    concurrency: int = 10,
    run_context: RunContext | None = None,
)
Reactor(
    workspaces: Workspaces,
    *,
    actions: Mapping[str, Action[D]],
    deps: D,
    holder: str | None = None,
    batch_size: int = 500,
    lease_ttl: timedelta = ...,
    concurrency: int = 10,
    run_context: RunContext | None = None,
)
Reactor(
    workspaces: Workspaces,
    *,
    actions: Mapping[str, Action[Any]] | None = None,
    deps: Any = None,
    holder: str | None = None,
    batch_size: int = 500,
    lease_ttl: timedelta = timedelta(seconds=30),
    concurrency: int = 10,
    run_context: RunContext | None = None,
)

Evaluates the rules of a set of Workspaces and executes their runs.

Any number of reactors, in any number of processes, can share the work: leases give each workspace one evaluator and each run one executor at a time.

Parameters:

Name Type Description Default
workspaces Workspaces

The workspaces, with their storage and rules.

required
actions Mapping[str, Action[Any]] | None

What rules run, by the name they refer to them by. Every code rule's action must be here, and its params must validate as the params model the action declares. Every action stored rules may run must be here too, declaring the params model StoredRules gives it.

None
deps Any

The application's dependencies, handed to every action in its Reaction.

None
holder str | None

This reactor's name in leases. Defaults to a new random id; give each process its own.

None
batch_size int

How many envelopes one transaction evaluates for one rule. The transaction holds the workspace's lock, so this bounds how long publishing can wait.

500
lease_ttl timedelta

How long a lease lasts unless renewed. A reactor that dies holding one releases it when it lapses, and its runs are attempted again.

timedelta(seconds=30)
concurrency int

How many runs this reactor executes at once.

10
run_context RunContext | None

Entered around each run attempt, inside its span, such as reflexr.langfuse.langfuse_run to attribute runs in Langfuse.

None

Raises:

Type Description
InvalidRule

If a code rule's action is not among actions, or its params do not validate as the action's params model.

ValueError

If an action stored rules may run is not among actions, or declares another params model than StoredRules gives it.

holder property

holder: str

This reactor's name in leases.

execute async

execute(*, limit: int = 100) -> int

Attempt the runs that are due and can start, up to limit of them.

Returns:

Type Description
int

How many attempts finished, whether they succeeded or failed.

settle async

settle(*, max_rounds: int = 100) -> Settled

Evaluate and execute until nothing more happens now.

Useful in tests and scripts. Runs waiting to retry later are left waiting. A round in which another reactor held a lease this one needed, on a workspace to evaluate or a run to attempt, is not idle, since the other may be about to let that work go: the reactor waits a few milliseconds, longer each time up to a tenth of a second, and tries again.

Raises:

Type Description
RuntimeError

If work is still happening after max_rounds rounds, as when rules keep triggering each other within the depth limit, or another reactor keeps holding a lease this one needs.

serve async

serve(
    *,
    poll_interval: timedelta = timedelta(seconds=1),
    stop: Event | None = None,
    grace: timedelta = _GRACE,
) -> None

Tick, evaluate and execute until stopped, checking for work every poll_interval.

A pass that fails, as when the database is briefly unreachable, is logged and the next one tries again: leases and transactions leave nothing half done.

Setting stop stops the reactor gracefully. Evaluation stops after its current batch, and attempts still waiting for a place do not start. Running attempts have grace to end; then their actions are cancelled, and each attempt is recorded as abandoned, a failed attempt that is retried under the rule's policy. Once every lease it held is released, serve returns. Without stop, it runs until cancelled.

Parameters:

Name Type Description Default
poll_interval timedelta

How long to wait between passes.

timedelta(seconds=1)
stop Event | None

Set it to stop the reactor, as a FastAPI lifespan does at shutdown.

None
grace timedelta

How long running attempts have to end once stop is set.

_GRACE

tick async

tick() -> int

Publish the ticks that are due, for every schedule and the workspaces it targets.

A schedule's first check in a workspace starts its timetable there; its first tick comes at its next time after that.

Returns:

Type Description
int

How many ticks were published.

evaluate async

evaluate(workspace: WorkspaceRef | None = None) -> int

Evaluate every enabled rule over new envelopes until each is caught up with the log.

Parameters:

Name Type Description Default
workspace WorkspaceRef | None

One workspace to evaluate. Defaults to every workspace with a log.

None

Returns:

Type Description
int

How many times rules fired. A workspace whose lease another reactor holds is

int

skipped.

Settled dataclass

Settled(
    ticks: int = 0, firings: int = 0, attempts: int = 0
)

What Reactor.settle did.

ticks class-attribute instance-attribute

ticks: int = 0

How many ticks schedules published.

firings class-attribute instance-attribute

firings: int = 0

How many times rules fired.

attempts class-attribute instance-attribute

attempts: int = 0

How many run attempts finished.

EVALUATION_LEASE module-attribute

EVALUATION_LEASE = 'evaluate'

The lease key a reactor holds on a workspace while it evaluates it.

REACTOR module-attribute

REACTOR = SystemActor(name='reactor')

The actor of the facts the reactor records, about the rules it evaluates and the runs it executes.

Actions

The action port, and what every action receives. See Actions.

Action

Bases: Protocol

What a rule runs: an async callable over a Reaction.

It may return the run's output: JSON-compatible data or a Pydantic model. Raising fails the attempt, which is retried under the rule's policy.

An action that takes params declares their model as a params attribute, as AgentAction, GraphAction and with_params do. The reactor validates each rule's params as that model when it is built, and again at each attempt, into the reaction's params.

Reaction

Reaction(
    *,
    workspace: Workspace,
    run: Run,
    rule: Rule,
    events: tuple[Envelope, ...],
    deps: D,
    params: BaseModel | None = None,
)

Everything an action needs to respond to a firing.

Attributes:

Name Type Description
workspace

The workspace, acting as the run's agent. Events published through it record the run as their cause and join its causal chain.

run

The run: its rule, scope, matched seqs, attempt and chain.

rule

The rule that fired.

events

The envelopes that made the rule fire, in order.

deps

The application's dependencies, given to the reactor.

params

The rule's params, validated as the action's params model, or None if the action declares none. params_as returns them typed.

run_id property

run_id: RunId

The run's id: the firing's id, and the action's idempotency key.

attempt property

attempt: int

Which attempt this is, from 1.

scope property

scope: dict[str, JsonValue]

The values of the rule's scope fields for this firing.

params_as

params_as(model: type[P]) -> P

Return the params as model, the action's params model.

An agent's tools and prompt, and a graph's steps, read their params this way; a function given to with_params gets them as its second argument.

Raises:

Type Description
TypeError

If the params are not a model: the action declares another model, or none.

checkpoint async

checkpoint(step: str, state: JsonValue) -> None

Save the run's progress after step, so a retry resumes after it.

The saved state is reaction.run.checkpoint on the next attempt.

Raises:

Type Description
InvalidState

If this attempt is no longer the run's current one; the action should stop.

emit async

emit(event: Event) -> Published

Publish an event caused by this run.

Its id is derived from the run, its last checkpoint, and how many events it emitted since, so a retried attempt that emits the same events again adds nothing to the log, and a graph resumed after a checkpoint never reuses an earlier event's id.

Raises:

Type Description
DepthExceeded

If the event would extend the causal chain beyond the limit.

with_params

with_params(
    fn: Callable[[Reaction[D], P], Awaitable[object]],
    params: type[P],
) -> Action[D]

Adapt a function of the reaction and its params into an action that declares params.

pyright checks that the function's second parameter takes a params::

class NotifyParams(BaseModel):
    thread_id: str


async def notify(reaction: Reaction[AppDeps], params: NotifyParams) -> None:
    await reaction.deps.bridge.notice(params.thread_id, reaction)


actions = {"notify": with_params(notify, NotifyParams)}

RunFailure

RunFailure(
    message: str, *, reason: str, permanent: bool = False
)

Bases: Exception

Raise from an action, a tool or a capability to fail the attempt with a reason.

The reason is a stable code, such as guardrail_blocked, recorded on the run and its facts and counted in reflexr.runs; the message is for people. A permanent failure is not retried, whatever the rule's policy, since retrying the same input would fail again. Other exceptions fail the attempt too, without a reason, and are retried.

Parameters:

Name Type Description Default
message str

What went wrong.

required
reason str

A stable code for why.

required
permanent bool

Whether retrying cannot help.

False

RunContext

A context entered around each run attempt, inside its span, given the attempt's reaction.

It is a port (ADR-0025): a backend that attributes a run in its own way, such as Langfuse's propagated trace attributes, implements it, and the reactor stays free of the backend. It matches artifactr's TurnContext.

Commands

Where Workspaces.execute, the one handler every surface hands commands to, remembers their results. See Deduplication.

CommandResults

Bases: Protocol

Where workspaces remember the results of commands: a port (ADR-0025).

InMemoryCommandResults remembers them in the process. An implementation over shared storage would deduplicate across processes too.

get async

get(key: CommandKey) -> CommandResult | None

Return the result remembered for a command, if there is one.

put async

put(key: CommandKey, result: CommandResult) -> None

Remember a command's result.

InMemoryCommandResults

InMemoryCommandResults(capacity: int = 10000)

The most recent results, in this process: the default CommandResults.

Parameters:

Name Type Description Default
capacity int

How many results to remember. Beyond it, the oldest is forgotten first, and a command whose result was forgotten is carried out again if it is repeated.

10000

get async

get(key: CommandKey) -> CommandResult | None

Return the result remembered for a command, if it is still remembered.

put async

put(key: CommandKey, result: CommandResult) -> None

Remember a command's result, forgetting the oldest beyond the capacity.

CommandKey dataclass

CommandKey(
    tenant_id: TenantId,
    workspace_id: WorkspaceId,
    participant: str,
    command_id: str,
)

What makes two commands the same: who sent them, where, and with which id.

tenant_id instance-attribute

tenant_id: TenantId

The tenant of the workspace the command was sent to.

workspace_id instance-attribute

workspace_id: WorkspaceId

The workspace it was sent to.

participant instance-attribute

participant: str

The sender, as its actor's participant: every action of one person or client.

command_id instance-attribute

command_id: str

The id the client chose for the command.

Schedules

Ticks on a timetable. See Schedules.

Schedule pydantic-model

Bases: BaseModel

When to publish ticks, and into which workspaces.

Give every or cron, not both::

Schedule(name="heartbeat-check", every=timedelta(seconds=30))
Schedule(name="weekday-standup", cron="0 9 * * mon-fri", timezone="Europe/Oslo")

Config:

  • frozen: True
  • extra: forbid

Fields:

Validators:

  • _check

cron pydantic-field

cron: str | None = None

A five-field cron expression, evaluated in timezone.

workspaces pydantic-field

workspaces: (
    Literal["all"]
    | tuple[tuple[TenantId, WorkspaceId], ...]
) = "all"

"all" for every workspace with a log, or (tenant, workspace) pairs.

catch_up pydantic-field

catch_up: Literal['skip', 'all'] = 'skip'

After downtime, publish only the latest missed tick ("skip"), or every one.

max_catch_up pydantic-field

max_catch_up: int = 100

The most ticks one catch-up publishes; older missed ticks are skipped.

after

after(moment: AwareDatetime) -> Iterator[datetime]

Yield the schedule's times strictly after moment, in order.

due

due(
    last: AwareDatetime, now: AwareDatetime
) -> list[datetime]

Return the ticks due after last and by now, as catch_up allows.

targets

targets(
    tenant_id: TenantId, workspace_id: WorkspaceId
) -> bool

Whether the schedule ticks in a workspace.

SCHEDULER module-attribute

SCHEDULER = SystemActor(name='scheduler')

The actor of the ticks schedules publish.

tick_id

tick_id(schedule: str, at: datetime) -> EventId

Return the id of a schedule's tick at at, the same in every replica.

Storage

The storage protocol, and the in-memory implementation. See Storage.

Storage

Bases: Protocol

Persistence for workspaces. Every method except discovery is scoped to one workspace.

A caller can be cancelled at any await, by a rule's timeout, a stopped reactor or a closed connection, so every method is cancel-safe: a cancelled caller leaves no lock or connection behind; cancellation may be deferred until the current statement ends. So a cancelled call may still have taken effect: a transaction may have committed, or acquire_lease taken the lease. A wait for a pooled connection cannot be interrupted either, so a task that holds a transaction must not await a task it cancelled. The event loop's shutdown is not covered: it cancels the storage's own tasks too, so stop the reactor with serve(stop=) first.

transaction

transaction(
    workspace: WorkspaceRef,
) -> AbstractAsyncContextManager[Transaction]

Begin a transaction: it commits if the block exits normally, and rolls back if not.

head_seq async

head_seq(workspace: WorkspaceRef) -> int

Return the log's latest seq, or 0 if it is empty.

read async

read(
    workspace: WorkspaceRef,
    *,
    after_seq: int = 0,
    before_seq: int | None = None,
    types: Collection[str] | None = None,
    limit: int | None = None,
    last: int | None = None,
) -> list[Envelope]

Return logged envelopes in the window after_seq < seq < before_seq, in order.

Parameters:

Name Type Description Default
workspace WorkspaceRef

The workspace whose log to read.

required
after_seq int

Only envelopes after this seq.

0
before_seq int | None

Only envelopes before this seq; None reads to the head.

None
types Collection[str] | None

Only envelopes of these event types; None reads every type.

None
limit int | None

At most this many: the first ones in the window that match.

None
last int | None

At most this many: the last ones in the window that match, still returned in order. This is the tail of the log.

None

Callers give at most one of limit and last, and no negative numbers; read checks.

subscribe

subscribe(
    workspace: WorkspaceRef, *, after_seq: int = 0
) -> AsyncGenerator[Envelope]

Yield envelopes after after_seq: the stored ones, then each one as it commits.

run async

run(workspace: WorkspaceRef, run_id: RunId) -> Run | None

Return a run, or None.

runs async

runs(
    workspace: WorkspaceRef,
    *,
    rule: RuleName | None = None,
    status: RunStatus | None = None,
    scope_key: ScopeKey | None = None,
    limit: int | None = None,
) -> list[Run]

Return runs, newest first, optionally filtered.

dead_letters async

dead_letters(
    workspace: WorkspaceRef, *, rule: RuleName | None = None
) -> list[EvaluationError]

Return the envelopes rules could not evaluate, oldest first.

progress async

progress(
    workspace: WorkspaceRef,
) -> dict[RuleName, RuleProgress]

Return every rule's progress in a workspace.

schedules async

schedules(workspace: WorkspaceRef) -> dict[str, datetime]

Return the time of each schedule's last tick in a workspace.

stored_rules async

stored_rules(workspace: WorkspaceRef) -> list[StoredRule]

Return a workspace's active stored rules, by name, in code-point order.

workspaces async

workspaces() -> list[WorkspaceRef]

Return every workspace with a log, for the reactor to evaluate.

due_runs async

due_runs(
    *, now: datetime, limit: int, policy: RunPolicy
) -> list[tuple[WorkspaceRef, Run]]

Return the runs to attempt now, across workspaces, oldest first.

These are the pending and retrying runs due by now that can start, and running runs whose run_lease has lapsed, because their executor stopped. A pending or retrying run of an ordered rule can start only if no earlier run of its scope, in firing order (fired_seq, then creation), holds the scope, so a scope's backlog yields one run however deep it is. Runs of the disabled rules are left out. So neither the runs waiting behind others nor those of disabled rules take up the limit.

acquire_lease async

acquire_lease(
    workspace: WorkspaceRef,
    key: str,
    holder: str,
    ttl: timedelta,
) -> bool

Take or renew an exclusive lease. Return False if another holder has it.

release_lease async

release_lease(
    workspace: WorkspaceRef, key: str, holder: str
) -> None

Release a lease if holder has it.

cursor async

cursor(workspace: WorkspaceRef, name: str) -> int

Return how far a named consumer of the log has got: the seq saved, or 0.

save_cursor async

save_cursor(
    workspace: WorkspaceRef, name: str, seq: int
) -> None

Save how far a named consumer of the log has got, such as a feedback mirror.

A cursor only moves forward: saving a seq below the saved one leaves it, so a consumer running in several processes cannot move it back.

Transaction

Bases: Protocol

One atomic unit of work on a workspace.

A transaction must be serialized with every other transaction on the same workspace from the moment it begins, so nothing it reads can change before it commits. In-memory storage holds a per-workspace lock; SQL storage locks the workspace's row. What a transaction writes is visible to its own later reads.

append async

append(entries: Sequence[Entry]) -> list[Envelope]

Append events to the log, assigning each the next seq and a ts.

ts is the later of the clock and the previous envelope's ts, so it never decreases within a workspace. Callers make publishing idempotent by looking an id up with envelope first; appending an id that is already in the log, or earlier in the same transaction, raises ValueError.

envelope async

envelope(event_id: EventId) -> Envelope | None

Return the envelope of an event id, or None.

head_seq async

head_seq() -> int

Return the log's latest seq, or 0 if it is empty.

read async

read(*, after_seq: int, limit: int) -> list[Envelope]

Return up to limit envelopes after after_seq, in order.

progress async

progress(rule: RuleName) -> RuleProgress | None

Return a rule's progress in this workspace, or None if it has none yet.

save_progress async

save_progress(
    rule: RuleName, progress: RuleProgress
) -> None

Save a rule's progress.

states async

states(
    rule: RuleName, keys: Collection[ScopeKey]
) -> dict[ScopeKey, ScopeState]

Return the stored state of the given scopes of a rule; missing scopes are left out.

save_states async

save_states(
    rule: RuleName, states: Mapping[ScopeKey, ScopeState]
) -> None

Save the state of scopes of a rule.

clear_states async

clear_states(rule: RuleName) -> None

Delete every scope state of a rule, when it is reset.

run async

run(run_id: RunId) -> Run | None

Return a run, or None.

save_runs async

save_runs(runs: Sequence[Run]) -> None

Insert or update runs.

scope_runs async

scope_runs(
    rule: RuleName, scope_key: ScopeKey
) -> list[Run]

Return every run of a rule's scope that has not succeeded, in firing order.

dead_letter async

dead_letter(errors: Sequence[EvaluationError]) -> None

Record envelopes a rule could not evaluate.

schedule async

schedule(name: str) -> datetime | None

Return the time of a schedule's last tick in this workspace, or None.

save_schedule async

save_schedule(name: str, at: datetime) -> None

Save the time of a schedule's last tick.

stored_rule async

stored_rule(name: RuleName) -> StoredRule | None

Return a stored rule, active or archived, or None if no rule of that name is stored.

stored_rules async

stored_rules() -> list[StoredRule]

Return the workspace's active stored rules, by name, in code-point order.

save_stored_rule async

save_stored_rule(stored: StoredRule) -> None

Save a stored rule's next version, which archives it if its status is archived.

Versions are numbered from 1 per name, so callers check an expected version by reading stored_rule first; saving any version but the next raises ValueError.

Entry dataclass

Entry(
    id: EventId,
    actor: Actor,
    event: Event,
    causation: Causation | None = None,
    correlation_id: str | None = None,
    traceparent: str | None = None,
)

An event to append. Storage assigns its seq and ts.

correlation_id class-attribute instance-attribute

correlation_id: str | None = None

The causal chain it belongs to; None starts a new chain with the event's own id.

WorkspaceRef dataclass

WorkspaceRef(
    tenant_id: TenantId, workspace_id: WorkspaceId
)

A tenant's workspace: the unit of isolation, ordering and locking.

(Not to be confused with a rule's Scope, which groups a rule's state within a workspace.)

run_lease

run_lease(run_id: RunId) -> str

Return the lease key an executor holds on a run while it runs it.

RUN_LEASE_PREFIX module-attribute

RUN_LEASE_PREFIX = 'run:'

The start of every run's lease key, which ends with the run's id.

RunPolicy dataclass

RunPolicy(
    disabled: frozenset[RuleName] = frozenset(),
    ordered: frozenset[RuleName] = frozenset(),
    blocking: frozenset[RuleName] = frozenset(),
)

What Storage.due_runs needs to know about the rules to find runs that can start.

Storage does not see the rules, so the executor describes them with of. A rule named in no set has all of its due runs returned, which is always safe: the executor checks every run again as it claims it, so a set that says too little costs claims that are declined, never a run that is not attempted.

disabled class-attribute instance-attribute

disabled: frozenset[RuleName] = frozenset()

Rules whose runs wait, and are left out altogether.

ordered class-attribute instance-attribute

Rules whose runs execute in firing order per scope (ordering="scope").

blocking class-attribute instance-attribute

blocking: frozenset[RuleName] = frozenset()

Ordered rules whose dead-lettered runs hold back their scope (on_dead_letter="block").

of classmethod

of(rules: Iterable[Rule]) -> Self

Describe the given rules.

holds

holds(run: Run) -> bool

Whether a run holds back the later runs of its scope (holds).

InMemoryStorage

InMemoryStorage(*, clock: Clock = utc_now)

Storage that keeps every workspace in process memory.

Parameters:

Name Type Description Default
clock Clock

Returns the current time, for envelope timestamps and lease expiry. Defaults to the system clock in UTC.

utc_now

transaction async

transaction(
    workspace: WorkspaceRef,
) -> AsyncGenerator[_Transaction]

Begin a transaction; it holds the workspace's lock until it ends.

head_seq async

head_seq(workspace: WorkspaceRef) -> int

Return the log's latest seq, or 0 if it is empty.

read async

read(
    workspace: WorkspaceRef,
    *,
    after_seq: int = 0,
    before_seq: int | None = None,
    types: Collection[str] | None = None,
    limit: int | None = None,
    last: int | None = None,
) -> list[Envelope]

Return logged envelopes in the window after_seq < seq < before_seq, in order.

Of those of types, if given: the first limit or the last last.

subscribe async

subscribe(
    workspace: WorkspaceRef, *, after_seq: int = 0
) -> AsyncGenerator[Envelope]

Yield stored envelopes after after_seq, then each new one as it commits.

run async

run(workspace: WorkspaceRef, run_id: RunId) -> Run | None

Return a run, or None.

runs async

runs(
    workspace: WorkspaceRef,
    *,
    rule: RuleName | None = None,
    status: RunStatus | None = None,
    scope_key: ScopeKey | None = None,
    limit: int | None = None,
) -> list[Run]

Return runs, newest first, optionally filtered.

dead_letters async

dead_letters(
    workspace: WorkspaceRef, *, rule: RuleName | None = None
) -> list[EvaluationError]

Return the envelopes rules could not evaluate, oldest first.

progress async

progress(
    workspace: WorkspaceRef,
) -> dict[RuleName, RuleProgress]

Return every rule's progress in a workspace.

schedules async

schedules(workspace: WorkspaceRef) -> dict[str, datetime]

Return the time of each schedule's last tick in a workspace.

stored_rules async

stored_rules(workspace: WorkspaceRef) -> list[StoredRule]

Return a workspace's active stored rules, by name, in code-point order.

workspaces async

workspaces() -> list[WorkspaceRef]

Return every workspace with a log.

due_runs async

due_runs(
    *, now: datetime, limit: int, policy: RunPolicy
) -> list[tuple[WorkspaceRef, Run]]

Return the runs to attempt now that can start, across workspaces, oldest first.

acquire_lease async

acquire_lease(
    workspace: WorkspaceRef,
    key: str,
    holder: str,
    ttl: timedelta,
) -> bool

Take or renew an exclusive lease. Return False if another holder has it.

release_lease async

release_lease(
    workspace: WorkspaceRef, key: str, holder: str
) -> None

Release a lease if holder has it.

cursor async

cursor(workspace: WorkspaceRef, name: str) -> int

Return how far a named consumer of the log has got: the seq saved, or 0.

save_cursor async

save_cursor(
    workspace: WorkspaceRef, name: str, seq: int
) -> None

Save how far a named consumer of the log has got; a cursor only moves forward.

Clock module-attribute

Clock = Callable[[], datetime]

Returns the current time; injectable so tests control time.

utc_now

utc_now() -> datetime

Return the current time in UTC: the default clock.