Skip to content

Architecture

Status: accepted design, being built in the phases tracked by RFC-0001. The code on main is 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 in adr/, proposals in rfcs/, and the wire protocol in protocol.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
  • evaluate folds 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_attempt applies a retry policy: when a failed run should run again, or that it should be dead-lettered.
  • due_ticks decides 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 Event subclass. reflexr's own facts form a closed union, so pyright checks match blocks for exhaustiveness: RuleFired, RuleErrored, RunStarted, RunProgressed, RunRetrying, RunSucceeded, RunDeadLettered, RunCancelled, RunSkipped, and Tick from schedules. An event type from a newer version validates as UnknownEvent and round-trips unchanged.
  • Actors match artifactr's kinds: UserActor, AgentActor (an agent or graph run, by rule and run), ExternalAgentActor (an MCP client) and SystemActor (the reactor and schedules), plus SourceActor for 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 with deps_type=Reaction[...]. The EventContext capability 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 def over a Reaction.

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 EventContext capability 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 a Stream bound to that tenant, stream and actor. Nothing below it accepts a raw tenant id (ADR-0004).
  • Sequencing. seq is assigned inside the append transaction, which holds the stream's lock (in SQL, the stream row, locked FOR UPDATE), giving a gap-free total order per stream. ts is assigned under the same lock as the later of the clock and the previous envelope's ts, 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 FunctionModel and TestModel; 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.