Skip to content

reflexr.core

reflexr's pure core: every decision as plain functions over immutable values.

Nothing in this package performs I/O or reads a clock. The workspace layer (or any host) loads state, asks core to decide, and saves what core returns in one transaction.

Events

Event types, their namespaces and registries, and the envelope each stored event travels in. See Events and envelopes.

Event pydantic-model

Bases: BaseModel

Base class for event types.

Subclass it through an abstract base that declares a namespace with event_namespace=, and give it ordinary Pydantic fields. Events are facts, so they are immutable, and unknown fields are rejected so that a producer's typo fails loudly.

Config:

  • frozen: True
  • extra: forbid

event_type class-attribute

event_type: str

The qualified name: the namespace, a :, and the local name, derived from the class name unless given as name=.

event_namespace class-attribute

event_namespace: str | None = None

The namespace, from the base that declared it.

UnknownEvent pydantic-model

Bases: Event

An event whose type is not registered in this process, kept as-is so it round-trips.

Stored envelopes can outlive the code that defined their types, and a newer producer can use types this version does not know. Their data is kept in model_extra.

Config:

  • frozen: True
  • extra: allow

Fields:

unknown_type pydantic-field

unknown_type: str

The type name the event was stored with.

AnyEvent module-attribute

AnyEvent = Annotated[
    Event,
    PlainValidator(_validate_event),
    PlainSerializer(_serialize_event),
    WithJsonSchema(
        {
            "type": "object",
            "properties": {"type": {"type": "string"}},
            "required": ["type"],
            "additionalProperties": True,
        }
    ),
]

Any event, validated through the registry by its "type".

Envelope pydantic-model

Bases: BaseModel

A stored event with its position and attribution. Its shape is also the wire shape.

Fields:

seq pydantic-field

seq: int

The event's position in its workspace's log, gap-free from 1.

ts pydantic-field

ts: AwareDatetime

When the event was appended. Rules measure time with it; it never decreases.

correlation_id pydantic-field

correlation_id: str

The id of the first event in this event's causal chain.

traceparent pydantic-field

traceparent: str | None = None

The W3C trace context of the span that published the event, so the runs it causes can link back to it.

depth property

depth: int

How many runs lie between this event and an event nothing caused.

event_type property

event_type: str

The type name of the event.

data cached property

data: dict[str, Any]

The event's fields as JSON-compatible values, which conditions read.

Causation pydantic-model

Bases: BaseModel

What caused an event that a run emitted, and how deep its causal chain is.

Config:

  • frozen: True

Fields:

depth pydantic-field

depth: int

How many runs lie between this event and an event nothing caused.

EventRegistry

EventRegistry()

Bases: Mapping[str, type['Event']]

A set of namespaces, read as the event types in them, by name.

A registry scopes which types a Workspaces accepts, so applications in one process stay apart. It is a view of the process's one table of types: a base's registry= puts its namespace in a registry, and add includes another's. reflexr's own namespace is in every registry.

add

add(*bases: type[Event]) -> None

Include the namespace each base declared, and so every type in it.

Raises:

Type Description
TypeError

If a base declares no namespace.

DEFAULT_REGISTRY module-attribute

DEFAULT_REGISTRY = EventRegistry()

The registry of every namespace whose base names no other one.

EventName module-attribute

EventName = Annotated[
    str,
    AfterValidator(check_event_name),
    WithJsonSchema(
        {"type": "string", "pattern": _EVENT_NAME}
    ),
]

A qualified event type name, checked with a hint when it has no namespace.

check_event_name

check_event_name(name: str) -> str

Return name if it is a qualified event type name.

Raises:

Type Description
ValueError

If it isn't. A name with no namespace is told the qualified names of the types with that local name, in whatever registry.

load_event

load_event(data: Mapping[str, Any]) -> Event

Validate event data, including its "type", as the registered type it names.

Data of an unregistered type becomes an UnknownEvent.

Raises:

Type Description
ValidationFailed

If there is no "type", or the data does not match its type.

dump_event

dump_event(
    event: Event,
    *,
    mode: Literal["json", "python"] = "json",
) -> dict[str, Any]

Return an event's data with its "type" first, as stored and sent on the wire.

type_of

type_of(event: Event) -> str

Return the type name of an event, including one kept as an UnknownEvent.

reflexr's own events

The facts reflexr appends about rules, runs, feedback and schedules. See reflexr's own events.

SystemEvent module-attribute

reflexr's facts about rules and runs. Each names the rule it is about.

SYSTEM_EVENTS module-attribute

Every event type reflexr itself appends.

about_rule

about_rule(event: Event) -> RuleName | None

Return the rule a reflexr event is about, or None for any other event.

RuleFired pydantic-model

Bases: _ReflexrEvent

A rule's condition held for one scope; a run of its action was created.

Fields:

matched pydantic-field

matched: tuple[int, ...]

The seq of every envelope that made the condition hold.

RuleErrored pydantic-model

Bases: _ReflexrEvent

A rule could not evaluate one envelope; it is dead-lettered for that rule alone.

Fields:

RuleReset pydantic-model

Bases: _ReflexrEvent

A rule's state was reset, because its definition changed or it was replayed.

Fields:

silent_through pydantic-field

silent_through: int = 0

Envelopes up to this seq rebuild the rule's state without recording firings.

RuleInstalled pydantic-model

Bases: _ReflexrEvent

A stored rule was installed or updated: its new version, whole, and where it came from.

The facts are a stored rule's history: the version that fired a run is the latest one before its firing.

