Workspaces and the log¶
A workspace is where a tenant's events live: one append-only log with a gap-free seq, and the progress, runs and dead letters of the rules that watch it. It is the unit of ordering, isolation and scale, exactly as in artifactr (ADR-0016). Every read and write goes through a Workspace handle. This page covers opening handles, publishing, causal chains, reading the log, operating runs and rules, and the ways a command can be rejected.
Opening a workspace¶
Workspaces holds everything the workspaces share, checks it once, and opens handles over one storage:
from reflexr import SourceActor, UserActor
from reflexr.workspace import InMemoryStorage, Workspaces
workspaces = Workspaces(
InMemoryStorage(),
events=[ServiceError, Deploy, Heartbeat, IncidentOpened],
rules=[error_spike, deploy_regression, heartbeat_lost],
)
monitoring = await workspaces.open("acme", "prod", actor=SourceActor(name="monitoring"))
ada = monitoring.as_actor(UserActor(id="ada", name="Ada"))
| Argument | Default | Meaning |
|---|---|---|
storage |
required | Where workspaces are kept (Storage) |
events |
every type in registry |
The event types clients may publish: an allowlist (Events and envelopes) |
emitted |
none | Event types only runs may publish, such as an incident an agent opens; clients are refused them |
registry |
DEFAULT_REGISTRY |
The namespaces whose types the workspaces accept, so applications in one process stay apart (Registries) |
rules |
none | The rules every workspace evaluates, checked when Workspaces is built |
predicates |
none | The Python predicates rules refer to, by name |
schedules |
none | The schedules that publish ticks into the workspaces |
clock |
UTC system time | The time for run transitions; tests pass a clock they control |
max_depth |
8 | The deepest causal chain an event may extend (Causation depth) |
tracer_provider, meter_provider |
OpenTelemetry's global ones | Where spans and metrics go (Observability) |
results |
the 10,000 most recent, in the process | Where execute remembers commands' results, so a retried command is carried out once (Deduplication) |
open(tenant_id, workspace_id, actor=...) is the only place a tenant id enters. The handle it returns is bound to that tenant, that workspace and one actor, and nothing on it can reach another tenant. Workspaces need no creating: opening one that has never been written to gives an empty workspace, and it appears to the reactor once it has a log. Handles are cheap, so open one per request or connection, for the actor making it. as_actor(actor) returns a handle on the same workspace acting as someone else, and every write through a handle is attributed to its actor. open also takes an authorize(tenant_id, workspace_id, actor) hook, and raises Forbidden if it refuses the actor the workspace: the surfaces pass theirs, so REST, the WebSocket and MCP refuse a workspace alike (Multi-tenancy and security).
Publishing events¶
publish appends one event, and publish_many several, atomically and in order:
published = await monitoring.publish(
ServiceError(service="auth", severity=8, message="token check failed"), id="alert-7"
)
batch = await monitoring.publish_many(
[Deploy(service="auth", version="2.4.1"), Heartbeat(service="auth")],
ids=["deploy-91", None], # None: generate an id
)
Both return Published results: the envelope that is in the log, and whether the id was a duplicate, so nothing was appended. Publishing is idempotent by id (Events and envelopes). Before appending, a handle checks that the event's type is in the allowlist (NotFound otherwise), that it is not one of reflexr's own events (Forbidden), that a handle no run caused is not publishing a run-only type from emitted (Forbidden), and, for a run's handle, that the causal chain is not too deep (DepthExceeded). Each publish is traced as a reflexr.publish {type} span, whose trace context is stored on the envelope (Observability).
The same publish is available to producers outside the process over REST and the WebSocket and MCP. Every surface hands commands to one handler, workspaces.execute(workspace, command, command_id=...), which you can call too:
from reflexr.core import Publish
result = await workspaces.execute(
ada, Publish(event=Heartbeat(service="billing"), id="hb-1"), command_id="c_1"
)
print(result.ok, result.outcome) # True type='published' seq=4 id='hb-1' duplicate=False
It returns the command's command_result and never raises: a rejection is in its rejection, as the payload REST sends. The first result for each command_id of a participant in a workspace is remembered, and a repeated id returns it without carrying anything out (Deduplication).
Causal chains¶
Every envelope belongs to a causal chain, named by its correlation_id: the id of the chain's first event. The chain is how reflexr connects a cause with everything it led to, and it is the session a run's traces are filed under (Sessions and causal chains).
- An event nothing caused starts a chain with its own id.
- A firing joins the chain of the latest envelope it matched, and its run carries that chain (ADR-0024). For an absence, that is the last envelope the rule saw.
- Events a run publishes continue its chain. The reactor gives each run a handle made with
caused_by(causation, correlation_id=...), so its events record the firing and run that caused them, one step deeper in the chain. - reflexr's facts about a firing or run, operators' actions on a run, and feedback on it all join the run's chain.
- A producer can join an existing chain with
publish(..., correlation_id="alert-7"), naming the chain by the id of its first event. An id that is not in the log is rejected withNotFound, and the id of a later event in a chain withValidationFailed, whose message names the chain that event belongs to, so the producer can join it by that id. Feedback on aChainTargetnames its chain the same way.
With a rule that runs triage on every severe error, and a triage action that opens an incident, the whole story of one alert is one chain:
from reflexr import F, Rule, by, on, run
from reflexr.workspace import Reaction, Reactor
async def triage(reaction: Reaction[None]) -> str:
await reaction.emit(IncidentOpened(service=str(reaction.scope["service"]), summary="errors"))
return "incident opened"
rule = Rule(
name="ops:error-spike",
when=on(ServiceError).where(F.severity >= 7),
scope=by(F.service),
then=run("triage"),
)
workspaces = Workspaces(InMemoryStorage(), events=[ServiceError, IncidentOpened], rules=[rule])
workspace = await workspaces.open("acme", "prod", actor=SourceActor(name="monitoring"))
await workspace.publish(ServiceError(service="auth", severity=8), id="alert-7")
await Reactor(workspaces, actions={"triage": triage}).settle()
for envelope in await workspace.read():
print(
envelope.seq,
envelope.event_type,
envelope.actor.kind,
envelope.correlation_id,
envelope.depth,
)
1 ops:service.error source alert-7 0
2 reflexr:rule_fired system alert-7 1
3 reflexr:run_started system alert-7 1
4 ops:incident.opened agent alert-7 1
5 reflexr:run_succeeded system alert-7 1
The incident's causation names the firing and run behind it, and its depth counts the runs between it and the alert. A rule that fired on oncall:incident.opened would run at depth 2, and so on up to the workspace's max_depth, beyond which publishing is rejected, so workflows that trigger each other stop (Loop and spend safety).
Reading the log¶
| Method | Returns |
|---|---|
read(after_seq=0, before_seq=None, types=None, limit=None, last=None) |
The envelopes in a window of the log, in order, optionally of some types: the first limit or the last last (below) |
subscribe(after_seq=0) |
An async iterator of the envelopes after after_seq, then each new one as it commits |
head_seq() |
The latest seq, or 0 if the log is empty |
run(run_id) |
A run, or raises NotFound |
runs(rule=None, status=None, scope_key=None, limit=None) |
Runs, newest first, optionally of one rule, status or scope |
dead_letters(rule=None) |
The envelopes rules could not evaluate, oldest first |
rule_progress() |
Each rule's RuleProgress: its cursor, generation and pending absence deadlines |
get_rule(rule) |
One of the workspace's rules, a WorkspaceRule: the rule, its origin, and a stored rule's version and provenance, or raises NotFound |
rule_statuses() |
A RuleStatus for each of the workspace's rules, code rules first: its origin, code or stored, and a stored rule's version, whether it is enabled, its cursor, its lag behind the head, its generation and its number of dead_letters. A rule that has not evaluated the workspace yet is at cursor 0. |
schedule_ticks() |
When each schedule last ticked in the workspace |
schedule_statuses() |
A ScheduleStatus for each schedule that targets the workspace: its last_tick and next_tick, or neither before its first check |
The statuses are what REST's GET /v1/workspaces/{workspace_id}/rules and GET /v1/workspaces/{workspace_id}/schedules return (REST endpoints), and what MCP's rule_status and schedule_status tools report as text (External agents over MCP), so the surfaces agree. get_rule is what GET /v1/workspaces/{workspace_id}/rules/{rule} and MCP's get_rule return.
subscribe yields the stored envelopes and then the live ones on one iterator, so nothing falls between catching up and following along:
Windows, types and the tail¶
read takes a window of the log, the envelopes with after_seq < seq < before_seq, and returns them oldest first. Without before_seq the window runs to the head. types keeps the envelopes of those event types. limit takes the first so many that match, and last the last so many, which is the tail of the log. Storage does the filtering, so SQL storage reads only the envelopes it returns, over an index on the event type (Storage).
# The five latest deploys, oldest first
deploys = await workspace.read(types=["ops:deploy.finished"], last=5)
# The five before those: page backwards from the oldest seq you have
earlier = await workspace.read(types=["ops:deploy.finished"], last=5, before_seq=deploys[0].seq)
# What happened between two points in the log
window = await workspace.read(after_seq=40, before_seq=60)
Give limit or last, not both: a read with both, or with a negative number, is refused with ValidationFailed. Over REST the same parameters are query parameters of GET /v1/workspaces/{workspace_id}/events, with type repeated for several types (REST endpoints), and over MCP they are the arguments of read_events (External agents over MCP).
A Run records its rule, scope and scope_key, the seq it fired at and the envelopes it matched, its chain, its status and attempts, the last error and its reason code, its output, a graph's latest checkpoint, and the trace id of each attempt. Its id is the firing's id. The reactor explains how runs move through their statuses.
Operating runs and rules¶
People and operators steer runs and rules through the same handles. Each operation is one transaction that applies the core's transition and appends the resulting event, attributed to the handle's actor and in the run's chain, so the log shows who intervened:
| Method | Applies to | Does | Appends |
|---|---|---|---|
retry_run(run_id) |
Any run that has not succeeded and is not running | Makes it runnable now; a dead-lettered, cancelled or skipped run gets a fresh retry budget | reflexr:run_requeued |
skip_run(run_id, reason=None) |
Pending, retrying or dead-lettered runs | Gives up on it, unblocking later runs of its scope | reflexr:run_skipped |
cancel_run(run_id, reason=None) |
Pending, running, retrying or dead-lettered runs | Cancels it; a running attempt is stopped when its executor next renews its lease, and its result is discarded | reflexr:run_cancelled |
replay_rule(rule, from_seq=0, mode="rebuild") |
A registered rule | Resets the rule to evaluate the log again after from_seq (Replaying a rule) |
reflexr:rule_reset |
[stuck] = await ada.runs(rule="ops:error-spike", status="dead")
await ada.retry_run(stuck.id) # try again, with a fresh budget
await ada.skip_run(stuck.id, reason="billing is being migrated") # or give up on it
An operation that does not fit the run's status is rejected with InvalidState: retrying a run that succeeded, for example. checkpoint_run belongs to actions, which call it through Reaction.checkpoint (Graphs); give_feedback records judgements of runs, firings and chains (Feedback and evaluation).
Rejections¶
A command that cannot be carried out raises a Rejection and changes nothing. Each has a stable code, which the surfaces send to clients, a message, and payload() with its typed details:
from reflexr.core import NotFound
try:
await ada.publish(Deploy(service="auth", version="2.4.2"), correlation_id="nope")
except NotFound as rejection:
print(rejection.payload())
| Rejection | code |
When | HTTP |
|---|---|---|---|
NotFound |
not_found |
Something the command names does not exist, such as a run, a rule or a chain, or an event type is not accepted (entity, id) |
404 |
InvalidState |
invalid_state |
The command does not fit the current state, such as skipping a run that succeeded | 409 |
ValidationFailed |
validation_failed |
The data does not validate, such as feedback that does not fit its type, a replay beyond the head of the log, or a correlation_id naming a later event of a chain rather than its first (errors) |
422 |
Forbidden |
forbidden |
The actor may not do this, such as publishing one of reflexr's own events | 403 |
DepthExceeded |
depth_exceeded |
A run's event would extend its causal chain beyond max_depth (depth, limit) |
422 |
UnsupportedProtocol |
unsupported_protocol |
A client asked for another protocol version | 400 |
The HTTP column is how REST reports each one, with the payload as the response's detail, or as the rejection of a POST /commands result. Over the WebSocket the payload travels in the command's result (Stream protocol). Over MCP only the message does: a rejection is a tool error carrying it, such as Error executing tool retry_run: run nope does not exist, or a failed resource read (External agents over MCP). A rule that refers to something that does not exist is not a rejection but an InvalidRule error, raised when the application starts (Checking rules).