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
|
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. |
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 |
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 |
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.
rules_in
async
¶
rules_in(
at: WorkspaceRef | Transaction,
) -> dict[RuleName, WorkspaceRule]
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 |
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
¶
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.
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
¶
Return a handle on the same workspace that acts as actor.
caused_by
¶
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 |
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
|
correlation_id
|
str | None
|
The causal chain every event joins, as in |
None
|
Raises:
| Type | Description |
|---|---|
ValueError
|
If |
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
¶
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
¶
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 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
¶
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'
|
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 |
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 |
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 |
ValidationFailed
|
If |
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 |
None
|
Returns:
| Type | Description |
|---|---|
RuleVersion
|
Its new version, or, if it is as given already, that version as a duplicate, whatever |
RuleVersion
|
|
Raises:
| Type | Description |
|---|---|
Forbidden
|
As for |
ValidationFailed
|
If the checks |
NotFound
|
If no stored rule has the name. |
InvalidState
|
If the rule is not at |
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 |
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 |
NotFound
|
If no stored rule has the name. |
InvalidState
|
If the rule is not at |
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 |
0
|
before_seq
|
int | None
|
Only envelopes before this |
None
|
types
|
Collection[str] | None
|
Only envelopes of these event types; |
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 |
None
|
Raises:
| Type | Description |
|---|---|
ValidationFailed
|
If both |
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
¶
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 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
¶
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
¶
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
¶
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.
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:Trueextra:forbid
Fields:
RuleStatus
pydantic-model
¶
Bases: BaseModel
A rule's progress in a workspace.
Config:
frozen:Trueextra:forbid
Fields:
ScheduleStatus
pydantic-model
¶
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
|
None
|
deps
|
Any
|
The application's dependencies, handed to every action in its
|
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
|
None
|
Raises:
| Type | Description |
|---|---|
InvalidRule
|
If a code rule's action is not among |
ValueError
|
If an action stored rules may run is not among |
execute
async
¶
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
¶
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 |
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 |
_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
¶
What Reactor.settle did.
EVALUATION_LEASE
module-attribute
¶
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 |
|
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
¶
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 |
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
¶
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
¶
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
¶
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
¶
RunContext = Callable[
[Reaction[Any]], AbstractAsyncContextManager[object]
]
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.
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.
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:Trueextra:forbid
Fields:
-
name(str) -
every(timedelta | None) -
cron(str | None) -
timezone(str) -
workspaces(Literal['all'] | tuple[tuple[TenantId, WorkspaceId], ...]) -
catch_up(Literal['skip', 'all']) -
max_catch_up(int)
Validators:
-
_check
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
¶
Yield the schedule's times strictly after moment, in order.
due
¶
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
¶
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 |
0
|
before_seq
|
int | None
|
Only envelopes before this |
None
|
types
|
Collection[str] | None
|
Only envelopes of these event types; |
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.
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
¶
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 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
¶
Return the envelope of an event id, or None.
read
async
¶
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.
scope_runs
async
¶
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
¶
Return the time of a schedule's last tick in this workspace, or None.
save_schedule
async
¶
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
¶
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
¶
Return the lease key an executor holds on a run while it runs it.
RUN_LEASE_PREFIX
module-attribute
¶
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
¶
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
¶
Ordered rules whose dead-lettered runs hold back their scope (on_dead_letter="block").
InMemoryStorage
¶
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.
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.
due_runs
async
¶
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
¶
Returns the current time; injectable so tests control time.