Fields:

spec pydantic-field

spec: dict[str, JsonValue]

The rule, as JSON.

provenance pydantic-field

provenance: dict[str, JsonValue]

Where the change came from, as its installer said.

RuleArchived pydantic-model

Bases: _ReflexrEvent

A stored rule was archived, and its unfinished runs cancelled.

Fields:

cancelled pydantic-field

cancelled: int

How many unfinished runs archiving cancelled.

RunStarted pydantic-model

Bases: _ReflexrEvent

An attempt of a run began.

Fields:

scope pydantic-field

scope: dict[str, JsonValue] = {}

The values of the rule's scope fields, as on reflexr:rule_fired.

RunProgressed pydantic-model

Bases: _ReflexrEvent

A graph run completed a step and saved a checkpoint.

Fields:

RunRetrying pydantic-model

Bases: _ReflexrEvent

An attempt failed, and the run will be tried again.

Fields:

reason pydantic-field

reason: str | None = None

A stable code for why the attempt failed, when the action gave one.

RunSucceeded pydantic-model

Bases: _ReflexrEvent

A run finished.

Fields:

RunDeadLettered pydantic-model

Bases: _ReflexrEvent

A run exhausted its retries, or failed in a way retrying cannot fix.

Fields:

reason pydantic-field

reason: str | None = None

A stable code for why the last attempt failed, when the action gave one.

RunCancelled pydantic-model

Bases: _ReflexrEvent

A run was cancelled.

Fields:

RunRequeued pydantic-model

Bases: _ReflexrEvent

Someone made a run runnable again. The envelope's actor records who.

Fields:

RunSkipped pydantic-model

Bases: _ReflexrEvent

A run was skipped without completing, unblocking its scope.

Fields:

FeedbackGiven pydantic-model

Bases: _ReflexrEvent

A person or an evaluator judged a run, a firing or a causal chain.

Fields:

value pydantic-field

value: dict[str, JsonValue]

The feedback's fields, as validated against its registered type.

Tick pydantic-model

Bases: _ReflexrEvent

A schedule's tick. Ticks also move rules' clocks forward in quiet workspaces.

Fields:

  • schedule (str)
  • at (AwareDatetime)

Actors

Who published an event or performed a command. See Actors.

Actor module-attribute

Any actor, discriminated by kind.

UserActor pydantic-model

Bases: BaseModel

A person, identified by the host application's user id.

Config:

  • frozen: True

Fields:

participant property

participant: str

A key that is equal for every action by the same person.

display_name property

display_name: str

How this actor is named to people and agents.

AgentActor pydantic-model

Bases: BaseModel

An agent or graph running an action for one firing of a rule.

Config:

  • frozen: True

Fields:

participant property

participant: str

A key that is equal for every run of the same rule's action.

display_name property

display_name: str

How this actor is named to people and agents.

ExternalAgentActor pydantic-model

Bases: BaseModel

An agent outside reflexr, connected over MCP.

Config:

  • frozen: True

Fields:

  • kind (Literal['external_agent'])
  • client_id (str)
  • name (str | None)

participant property

participant: str

A key that is equal for every action by the same client.

display_name property

display_name: str

How this actor is named to people and agents.

SystemActor pydantic-model

Bases: BaseModel

reflexr itself (the reactor, schedules), or the application acting on its own behalf.

Config:

  • frozen: True

Fields:

participant property

participant: str

A key that is equal for every action by the same system component.

display_name property

display_name: str

How this actor is named to people and agents.

SourceActor pydantic-model

Bases: BaseModel

A system that publishes events, such as a monitoring service or a webhook.

Config:

  • frozen: True

Fields:

participant property

participant: str

A key that is equal for every event from the same source.

display_name property

display_name: str

How this actor is named to people and agents.

EvaluatorActor pydantic-model

Bases: BaseModel

An evaluator: a judge or decision model whose verdicts are recorded as feedback.

Config:

  • frozen: True

Fields:

version pydantic-field

version: str

The evaluator's version, such as a hash of a trained judge, so verdicts never mix.

participant property

participant: str

A key that is equal for every verdict of the same evaluator version.

display_name property

display_name: str

How this actor is named to people and agents.

Rules

A rule's condition, scope, action and policies. See Rules.

Rule pydantic-model

Bases: BaseModel

When when holds for a scope, run then.

Config:

  • frozen: True
  • extra: forbid

Fields:

Validators:

  • _unreserved → name

name pydantic-field

name: RuleName

A namespace, a : and a name, such as oncall:triage. reflexr is reserved.

ordering pydantic-field

ordering: Literal['scope', 'none'] = 'scope'

scope: runs of one scope execute in firing order. none: in parallel.

on_dead_letter pydantic-field

on_dead_letter: Literal['continue', 'block'] = 'continue'

Whether later runs of a scope continue past a dead-lettered run, or wait for an operator.

start pydantic-field

start: Literal['now', 'beginning'] = 'now'

Where a rule starts in a workspace whose log already has envelopes.

timeout pydantic-field

timeout: timedelta | None = None

How long one attempt of the action may take before it is cancelled and counts as failed.

enabled pydantic-field

enabled: bool = True

Whether the reactor evaluates the rule and executes its runs.

A disabled rule stays registered and checked, but its cursor holds, it records no firings, and its pending and retrying runs wait. Enabling it again resumes from its cursor.

definition

definition() -> str

Return a hash of what the rule decides (its condition and scope).

When it changes, the rule's state no longer applies and is reset. Changes to the action or its params, retries, ordering or whether it is enabled do not reset it.

