Architecture¶
Status: accepted design, being built in the phases tracked by RFC-0001. The code on
mainis still the pre-rebuild template until phase 0 removes it. This document is evergreen: it is updated in the same pull request as the code that changes it, and the table below shows what exists today. Decisions are recorded inadr/, proposals inrfcs/, and the wire protocol inprotocol.md.
| Package | Status |
|---|---|
reflexr.core | Planned (phase 1) |
reflexr.stream | Planned (phase 2) |
reflexr.agent | Planned (phase 3) |
reflexr.sql | Planned (phase 4) |
reflexr.fastapi, reflexr.mcp | Planned (phase 5) |
examples/oncall | Planned (phase 6) |
What reflexr is¶
reflexr is a Python library for reactive agent workflows. Applications publish events into streams. Rules watch the streams, and when a rule's condition holds (three errors from one service within a minute, a deploy followed by a spike, a heartbeat that stops), it runs a workflow: a pydantic-ai agent, a pydantic-graph graph, or a plain async function.
It is the sibling of artifactr, and the two cover different ways of working with LLM agents:
| artifactr | reflexr | |
|---|---|---|
| Trigger | A person talks to an agent | Something happens in the world |
| Shape | Live chat, with shared artifacts both sides edit | Rules over event streams, running workflows in the background |
| Unit | A workspace of artifacts and threads | A stream of events, with the rules that watch it |
| Agent's job | Collaborate, turn by turn | Decide and act, then report |
Neither library imports the other. They share conventions (envelopes, actors, tenancy, storage protocols, pydantic-ai integration, protocol shape), so that a third system can bring them together over a shared context with thin adapters (ADR-0003).
Goals¶
- Rules are typed, serializable data. Conditions compose, rules have a JSON Schema, and an agent can author one as easily as a person.
- Deciding is deterministic. The same log and the same rules always produce the same firings, so rules are testable without I/O and replayable after a fix.
- Rules are isolated. Each has its own cursor, state, retries and dead letters, so one failing rule never re-runs or blocks another.
- Acting is durable. Every firing's workflow runs at least once, with an idempotency key, and graph workflows resume from their last completed step after a crash.
- Multi-tenant, with many streams, scaled across processes with storage leases.
- Idiomatic use of Pydantic, pydantic-ai, pydantic-graph, SQLAlchemy, FastAPI and the MCP SDK, rather than parallel abstractions next to them.
- Small: seven core concepts, readable in an afternoon.
Non-goals (for now)¶
- A general stream processor. There are no joins across streams, and a stream's appends are serialized, so throughput per stream is modest; scale comes from many streams.
- Exactly-once external effects. Actions run at least once and receive an idempotency key.
- Connectors for every broker. Applications publish over the API or in process; Kafka or SQS consumers can be written against the same
publish. - A frontend.
Core concepts¶
| Concept | What it is |
|---|---|
| Event | A Pydantic model subclass for one kind of fact (service.error, deploy.finished), registered by name. Stored in an envelope with its seq, time, actor and causation. |
| Stream | A tenant-scoped, append-only log of events. The unit of ordering, isolation and scale. |
| Rule | Typed data: a condition over events (when), a scope that partitions its state and ordering (scope), and the action it runs (then). |
| Firing | The durable record that a rule's condition held for one scope at one point in the log, with the events that matched. Written atomically with the rule's cursor and state. |
| Action | What a firing runs: a pydantic-ai agent, a pydantic-graph graph or an async function, registered under a name that rules refer to. |
| Run | One execution of an action for one firing. Retried until it succeeds or is dead-lettered; its lifecycle is recorded in the log. |
| Reactor | The runtime that evaluates rules and executes runs for a set of streams, in any number of processes. |
Supporting types: Condition (the stages of a rule's when), Scope, Schedule, Actor, and Reaction (the dependencies an action receives: the stream, the firing, the run and the application's own deps).
Layers¶
graph TD
app["Your application"] --> fastapi["reflexr.fastapi<br/>ingest, REST, WebSocket"]
app --> mcp["reflexr.mcp<br/>MCP server"]
app --> agent["reflexr.agent<br/>agent and graph actions"]
fastapi --> stream
mcp --> stream
agent --> stream["reflexr.stream<br/>streams, storage protocol, Reactor"]
sql["reflexr.sql<br/>SQLAlchemy storage"] --> stream
stream --> core["reflexr.core<br/>pure, synchronous rules"] Dependencies point one way. Each layer is usable without the ones above it, and a test enforces the imports.
| Package | Depends on | Responsibility |
|---|---|---|
reflexr.core | pydantic | Events and envelopes, actors, conditions and their reducers, rules, evaluation, retry policy, schedules. Pure, synchronous, no I/O. |
reflexr.stream | core | Streams, Stream, the storage protocol, in-memory storage, the Reactor, function actions, the schedule runner. |
reflexr.agent | stream, pydantic-ai, pydantic-graph | Reaction, agent actions with the EventContext capability, graph actions with checkpoints. |
reflexr.sql (extra) | stream, SQLAlchemy 2 async, Alembic | Durable storage on PostgreSQL and SQLite, and its migrations. |
reflexr.fastapi (extra) | stream, FastAPI | HTTP ingest, REST reads and administration, and the WebSocket stream protocol. |
reflexr.mcp (extra) | stream, mcp | Publishing, reading and administration as MCP tools. |
reflexr.core: sans-IO¶
Core holds every decision as plain functions over immutable values (ADR-0001). The host (the stream layer, or any other) loads state, asks core to decide, and saves what core returns in its own transaction:
states = load_states(rule, scopes_of(batch)) # the host's I/O
evaluation = core.evaluate(rule, states, batch) # new states, firings, errors
save(evaluation, cursor=batch[-1].seq) # one transaction: states, cursor, firings
evaluatefolds envelopes into a rule's per-scope state and returns the new state, the firings, and any evaluation errors. It reads nothing but its arguments; time is the envelopes'ts.next_attemptapplies a retry policy: when a failed run should run again, or that it should be dead-lettered.due_ticksdecides which ticks of a schedule are due between two instants.
Because core is pure, its behaviour is pinned by a conformance suite of JSON fixtures: given a rule and a sequence of envelopes, expect these firings and this state, or this error. The fixtures are the language-neutral specification.
Events and envelopes¶
An event type is a Pydantic model that subclasses Event. Defining the class registers it under a name derived from the class name (ServiceError becomes service_error), or the name given explicitly, the same convention as artifactr's artifact types and pydantic-ai's CustomEvent:
class ServiceError(Event, name="service.error"):
service: str
severity: int
message: str
class Deploy(Event, name="deploy.finished"):
service: str
version: str
A stored event is an Envelope, which is also its wire shape:
class Envelope(BaseModel):
seq: int # gap-free position in the stream's log, from 1
id: str # the event's id: publishing the same id again is a no-op
ts: datetime # when it was appended; the clock rules use
stream_id: str
actor: Actor # who published it
causation: Causation | None # the firing and run that emitted it, and the chain's depth
correlation_id: str # the id of the first event in its causal chain
event: Event # discriminated by `type`
- Publishing is idempotent by event id within a stream, so producers can retry, and a run that retries does not duplicate the events it emits.
- Application events form an open family: any registered
Eventsubclass. reflexr's own facts form a closed union, so pyright checksmatchblocks for exhaustiveness:RuleFired,RuleErrored,RunStarted,RunProgressed,RunRetrying,RunSucceeded,RunDeadLettered,RunCancelled,RunSkipped, andTickfrom schedules. An event type from a newer version validates asUnknownEventand round-trips unchanged. - Actors match artifactr's kinds:
UserActor,AgentActor(an agent or graph run, by rule and run),ExternalAgentActor(an MCP client) andSystemActor(the reactor and schedules), plusSourceActorfor systems that publish events, such as a monitoring service. - Type allowlist.
Streams(storage, events=[ServiceError, Deploy, Heartbeat])rejects publishing any other type, even one registered elsewhere in the process.
Rules¶
A rule is a Pydantic model: what to watch, how to partition it, and what to run (ADR-0006).
from datetime import timedelta
from reflexr import F, Rule, on, run
error_spike = Rule(
name="error-spike",
when=on(ServiceError)
.where(F.severity >= 7)
.count(at_least=3, within=timedelta(minutes=1))
.throttle(at_most=1, per=timedelta(minutes=15)),
scope=F.service,
then=run(triage),
)
Serialized, the same rule is plain JSON, with a schema generated from the models:
{
"name": "error-spike",
"when": {
"filter": {"kind": "all", "of": [
{"kind": "on", "types": ["service.error"]},
{"kind": "where", "field": "severity", "op": "ge", "value": 7}
]},
"pattern": {"kind": "count", "at_least": 3, "within": "PT1M"},
"throttle": {"at_most": 1, "per": "PT15M"}
},
"scope": {"fields": ["service"]},
"then": {"action": "triage"}
}
Conditions¶
A condition is a pipeline of stages, evaluated per scope:
graph LR
envelope --> filter["filter<br/>(stateless)"]
filter --> dedupe["dedupe<br/>(optional)"]
dedupe --> pattern["pattern<br/>(stateful)"]
pattern --> throttle["throttle<br/>(optional)"]
throttle --> firing | Stage | Kinds | State |
|---|---|---|
| filter | on(types), where(field op value), all, any, not, and predicate(name), a registered Python function as an escape hatch | None |
| dedupe | Drop an event whose key was seen within a window | Recent keys |
| pattern | each (the default: every event that gets through fires), count(at_least, within), sequence(steps, within) (A then B), absence(within) (nothing matched for a while) | Per pattern |
| throttle | At most N firings per period, a cooldown that bounds spend | Recent firings |
Field references (F.severity, F.labels.env) are checked against the event types the filter admits when the rule is built, so a typo fails at startup rather than silently never matching. Operators are eq, ne, lt, le, gt, ge, in, contains, matches and exists.
Scopes¶
scope=F.service gives a rule independent state and ordering per service: three errors from auth and two from billing are two counts, and a slow run for auth never delays billing. The default scope is the whole stream. An envelope that passes a rule's filter but lacks its scope fields is an evaluation error for that rule.
Time¶
Rules measure time by the log: an envelope's ts, assigned when it is appended (ADR-0007). Evaluation never reads a clock, so replaying a log reproduces its firings exactly.
Every envelope advances a rule's clock, including envelopes its filter rejects: time is a separate input to the stateful stages, not an event they have to match. When the clock passes a deadline, such as the end of an absence window, it applies to every scope that already has state. So a Tick, which matches no heartbeat filter and has no service field, still lets absence fire for each service the rule has seen. A schedule appending a Tick every so often guarantees that time keeps moving in a quiet stream.
Deciding and acting¶
Processing a stream has two halves with different guarantees (ADR-0005):
sequenceDiagram
participant P as Producer
participant S as Stream log
participant E as Evaluator (per stream, leased)
participant C as core.evaluate
participant X as Executor (per run, leased)
participant A as Action
P->>S: publish(event) assigns seq
E->>S: read after the rule's cursor
E->>C: evaluate(rule, states, envelopes)
C-->>E: states', firings
E->>S: one transaction: states', cursor, RuleFired, pending runs
X->>S: claim a runnable run (lease)
X->>A: run(action, Reaction)
alt succeeds
X->>S: RunSucceeded
else raises
X->>S: RunRetrying with next_attempt_at, or RunDeadLettered
end
A-->>S: emitted events (causation = firing and run) Deciding is exact. One evaluator holds each stream's lease at a time and evaluates every rule in one pass over new envelopes. For each rule it loads the state of the scopes involved, calls core.evaluate, and saves the new state, the advanced cursor, the RuleFired envelopes and the pending runs in a single transaction. A crash before the commit means the same envelopes are evaluated again against the same state, so each envelope affects each rule exactly once. An evaluation error (a predicate that raises, a missing scope field) is recorded as RuleErrored and dead-lettered for that rule alone; the rule moves on.
Each rule has its own cursor. A rule added later, or reset for replay, catches up on its own without holding back the others.
Acting is at least once. Executors claim runnable runs under leases that they renew while working, so a crashed executor's run is picked up again when its lease lapses. With the default ordering="scope", runs of one rule and scope execute in firing order: a run that fails and is waiting to retry holds back later runs of the same scope, and nothing else. A run that exhausts its retry policy is dead-lettered; later runs of its scope continue (on_dead_letter="continue", the default) or wait for someone to retry or skip it ("block").
The firing id doubles as the run id and the action's idempotency key, and events a run emits get ids derived from it, so a retried run does not publish duplicates.
Actions¶
An action is what a firing runs (ADR-0008). Every kind receives the same Reaction: the stream (bound to the run's actor), the firing and its matched events, the run id and attempt, and the application's deps.
# A function
async def page(reaction: Reaction[AppDeps]) -> None:
await reaction.deps.pager.notify(reaction.firing.scope, reaction.run_id)
# A pydantic-ai agent: plain Agent, reflexr's capability adds the firing and event tools
triage_agent = Agent(
"anthropic:claude-sonnet-5-5",
deps_type=Reaction[AppDeps],
output_type=Triage,
capabilities=[EventContext(emit=[IncidentOpened], lookback=timedelta(hours=1))],
)
triage = AgentAction(triage_agent, name="triage")
# A pydantic-graph graph, checkpointed after every step
runbook = GraphAction(runbook_graph, name="runbook", state=RunbookState, inputs=from_triage)
- Agents are plain pydantic-ai
Agents withdeps_type=Reaction[...]. TheEventContextcapability renders the firing (the rule, the scope, the matched events) into the instructions, and gives the agent tools to read back through the stream and to emit events of the types it is allowed. The run's output can be emitted as an event, or handled by the action. - Graphs are pydantic-graph graphs built with
GraphBuilder. reflexr drives them step by step and saves the graph state and pending tasks to the run after every step. A retry, or another executor after a crash, resumes from the last completed step instead of starting over (ADR-0009). - Functions are
async defover aReaction.
Rules refer to actions by name ({"action": "triage"}), because functions and agents are not data. run(triage) names the action and lets the reactor register it; rules loaded from JSON resolve names against the actions the application registers.
Safety¶
LLM workflows triggered by events can loop and can spend (ADR-0010):
- Causation depth. An event emitted by a run carries the depth of its causal chain. Publishing beyond the stream's limit (8 by default) is rejected, so a rule whose action triggers itself stops instead of running away.
- Throttles on rules cap how often a rule can fire per scope.
- Emit allowlists on the
EventContextcapability limit which event types an agent can publish. - Concurrency limits per reactor, per rule and per stream bound how many runs execute at once. Agent actions accept pydantic-ai
UsageLimits. - Tenancy: every handle is scoped to one tenant and stream.
Tenancy and concurrency¶
- Scoped handles.
await streams.open(tenant_id, stream_id, actor=...)returns aStreambound to that tenant, stream and actor. Nothing below it accepts a raw tenant id (ADR-0004). - Sequencing.
seqis assigned inside the append transaction, which holds the stream's lock (in SQL, the stream row, lockedFOR UPDATE), giving a gap-free total order per stream.tsis assigned under the same lock as the later of the clock and the previous envelope'sts, so it never decreases within a stream, whatever the clock skew between processes. Appends to one stream are serialized; many streams scale out. - Leases. One evaluator per stream, one executor per run, and one schedule runner per schedule, each a storage lease with a time-to-live that its holder renews and that lapses if it dies. Any number of reactor processes can share the work.
Schedules¶
A schedule appends events on a timetable: Schedule(name="heartbeat-check", every=timedelta(seconds=30), streams=[...]), or a cron expression with a time zone. Each tick's event id is derived from the schedule and the tick's time, so replicas that race publish it once, and a lease keeps them from racing in the first place. Missed ticks are skipped or caught up, per schedule.
Surfaces¶
Every surface is a thin adapter over a Stream handle (ADR-0011). The wire formats are in protocol.md.
| Surface | Package | What it offers |
|---|---|---|
| HTTP ingest and REST | reflexr.fastapi | Publish events (one or a batch, idempotent). Read the log. Inspect and administer rules (cursor, lag, replay), runs (retry, skip, cancel) and dead letters. |
| WebSocket | reflexr.fastapi | A resumable subscription to the log with the same hello, replay and close-code shape as artifactr's thread protocol, plus publish and administration commands. |
| Schedules | reflexr.stream | Cron and interval schedules that publish into streams. |
| MCP | reflexr.mcp | Publishing, reading and administration as MCP tools, so external agents can feed and operate streams. |
Authentication is the host's: each surface takes a resolver that returns the tenant and actor for a request.
Storage protocol¶
reflexr.stream defines the storage protocol, and the same behaviour suite runs against every implementation:
InMemoryStorage, for tests, examples and single-process prototypes.SqlStorage(reflexr.sql), on PostgreSQL and SQLite with SQLAlchemy 2's asyncio extension, with Alembic migrations shipped in the package.
A transaction is scoped to one stream. Within it the host appends envelopes, loads and saves rule cursors and states, creates and updates runs, and dead-letters evaluation errors. Outside transactions, storage serves reads (read, subscribe, runs, dead letters), leases with an injectable clock, and schedule state.
Data model¶
erDiagram
TENANT ||--o{ STREAM : owns
STREAM ||--o{ EVENT : "log ordered by seq"
STREAM ||--o{ RULE_CURSOR : "one per rule"
RULE_CURSOR ||--o{ RULE_STATE : "one per scope"
STREAM ||--o{ RUN : "one per firing"
RUN }o--|| EVENT : "fired at"
STREAM ||--o{ DEAD_LETTER : "per rule" - Rules live in code (v0.1) and are identified by name. A rule's cursor row stores a hash of its definition; changing the definition resets its state and is recorded in the log.
- Rule state is a JSON document per rule and scope, owned by the rule's condition, saved with its cursor.
- Runs hold their status, attempts, next attempt time, lease, and for graphs the latest checkpoint.
Aligned with artifactr¶
| Convention | artifactr | reflexr |
|---|---|---|
| Tenancy | workspaces.open(tenant, workspace, actor=) | streams.open(tenant, stream, actor=) |
| Log | Per-workspace envelopes with a gap-free seq | Per-stream envelopes with a gap-free seq |
| Actors | user, agent, external agent, system | The same, plus source |
| Types | Artifact subclasses registered by name | Event subclasses registered by name |
| Core | Sans-IO commit, conformance fixtures | Sans-IO evaluate, conformance fixtures |
| Storage | Protocol, in-memory and SQL, leases, one behaviour suite | The same |
| pydantic-ai | A capability plus a deps type (ArtifactWorkspace, Session) | A capability plus a deps type (EventContext, Reaction) |
| WebSocket | hello, replay, replay_complete, close codes | The same shape |
| MCP | Tenant in resource URIs | The same |
Dependencies¶
Python 3.12+. Runtime: pydantic (core); pydantic-ai-slim and pydantic-graph (agent). Extras: sql (sqlalchemy[asyncio], alembic), postgres and sqlite (drivers), fastapi, mcp. Tooling: uv, ruff, pyright in strict mode, pytest, and Zensical with mkdocstrings for the documentation site.
Testing¶
- Core: the conformance fixtures, plus property tests that replaying any log reproduces its firings.
- Stream: one behaviour suite (publishing, evaluation, runs, retries, leases, schedules) against every storage.
- Agent: scripted models with pydantic-ai's
FunctionModelandTestModel; graph runs interrupted and resumed. - Surfaces: contract tests for every frame and endpoint, and an MCP client round trip.
- Reference implementation: the real server and client end to end, with a scripted model.
Coverage is 100% of lines and branches, and pyright runs in strict mode with no suppressions (ADR-0013).
Decisions¶
| ADR | Decision |
|---|---|
| 0001 | A library with a sans-IO core, replacing the template |
| 0002 | The name reflexr |
| 0003 | An independent sibling of artifactr, with aligned conventions |
| 0004 | Tenant-scoped streams with one log each |
| 0005 | Per-rule cursors: decide exactly, act at least once |
| 0006 | Rules as typed, serializable data |
| 0007 | Rule state as pure reducers, with the log as the clock |
| 0008 | Actions: agents, graphs and functions over one Reaction |
| 0009 | Graph checkpoints at step boundaries |
| 0010 | Loop and spend safety |
| 0011 | Surfaces: ingest, REST, WebSocket, schedules and MCP |
| 0012 | Trunk-based development with RFCs, ADRs and evergreen docs |
| 0013 | Quality gates |
| 0014 | MIT license |
| 0015 | Reference implementation: incident response |
Open questions¶
- Managing rules at runtime. Rules are data from day one, but v0.1 registers them in code. Storing versioned rules per tenant, and installing them through the API, is the next step, and what lets an agent's proposed rule go live once a person accepts it.
- Pausing runs for a person. A workflow that needs approval mid-way could pause as pydantic-ai deferred tools do. Until then, a run can emit an event and a second rule can continue when the answer arrives.
- Live output from runs. Token-level streaming of agent runs over the WebSocket, as artifactr streams its runs.
- Resuming inside parallel branches. Graph checkpoints are proven for sequential steps; a run that crashes inside a fork restarts from the last checkpoint before it.
- Retention. Compacting old envelopes, and what a rule replaying past the retained log starts from.
- Hot streams. A stream's appends are serialized. Partitioning one logical stream across several, by key, may be needed for high-volume sources.