check

check(
    *,
    events: Mapping[str, type[Event]],
    actions: Collection[str] | None = None,
    predicates: Collection[str] = (),
) -> None

Confirm that everything the rule refers to exists.

Parameters:

Name Type Description Default
events Mapping[str, type[Event]]

The event types the rule may watch, by name. reflexr's own events are always available.

required
actions Collection[str] | None

The names of the registered actions, or None to leave the action for the runtime to check.

None
predicates Collection[str]

The names of the registered predicates.

()

Raises:

Type Description
InvalidRule

Listing every problem found.

Scope pydantic-model

Bases: BaseModel

The event fields that partition a rule's state and ordering. None: the whole workspace.

Config:

  • frozen: True
  • extra: forbid

Fields:

of

of(
    data: Mapping[str, Any],
) -> tuple[ScopeKey, dict[str, JsonValue]] | None

Return the scope key and values of an event's data, or None if a field is missing.

by

by(*fields: FieldRef | str) -> Scope

Scope a rule by event fields: by(F.service), by(F.service, F.region).

scope_key

scope_key(values: list[JsonValue]) -> ScopeKey

Return the canonical key of scope values: their JSON, with sorted object keys.

ActionRef pydantic-model

Bases: BaseModel

The action a rule runs, by name, and the params the rule passes it.

Config:

  • frozen: True
  • extra: forbid

Fields:

params pydantic-field

params: dict[str, JsonValue]

Values for the action's params model, validated against it when the reactor is built and at each attempt.

run

run(
    action: str | Named, /, **params: JsonValue
) -> ActionRef

Refer to an action by its name, or by the action itself, with its params.

run(triage), or run("notify", thread_id="thr_4") for an action that declares a params model.

load_params

load_params(
    rule: Rule, model: type[BaseModel] | None
) -> BaseModel | None

Validate a rule's params as its action's params model.

The params are JSON, so they are validated as JSON: a strict model accepts what a rule written as JSON can say, such as a date as a string.

Parameters:

Name Type Description Default
rule Rule

The rule.

required
model type[BaseModel] | None

The params model the rule's action declares, or None if it declares none.

required

Returns:

Type Description
BaseModel | None

The validated params, or None for an action that declares no model.

Raises:

Type Description
InvalidRule

Listing every field that does not validate, by its path in the rule's JSON, such as then.params.thread_id, or that then.params must be empty when the action declares no model.

Named

Bases: Protocol

Anything with a name, such as an action.

name property

name: str

The name a rule refers to it by.

RetryPolicy pydantic-model

Bases: BaseModel

How a failed run is retried: exponential backoff, up to max_attempts attempts.

Config:

  • frozen: True
  • extra: forbid

Fields:

backoff pydantic-field

backoff: timedelta = timedelta(seconds=1)

The delay after the first failure.

delay

delay(attempts: int) -> timedelta

Return the delay after the attempts-th failed attempt.

next_attempt

next_attempt(
    attempts: int, now: AwareDatetime
) -> datetime | None

Return when to try again after attempts failed attempts, or None to give up.

Stored rules

What rules installed at runtime may do, the limits every one is held to, and how a workspace keeps one. See RFC-0003.

StoredRules dataclass

StoredRules(
    allow: Callable[
        [TenantId, WorkspaceId, Actor, RuleChange],
        Awaitable[bool],
    ],
    actions: Mapping[str, type[BaseModel] | None],
    namespaces: Set[str],
)

What stored rules may do: who may change them, the actions they run, their namespaces.

Everything else is fixed by core's limits, the same for every tenant.

Parameters:

Name Type Description Default
allow Callable[[TenantId, WorkspaceId, Actor, RuleChange], Awaitable[bool]]

Decides whether an actor may make a change to a workspace's stored rules. Shaped like Authorize, with the change, and asked after it.

required
actions Mapping[str, type[BaseModel] | None]

The only actions stored rules may run, by name, each mapped to the params model it declares, or to None if it declares none.

required
namespaces Set[str]

The rule namespaces stored rules may be named in, such as {"chat"}.

required

Raises:

Type Description
ValueError

If namespaces includes reflexr, which is reserved.

RuleChange pydantic-model

Bases: BaseModel

A change to a stored rule, which the configuration's allow hook allows or refuses.

Config:

  • frozen: True
  • extra: forbid

Fields:

kind pydantic-field

kind: Literal['install', 'update', 'archive']

Which command makes the change: install_rule, update_rule or archive_rule.

rule pydantic-field

rule: RuleName

The stored rule's name.

spec pydantic-field

spec: Rule | None = None

The rule as installed or updated. None when archiving.

check_stored

check_stored(rule: Rule, config: StoredRules) -> list[str]

Return every problem with a rule as a stored rule, under config and core's limits.

It is pure, beside Rule.check, so a draft can be checked with both before it is proposed. The problems follow the rule's fields, each once, and an empty list means the rule may be stored.

check_provenance

check_provenance(
    provenance: Mapping[str, JsonValue],
) -> list[str]

Return the problem with a change's provenance, if its JSON is over core's limit.

Provenance is opaque to reflexr, so its size is all there is to check. It is measured as a rule's JSON is: compact, in UTF-8 bytes.

StoredRule pydantic-model

Bases: BaseModel

One stored rule as its workspace keeps it: the current version, and where it came from.

Storage keeps only the current version, and the log's facts hold the history. Every change is a new version, numbered from 1 per workspace and name, and archiving is one too.

Config:

  • frozen: True
  • extra: forbid

Fields:

rule pydantic-field

rule: Rule

The rule, whose name is the stored rule's.

version pydantic-field

version: int

The current version, from 1.

status pydantic-field

status: Literal['active', 'archived'] = 'active'

Whether the rule runs. An archived rule stays stored, so one installed again under its name continues its versions.

provenance pydantic-field

provenance: dict[str, JsonValue]

Where the change came from, as its installer says: opaque JSON that reflexr does not read.

MAX_STORED_FIRINGS_PER_HOUR module-attribute

MAX_STORED_FIRINGS_PER_HOUR = 60

The most firings a stored rule's throttle may allow in an hour.

MAX_STORED_ATTEMPTS module-attribute

MAX_STORED_ATTEMPTS = 5

The most attempts a stored rule's retry policy may make of a run.

MAX_STORED_TIMEOUT module-attribute

MAX_STORED_TIMEOUT = timedelta(minutes=5)

The longest timeout a stored rule may give its action.

MAX_STORED_WINDOW module-attribute

MAX_STORED_WINDOW = timedelta(days=1)

The longest window a stored rule may use: a pattern's, a dedupe's, or a throttle's period.

MAX_STORED_DESCRIPTION module-attribute

MAX_STORED_DESCRIPTION = 2000

The most characters a stored rule's description may have.

MAX_STORED_BYTES module-attribute

MAX_STORED_BYTES = 16 * 1024

The most bytes a stored rule's JSON may have.

MAX_STORED_PROVENANCE module-attribute

MAX_STORED_PROVENANCE = 4 * 1024

The most bytes a change's provenance may have, as JSON.

MAX_STORED_RULES module-attribute

MAX_STORED_RULES = 50

The most active stored rules a workspace may have.

Conditions

The builder, and the stages it produces. See Conditions.

on

on(*types: type[Event] | str) -> Condition

Start a condition on envelopes of the given event types, by class or by name.

sequence

sequence(*steps: Condition, within: timedelta) -> Condition

Fire when envelopes match steps in order within within.

Each step is a filter-only condition, such as on(Deploy) or on(ServiceError).where(F.severity >= 7).

F module-attribute

F = _FieldRoot()

The root of field references: F.severity >= 7.

field

field(path: str) -> FieldRef

Return a reference to the field at a dotted path, for names F cannot spell.

Examples are fields named like FieldRef methods (field("contains")).

FieldRef

FieldRef(path: tuple[str, ...])

A reference to an event field, built from F: F.severity, F.labels.env.

Comparisons build WhereFilter s. Equality is spelled eq, because == has to return a bool; where(service="auth") is the shorthand.

eq

eq(value: JsonValue) -> WhereFilter

The field equals value.

ne

ne(value: JsonValue) -> WhereFilter

The field exists and does not equal value.

is_in

is_in(values: list[JsonValue]) -> WhereFilter

The field equals one of values.

contains

contains(value: JsonValue) -> WhereFilter

The field is a list containing value, or a string containing it.

matches

matches(pattern: str) -> WhereFilter

The field is a string that the regular expression pattern searches.

exists

exists(exists: bool = True) -> WhereFilter

The field is present (or, with False, absent).

Condition pydantic-model

Bases: _Stage

A rule's when: a filter, and optionally dedupe, a pattern and a throttle.

Fields:

where

where(*filters: Filter, **equals: JsonValue) -> Condition

Narrow the filter: every filter given, and every field equal to its keyword.

distinct

distinct(
    *key: FieldRef | str, within: timedelta
) -> Condition

Drop an envelope whose key fields were already seen within within.

count

count(*, at_least: int, within: timedelta) -> Condition

Fire when at_least envelopes pass within within.

absent

absent(*, within: timedelta) -> Condition

Fire when nothing has passed for within since the last envelope that did.

at_most

at_most(times: int, *, per: timedelta) -> Condition

Allow at most times firings per scope in any period of per.

Filter

Filter = Annotated[
    OnFilter
    | WhereFilter
    | AllFilter
    | AnyFilter
    | NotFilter
    | PredicateFilter,
    Field(discriminator=kind),
]

Which envelopes a rule considers.

OnFilter pydantic-model

Bases: _Stage

Envelopes whose event is one of types.

Fields:

WhereFilter pydantic-model

Bases: _Stage

Envelopes whose event field field compares with value by op.

A missing field matches nothing, except exists with value=False.

Fields:

Validators:

  • _check_value

field pydantic-field

field: str

A dotted path into the event's fields, such as labels.env.

Op module-attribute

Op = Literal[
    "eq",
    "ne",
    "lt",
    "le",
    "gt",
    "ge",
    "in",
    "contains",
    "matches",
    "exists",
]

How a WhereFilter compares a field with its value.

AllFilter pydantic-model

Bases: _Stage

Envelopes that match every filter in of.

Fields:

AnyFilter pydantic-model

Bases: _Stage

Envelopes that match at least one filter in of.

Fields:

NotFilter pydantic-model

Bases: _Stage

Envelopes that do not match filter.

Fields:

PredicateFilter pydantic-model

Bases: _Stage

Envelopes whose event the registered Python predicate name accepts.

An escape hatch for what the other filters cannot express. Predicates must be pure.

Fields:

Dedupe pydantic-model

Bases: _Stage

Drop an envelope whose key fields were already seen within within.

Fields:

Pattern

Pattern = Annotated[
    EachPattern
    | CountPattern
    | SequencePattern
    | AbsencePattern,
    Field(discriminator=kind),
]

When the envelopes a rule considers make it fire.

EachPattern pydantic-model

Bases: _Stage

Fire for every envelope that passes the filter.

Fields:

CountPattern pydantic-model

Bases: _Stage

Fire when at_least envelopes pass within within. A firing consumes them.

Fields:

SequencePattern pydantic-model

Bases: _Stage

Fire when envelopes match steps in order, all within within of the first.

An envelope that matches the first step while a sequence is in progress (but not the step it is waiting for) starts the sequence again from it, so the latest start counts.

Fields:

AbsencePattern pydantic-model

Bases: _Stage

Fire once when nothing has passed the filter for within since the last envelope did.

Only scopes that have been seen can go quiet. Time passes with every envelope in the log, whether or not it passes the filter.

Fields:

Throttle pydantic-model

Bases: _Stage

Allow at most at_most firings per scope in any period of per.

Fields:

Evaluation

The host contract: what to load, and how a batch of envelopes is decided (ADR-0017). The reactor uses these; call them directly only when you write a host of your own, or to test a rule without storage.

needs

needs(
    rule: Rule,
    progress: RuleProgress,
    envelopes: Sequence[Envelope],
    *,
    predicates: Predicates = NO_PREDICATES,
) -> frozenset[ScopeKey]

Return the scopes whose state evaluating envelopes reads.

These are the scopes of envelopes that pass the rule's filter, and the scopes whose absence deadline passes during the batch. Scopes the host has no state for are new.

evaluate

evaluate(
    rule: Rule,
    progress: RuleProgress,
    states: Mapping[ScopeKey, ScopeState],
    envelopes: Sequence[Envelope],
    *,
    predicates: Predicates = NO_PREDICATES,
    max_depth: int | None = None,
) -> Evaluation

Fold envelopes into a rule's state and decide its firings.

Parameters:

Name Type Description Default
rule Rule

The rule.

required
progress RuleProgress

Its progress, whose definition must be the rule's.

required
states Mapping[ScopeKey, ScopeState]

The state of the scopes needs returned that exist.

required
envelopes Sequence[Envelope]

New envelopes, in seq order, after progress.cursor.

required
predicates Predicates

The registered predicates.

NO_PREDICATES
max_depth int | None

The deepest causal chain a firing may extend. A firing whose facts would be deeper is refused and recorded as an evaluation error, so rules that trigger each other stop. None sets no limit.

None

Raises:

Type Description
ValueError

If the progress belongs to another definition, or the envelopes are out of order or already evaluated.

NotLoaded

If a scope whose deadline passes was not loaded.

begin

begin(rule: Rule, *, head_seq: int) -> RuleProgress

Return where a rule starts in a workspace whose log is at head_seq.

reset

reset(
    rule: Rule,
    progress: RuleProgress,
    *,
    from_seq: int,
    replayed: bool = False,
    silent_through: int = 0,
) -> tuple[RuleProgress, RuleReset]

Reset a rule's state to evaluate again from from_seq onwards.

The host deletes the rule's scope states along with saving the new progress. Resetting is how a changed definition takes effect, and how a rule is replayed.

Parameters:

Name Type Description Default
rule Rule

The rule, as it is now.

required
progress RuleProgress

Its current progress.

required
from_seq int

Evaluate again from the envelope after this seq.

required
replayed bool

Whether an operator asked for the reset, rather than a changed definition.

False
silent_through int

Rebuild state without recording firings or errors up to this seq, typically the head of the log, so a replay recomputes state without acting on the past again.

0

matches

matches(
    filter: Filter,
    envelope: Envelope,
    predicates: Predicates,
) -> bool

Whether an envelope passes a filter. Predicates may raise.

Predicates

Predicates = Mapping[str, Callable[[Event], bool]]

Registered Python predicates, by name. They must be pure.

NO_PREDICATES module-attribute

NO_PREDICATES: Predicates = MappingProxyType({})

NotLoaded

Bases: LookupError

The host did not load a scope that needs asked for.

State

What a rule remembers between batches, and what evaluating a batch decides.

RuleProgress pydantic-model

Bases: _State

Where a rule is in a workspace's log, stored with its cursor.

Fields:

cursor pydantic-field

cursor: int = 0

The seq of the last envelope evaluated.

generation pydantic-field

generation: int = 0

Incremented when the rule's state is reset, so firing ids never repeat.

definition pydantic-field

definition: str = ''

The definition the state belongs to.

deadlines pydantic-field

deadlines: dict[ScopeKey, AwareDatetime] = {}

When each waiting scope's absence window ends.

silent_through pydantic-field

silent_through: int = 0

Envelopes up to this seq advance the rule's state without recording firings or errors: a rebuild after a replay.

ScopeState pydantic-model

Bases: _State

Everything a rule remembers about one scope.

Fields:

scope pydantic-field

scope: dict[str, JsonValue]

The scope's field values.

DedupeState pydantic-model

Bases: _State

When each recent key was last seen.

Fields:

PatternState

PatternState = Annotated[
    EachState | CountState | SequenceState | AbsenceState,
    Field(discriminator=kind),
]

A pattern's state for one scope.

EachState pydantic-model

Bases: _State

The each pattern keeps nothing.

Fields:

CountState pydantic-model

Bases: _State

The envelopes counted within the window so far.

Fields:

SequenceState pydantic-model

Bases: _State

How far through its steps a sequence is, and the envelopes that matched them.

Fields:

AbsenceState pydantic-model

Bases: _State

The last envelope that passed; the rule's deadline for this scope follows from it.

Fields:

ThrottleState pydantic-model

Bases: _State

When the scope fired recently.

Fields:

  • fired (tuple[AwareDatetime, ...])

Match pydantic-model

Bases: _State

An envelope that counted towards a pattern.

Fields:

correlation_id pydantic-field

correlation_id: str

The causal chain the envelope belongs to.

Firing pydantic-model

Bases: _State

A rule's condition held for one scope; its run is created with the same id.

Fields:

seq pydantic-field

seq: int

The envelope at which the rule fired.

at pydantic-field

at: AwareDatetime

That envelope's time.

matched pydantic-field

matched: tuple[int, ...]

The envelopes that made the condition hold.

depth pydantic-field

depth: int = 0

The deepest causation depth among the matched envelopes.

correlation_id pydantic-field

correlation_id: str

The causal chain the firing joins: that of the latest envelope it matched.

causation property

causation: Causation

The causation of reflexr's facts about this firing.

It is one step deeper than the envelopes the firing matched.

EvaluationError pydantic-model

Bases: _State

A rule could not evaluate one envelope.

Fields:

Fact pydantic-model

Bases: _State

An event evaluation decided to append, with the causal chain and depth it belongs to.

Fields:

Evaluation pydantic-model

Bases: _State

What evaluating a batch of envelopes decided, for the host to save in one transaction.

Fields:

states pydantic-field

states: dict[ScopeKey, ScopeState] = {}

The scopes whose state changed.

facts pydantic-field

facts: tuple[Fact, ...] = ()

The facts to append to the log, in the order they happened.

Runs

A firing's run and its lifecycle, as pure transitions. See The reactor.

Run pydantic-model

Bases: BaseModel

One firing's action: its status, attempts, and a graph's latest checkpoint.

Config:

  • frozen: True
  • extra: forbid

Fields:

id pydantic-field

id: RunId

The firing's id, which is also the action's idempotency key.

correlation_id pydantic-field

correlation_id: str

The causal chain the run belongs to, which is its session in traces.

reason pydantic-field

reason: str | None = None

A stable code for why the last attempt failed, such as guardrail_blocked.

checkpoint pydantic-field

checkpoint: JsonValue = None

A graph run's state and pending steps after its last completed step.

step pydantic-field

step: str | None = None

The last step a graph run completed.

checkpoints pydantic-field

checkpoints: int = 0

How many checkpoints the run has saved, across its attempts.

trace_ids pydantic-field

trace_ids: tuple[TraceId, ...] = ()

The trace id of each attempt, so feedback on the run can be attached to its traces.

causation property

causation: Causation

The causation of the events the run emits, and of reflexr's facts about it.

It is one step deeper than the envelopes that made the rule fire.

RunStatus module-attribute

RunStatus = Literal[
    "pending",
    "running",
    "retrying",
    "succeeded",
    "dead",
    "cancelled",
    "skipped",
]

Where a run is in its lifecycle.

FINISHED module-attribute

FINISHED: frozenset[RunStatus] = frozenset(
    {"succeeded", "dead", "cancelled", "skipped"}
)

Statuses in which a run does not run again unless retried.

WAITING module-attribute

WAITING: frozenset[RunStatus] = frozenset(
    {"pending", "retrying"}
)

Statuses in which a run waits for its next attempt to start.

HOLDING module-attribute

HOLDING: frozenset[RunStatus] = frozenset(
    {"pending", "retrying", "running"}
)

Statuses in which a run of a rule ordered by scope holds back the later runs of its scope.

create_run

create_run(firing: Firing, *, now: AwareDatetime) -> Run

Return the pending run for a firing.

start

start(
    run: Run,
    *,
    now: AwareDatetime,
    trace_id: str | None = None,
) -> tuple[Run, RunStarted]

Begin an attempt of a pending or retrying run, recording its trace if there is one.

succeed

succeed(
    run: Run,
    *,
    now: AwareDatetime,
    output: JsonValue = None,
) -> tuple[Run, RunSucceeded]

Finish a running run.

fail

fail(
    run: Run,
    rule: Rule,
    *,
    now: AwareDatetime,
    error: str,
    reason: str | None = None,
    permanent: bool = False,
) -> tuple[Run, RunRetrying | RunDeadLettered]

Record a failed attempt: retry later under the rule's policy, or dead-letter the run.

Parameters:

Name Type Description Default
run Run

The running run.

required
rule Rule

Its rule, whose retry policy applies.

required
now AwareDatetime

The current time.

required
error str

What went wrong, for people.

required
reason str | None

A stable code for why, for metrics and dashboards.

None
permanent bool

Whether retrying cannot help, as when a guardrail blocked the input: the run is dead-lettered whatever its policy allows.

False

checkpoint

checkpoint(
    run: Run,
    *,
    now: AwareDatetime,
    step: str,
    state: JsonValue,
) -> tuple[Run, RunProgressed]

Save a running graph run's state after it completed step.

cancel

cancel(
    run: Run,
    *,
    now: AwareDatetime,
    reason: str | None = None,
) -> tuple[Run, RunCancelled]

Cancel a run that has not finished.

skip

skip(
    run: Run,
    *,
    now: AwareDatetime,
    reason: str | None = None,
) -> tuple[Run, RunSkipped]

Give up on a run that is waiting, unblocking later runs of its scope.

retry

retry(
    run: Run, *, now: AwareDatetime
) -> tuple[Run, RunRequeued]

Make a run runnable now, with a fresh retry budget if it had finished.

runnable

runnable(
    runs: Sequence[Run], rule: Rule, *, now: AwareDatetime
) -> list[Run]

Return the runs of one rule and scope that may start now.

Parameters:

Name Type Description Default
runs Sequence[Run]

Every run of the scope that has not succeeded, in firing order.

required
rule Rule

Their rule, whose ordering and dead-letter policy apply.

required
now AwareDatetime

The current time.

required

holds

holds(run: Run, *, blocking: bool) -> bool

Whether a run of a rule ordered by scope holds back the later runs of its scope.

It does while it is waiting or running, and a dead-lettered one does too if its rule blocks (on_dead_letter="block").

Feedback

Typed feedback on runs, firings and causal chains. See Feedback and evaluation.

Feedback pydantic-model

Bases: BaseModel

Base class for feedback types.

Field types decide how each field is scored: bounded numbers are numeric, bool is a yes/no, Literal and Enum are categories, and str is free text.

Config:

  • frozen: True
  • extra: forbid

feedback_type class-attribute

feedback_type: str

The registered type name, derived from the class name unless given as name=.

targets class-attribute

What this type of feedback can be given on.

FeedbackTarget

FeedbackTarget = Annotated[
    RunTarget | FiringTarget | ChainTarget,
    Field(discriminator=kind),
]

What a piece of feedback is about.

RunTarget pydantic-model

Bases: BaseModel

Feedback on one run: a workflow execution.

Config:

  • frozen: True
  • extra: forbid

Fields:

FiringTarget pydantic-model

Bases: BaseModel

Feedback on one firing: whether the rule should have fired.

Config:

  • frozen: True
  • extra: forbid

Fields:

ChainTarget pydantic-model

Bases: BaseModel

Feedback on a causal chain: everything one triggering event led to.

Config:

  • frozen: True
  • extra: forbid

Fields:

TargetKind

TargetKind = Literal['run', 'firing', 'chain']

What feedback can be about: a run (a workflow execution), a firing, or a causal chain.

TARGET_KINDS module-attribute

TARGET_KINDS: frozenset[TargetKind] = frozenset(
    {"run", "firing", "chain"}
)

feedback_types

feedback_types() -> Mapping[str, type[Feedback]]

Return a read-only view of every registered feedback type, by name.

load_feedback

load_feedback(
    feedback_type: str,
    target: RunTarget | FiringTarget | ChainTarget,
    value: Mapping[str, Any],
) -> Feedback

Validate feedback of a registered type for a target.

Raises:

Type Description
NotFound

If no feedback type is registered under that name.

ValidationFailed

If the type cannot be given on that kind of target, or the value does not validate.

Rejections

The ways a command can fail, and the error for an invalid rule. See Rejections.

Rejection

Rejection(message: str)

Bases: Exception

Base class for every reason reflexr refuses a command.

details

details() -> dict[str, JsonValue]

Return the rejection's typed fields; subclasses extend this.

payload

payload() -> dict[str, JsonValue]

Return the rejection as JSON-compatible data for the wire.

NotFound

NotFound(entity: str, id: str)

Bases: Rejection

Something the command refers to does not exist.

details

details() -> dict[str, JsonValue]

Return what was missing.

InvalidState

InvalidState(message: str)

Bases: Rejection

The command does not apply to the current state, such as retrying a finished run.

ValidationFailed

ValidationFailed(message: str, errors: list[JsonValue])

Bases: Rejection

The command's data does not validate, such as an event that does not match its type.

details

details() -> dict[str, JsonValue]

Return Pydantic's validation errors.

Forbidden

Forbidden(message: str)

Bases: Rejection

The actor may not perform the command.

DepthExceeded

DepthExceeded(depth: int, limit: int)

Bases: Rejection

An event would extend a causal chain beyond the workspace's limit.

details

details() -> dict[str, JsonValue]

Return the depth and the limit.

UnsupportedProtocol

UnsupportedProtocol(message: str)

Bases: Rejection

The client asked for a protocol version this server does not speak.

InvalidRule

InvalidRule(rule: str, problems: list[str])

Bases: ValueError

A rule refers to event types, fields, predicates or actions that do not exist.

Raised when rules are registered, so a mistake fails at startup rather than silently never matching.

Protocol

The stream protocol's commands, outcomes, frames and resume rule. See the stream protocol.

PROTOCOL module-attribute

PROTOCOL: Final = 'reflexr.v1'

The protocol version this library speaks.

Command

Anything a client can ask a workspace to do.

Publish pydantic-model

Bases: _Model

Publish an event. Publishing an id already in the log appends nothing.

Fields:

correlation_id pydantic-field

correlation_id: str | None = None

The causal chain to join: the id of its first event. By default, a new chain.

GiveFeedback pydantic-model

Bases: _Model

Give typed feedback on a run, a firing or a causal chain.

Fields:

RetryRun pydantic-model

Bases: _Model

Make a run runnable now.

Fields:

SkipRun pydantic-model

Bases: _Model

Give up on a waiting or dead-lettered run.

Fields:

CancelRun pydantic-model

Bases: _Model

Cancel a run that has not finished, stopping its action if it is running.

Fields:

ReplayRule pydantic-model

Bases: _Model

Reset a rule to evaluate the log again, rebuilding its state or firing again.

Fields:

InstallRule pydantic-model

Bases: _Model

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

Fields:

provenance pydantic-field

provenance: dict[str, JsonValue] = {}

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

UpdateRule pydantic-model

Bases: _Model

Replace an active stored rule with a new version, reset if its condition or scope changed.

Fields:

expected_version pydantic-field

expected_version: int | None = None

The version the change replaces. If the rule is at another, someone else changed it.

ArchiveRule pydantic-model

Bases: _Model

Archive a stored rule, and cancel its unfinished runs.

Fields:

Outcome

What a command did.

PublishedOutcome pydantic-model

Bases: _Model

What publishing did: the event's position, and whether it was already there.

Fields:

RecordedOutcome pydantic-model

Bases: _Model

Where a fact the command recorded, such as feedback, was appended.

Fields:

RunOutcome pydantic-model

Bases: _Model

A run after the command changed it.

Fields:

RuleOutcome pydantic-model

Bases: _Model

A rule's progress after the command reset it.

Fields:

RuleVersionOutcome pydantic-model

Bases: _Model

A stored rule's version after the command changed it, or found it changed already.

Fields:

seq pydantic-field

seq: int

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

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

duplicate pydantic-field

duplicate: bool = False

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

Hello pydantic-model

Bases: _Model

The first frame a client sends on a connection.

Fields:

resume_after_seq pydantic-field

resume_after_seq: int = 0

The last seq the client has: everything after it is replayed.

from_head pydantic-field

from_head: bool = False

Start at the head of the log instead, replaying nothing. resume_after_seq must then be 0; a client that reconnects resumes from welcome.head_seq with it.

types pydantic-field

types: tuple[EventName, ...] | None = None

The event types to receive; None means every type. replay_complete always arrives, since it is decided on the whole log.

Welcome pydantic-model

Bases: _Model

The server's answer to hello.

Fields:

EventFrame pydantic-model

Bases: Envelope

An envelope from the log.

Config:

  • frozen: True

Fields:

ReplayComplete pydantic-model

Bases: _Model

Every event up to up_to_seq has been replayed; what follows is live.

Fields:

  • type (Literal['replay_complete'])
  • up_to_seq (int)

CommandFrame pydantic-model

Bases: _Model

A client's command, with an id that correlates it with its result.

command_id is also an idempotency key: the server deduplicates repeated ids.

Fields:

CommandResult pydantic-model

Bases: _Model

The result of one command frame.

Fields:

ErrorFrame pydantic-model

Bases: _Model

A frame the server could not understand.

Fields:

ClientFrame

ClientFrame = Annotated[
    Hello | CommandFrame, Field(discriminator=type)
]

Any frame a client sends.

ServerFrame

ServerFrame = Annotated[
    Welcome
    | EventFrame
    | ReplayComplete
    | CommandResult
    | ErrorFrame,
    Field(discriminator=type),
]

Any frame the server sends.

resume

resume(hello: Hello, *, head_seq: int) -> ResumePlan

Decide where replay starts for a connecting client.

Parameters:

Name Type Description Default
hello Hello

The client's hello frame.

required
head_seq int

The workspace log's latest seq (0 if it is empty).

required

Raises:

Type Description
UnsupportedProtocol

If the client speaks another protocol version.

ValidationFailed

If the client asks to start at the head and to resume.

ResumePlan pydantic-model

Bases: _Model

Where a connection's replay starts.

Fields:

replay_after pydantic-field

replay_after: int

Replay every event with a seq greater than this.

reset pydantic-field

reset: bool = False

Whether the client must discard what it has and rebuild from the replay.

Identifiers

Identifiers are plain strings. These aliases say what a string identifies, and the factories generate or derive ids with a short type prefix.

TenantId

TenantId = str

Identifies a tenant: the top-level isolation boundary.

WorkspaceId

WorkspaceId = str

Identifies a workspace within a tenant. A workspace has one event log.

EventId

EventId = str

Identifies an event within a workspace. Publishing an id that already exists is a no-op.

RuleName

RuleName = str

Identifies a rule. Rules are named by their authors, not generated.

ScopeKey

ScopeKey = str

Identifies one scope of a rule: the canonical JSON of its scope fields' values.

FiringId

FiringId = str

Identifies a firing. It is also the id of the run the firing starts.

RunId

RunId = str

Identifies a run: the execution of one firing's action.

TraceId

TraceId = str

A W3C trace id, as 32 lowercase hex digits: the trace of one run attempt.

new_id

new_id(prefix: str) -> str

Return a new random identifier with the given prefix.

Parameters:

Name Type Description Default
prefix str

A short type tag, such as "evt".

required

Returns:

Type Description
str

An identifier such as evt_3f9c2a1b7d4e5f60.

new_event_id

new_event_id() -> EventId

Return a new event identifier.

firing_id

firing_id(
    rule: RuleName,
    generation: int,
    scope: ScopeKey,
    seq: int,
) -> FiringId

Return the id of the firing of rule for scope at seq.

Parameters:

Name Type Description Default
rule RuleName

The rule that fired.

required
generation int

The rule's generation, which changes when its state is reset.

required
scope ScopeKey

The scope that fired.

required
seq int

The envelope at which it fired.

required

Returns:

Type Description
FiringId

A deterministic identifier such as fir_0b6a3c9e1f2d4a58.

derived_event_id

derived_event_id(
    run: RunId, index: int, *, segment: int = 0
) -> EventId

Return the id of the index-th event a run emits after its segment-th checkpoint.

Every attempt of a run derives the same ids in the same order from where it starts, so a retried run that emits again publishes duplicates, which the log ignores. A run resumed from its n-th checkpoint emits in segment n, so its events never collide with the ones emitted before that checkpoint.