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:Trueextra:forbid
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:Trueextra:allow
Fields:
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(int) -
id(EventId) -
ts(AwareDatetime) -
workspace_id(WorkspaceId) -
actor(Actor) -
causation(Causation | None) -
correlation_id(str) -
traceparent(str | None) -
event(AnyEvent)
ts
pydantic-field
¶
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.
Causation
pydantic-model
¶
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.
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
¶
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
¶
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 |
dump_event
¶
Return an event's data with its "type" first, as stored and sent on the wire.
type_of
¶
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
¶
SystemEvent = (
RuleFired
| RuleErrored
| RuleReset
| RuleInstalled
| RuleArchived
| RunStarted
| RunProgressed
| RunRetrying
| RunSucceeded
| RunDeadLettered
| RunCancelled
| RunRequeued
| RunSkipped
)
reflexr's facts about rules and runs. Each names the rule it is about.
SYSTEM_EVENTS
module-attribute
¶
SYSTEM_EVENTS: tuple[type[Event], ...] = (
RuleFired,
RuleErrored,
RuleReset,
RuleInstalled,
RuleArchived,
RunStarted,
RunProgressed,
RunRetrying,
RunSucceeded,
RunDeadLettered,
RunCancelled,
RunRequeued,
RunSkipped,
FeedbackGiven,
Tick,
)
Every event type reflexr itself appends.
about_rule
¶
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:
RuleErrored
pydantic-model
¶
RuleReset
pydantic-model
¶
Bases: _ReflexrEvent
A rule's state was reset, because its definition changed or it was replayed.
Fields:
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:
RuleArchived
pydantic-model
¶
RunStarted
pydantic-model
¶
RunProgressed
pydantic-model
¶
RunRetrying
pydantic-model
¶
Bases: _ReflexrEvent
An attempt failed, and the run will be tried again.
Fields:
RunSucceeded
pydantic-model
¶
RunDeadLettered
pydantic-model
¶
RunCancelled
pydantic-model
¶
RunRequeued
pydantic-model
¶
RunSkipped
pydantic-model
¶
FeedbackGiven
pydantic-model
¶
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
¶
Actor = Annotated[
UserActor
| AgentActor
| ExternalAgentActor
| SystemActor
| SourceActor
| EvaluatorActor,
Field(discriminator="kind"),
]
Any actor, discriminated by kind.
UserActor
pydantic-model
¶
Bases: BaseModel
A person, identified by the host application's user id.
Config:
frozen:True
Fields:
AgentActor
pydantic-model
¶
Bases: BaseModel
An agent or graph running an action for one firing of a rule.
Config:
frozen:True
Fields:
ExternalAgentActor
pydantic-model
¶
SystemActor
pydantic-model
¶
Bases: BaseModel
reflexr itself (the reactor, schedules), or the application acting on its own behalf.
Config:
frozen:True
Fields:
SourceActor
pydantic-model
¶
Bases: BaseModel
A system that publishes events, such as a monitoring service or a webhook.
Config:
frozen:True
Fields:
EvaluatorActor
pydantic-model
¶
Bases: BaseModel
An evaluator: a judge or decision model whose verdicts are recorded as feedback.
Config:
frozen:True
Fields:
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:Trueextra:forbid
Fields:
-
name(RuleName) -
description(str) -
when(Condition) -
scope(Scope) -
then(ActionRef) -
retry(RetryPolicy) -
ordering(Literal['scope', 'none']) -
on_dead_letter(Literal['continue', 'block']) -
start(Literal['now', 'beginning']) -
timeout(timedelta | None) -
enabled(bool)
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:Trueextra:forbid
Fields:
by
¶
Scope a rule by event fields: by(F.service), by(F.service, F.region).
scope_key
¶
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:Trueextra:forbid
Fields:
run
¶
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
¶
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 |
Named
¶
RetryPolicy
pydantic-model
¶
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 |
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 |
required |
Raises:
| Type | Description |
|---|---|
ValueError
|
If |
RuleChange
pydantic-model
¶
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
¶
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:Trueextra:forbid
Fields:
MAX_STORED_FIRINGS_PER_HOUR
module-attribute
¶
The most firings a stored rule's throttle may allow in an hour.
MAX_STORED_ATTEMPTS
module-attribute
¶
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
¶
The most characters a stored rule's description may have.
MAX_STORED_BYTES
module-attribute
¶
The most bytes a stored rule's JSON may have.
MAX_STORED_PROVENANCE
module-attribute
¶
The most bytes a change's provenance may have, as JSON.
MAX_STORED_RULES
module-attribute
¶
The most active stored rules a workspace may have.
Conditions¶
The builder, and the stages it produces. See Conditions.
on
¶
Start a condition on envelopes of the given event types, by class or by name.
sequence
¶
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).
field
¶
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
¶
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.
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.
Condition
pydantic-model
¶
Bases: _Stage
A rule's when: a filter, and optionally dedupe, a pattern and a throttle.
Fields:
where
¶
Narrow the filter: every filter given, and every field equal to its keyword.
distinct
¶
Drop an envelope whose key fields were already seen within within.
count
¶
Fire when at_least envelopes pass within within.
absent
¶
Fire when nothing has passed for within since the last envelope that did.
Filter
¶
Filter = Annotated[
OnFilter
| WhereFilter
| AllFilter
| AnyFilter
| NotFilter
| PredicateFilter,
Field(discriminator=kind),
]
Which envelopes a rule considers.
OnFilter
pydantic-model
¶
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
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
¶
AnyFilter
pydantic-model
¶
NotFilter
pydantic-model
¶
PredicateFilter
pydantic-model
¶
Dedupe
pydantic-model
¶
Pattern
¶
Pattern = Annotated[
EachPattern
| CountPattern
| SequencePattern
| AbsencePattern,
Field(discriminator=kind),
]
When the envelopes a rule considers make it fire.
EachPattern
pydantic-model
¶
CountPattern
pydantic-model
¶
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
¶
Throttle
pydantic-model
¶
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 |
required |
states
|
Mapping[ScopeKey, ScopeState]
|
The state of the scopes |
required |
envelopes
|
Sequence[Envelope]
|
New envelopes, in |
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
|
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 |
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
|
0
|
matches
¶
matches(
filter: Filter,
envelope: Envelope,
predicates: Predicates,
) -> bool
Whether an envelope passes a filter. Predicates may raise.
Predicates
¶
Registered Python predicates, by name. They must be pure.
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(int) -
generation(int) -
definition(str) -
deadlines(dict[ScopeKey, AwareDatetime]) -
silent_through(int)
ScopeState
pydantic-model
¶
Bases: _State
Everything a rule remembers about one scope.
Fields:
-
scope(dict[str, JsonValue]) -
dedupe(DedupeState | None) -
pattern(PatternState) -
throttle(ThrottleState | None)
DedupeState
pydantic-model
¶
PatternState
¶
PatternState = Annotated[
EachState | CountState | SequenceState | AbsenceState,
Field(discriminator=kind),
]
A pattern's state for one scope.
EachState
pydantic-model
¶
CountState
pydantic-model
¶
SequenceState
pydantic-model
¶
AbsenceState
pydantic-model
¶
ThrottleState
pydantic-model
¶
Match
pydantic-model
¶
Firing
pydantic-model
¶
Bases: _State
A rule's condition held for one scope; its run is created with the same id.
Fields:
EvaluationError
pydantic-model
¶
Fact
pydantic-model
¶
Bases: _State
An event evaluation decided to append, with the causal chain and depth it belongs to.
Fields:
-
event(RuleFired | RuleErrored) -
correlation_id(str) -
causation(Causation | None)
Evaluation
pydantic-model
¶
Bases: _State
What evaluating a batch of envelopes decided, for the host to save in one transaction.
Fields:
-
progress(RuleProgress) -
states(dict[ScopeKey, ScopeState]) -
firings(tuple[Firing, ...]) -
errors(tuple[EvaluationError, ...]) -
facts(tuple[Fact, ...])
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:Trueextra:forbid
Fields:
-
id(RunId) -
rule(RuleName) -
scope_key(ScopeKey) -
scope(dict[str, JsonValue]) -
fired_seq(int) -
matched(tuple[int, ...]) -
correlation_id(str) -
depth(int) -
status(RunStatus) -
attempts(int) -
next_attempt_at(AwareDatetime) -
created_at(AwareDatetime) -
updated_at(AwareDatetime) -
error(str | None) -
reason(str | None) -
checkpoint(JsonValue) -
step(str | None) -
checkpoints(int) -
output(JsonValue) -
trace_ids(tuple[TraceId, ...])
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
¶
A graph run's state and pending steps after its last completed step.
checkpoints
pydantic-field
¶
checkpoints: int = 0
How many checkpoints the run has saved, across its attempts.
trace_ids
pydantic-field
¶
The trace id of each attempt, so feedback on the run can be attached to its traces.
RunStatus
module-attribute
¶
RunStatus = Literal[
"pending",
"running",
"retrying",
"succeeded",
"dead",
"cancelled",
"skipped",
]
Where a run is in its lifecycle.
FINISHED
module-attribute
¶
Statuses in which a run does not run again unless retried.
WAITING
module-attribute
¶
Statuses in which a run waits for its next attempt to start.
HOLDING
module-attribute
¶
Statuses in which a run of a rule ordered by scope holds back the later runs of its scope.
create_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
¶
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
¶
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:Trueextra:forbid
feedback_type
class-attribute
¶
feedback_type: str
The registered type name, derived from the class name unless given as name=.
targets
class-attribute
¶
targets: frozenset[TargetKind]
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
¶
FiringTarget
pydantic-model
¶
ChainTarget
pydantic-model
¶
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
¶
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.
NotFound
¶
ValidationFailed
¶
DepthExceeded
¶
InvalidRule
¶
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
¶
Command = Annotated[
Publish
| GiveFeedback
| RetryRun
| SkipRun
| CancelRun
| ReplayRule
| InstallRule
| UpdateRule
| ArchiveRule,
Field(discriminator=type),
]
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:
-
type(Literal['publish']) -
event(AnyEvent) -
id(EventId | None) -
correlation_id(str | None)
GiveFeedback
pydantic-model
¶
Bases: _Model
Give typed feedback on a run, a firing or a causal chain.
Fields:
-
type(Literal['give_feedback']) -
feedback_type(str) -
target(FeedbackTarget) -
value(dict[str, JsonValue])
RetryRun
pydantic-model
¶
SkipRun
pydantic-model
¶
CancelRun
pydantic-model
¶
ReplayRule
pydantic-model
¶
InstallRule
pydantic-model
¶
Bases: _Model
Install a stored rule: its first version, or the next version of an archived one.
Fields:
-
type(Literal['install_rule']) -
rule(Rule) -
provenance(dict[str, JsonValue])
UpdateRule
pydantic-model
¶
Bases: _Model
Replace an active stored rule with a new version, reset if its condition or scope changed.
Fields:
ArchiveRule
pydantic-model
¶
Outcome
¶
Outcome = Annotated[
PublishedOutcome
| RecordedOutcome
| RunOutcome
| RuleOutcome
| RuleVersionOutcome,
Field(discriminator=type),
]
What a command did.
PublishedOutcome
pydantic-model
¶
RecordedOutcome
pydantic-model
¶
RunOutcome
pydantic-model
¶
RuleOutcome
pydantic-model
¶
Bases: _Model
A rule's progress after the command reset it.
Fields:
-
type(Literal['rule']) -
rule(RuleName) -
progress(RuleProgress)
RuleVersionOutcome
pydantic-model
¶
Bases: _Model
A stored rule's version after the command changed it, or found it changed already.
Fields:
Hello
pydantic-model
¶
Bases: _Model
The first frame a client sends on a connection.
Fields:
-
type(Literal['hello']) -
protocol(str) -
resume_after_seq(int) -
from_head(bool) -
types(tuple[EventName, ...] | None)
resume_after_seq
pydantic-field
¶
resume_after_seq: int = 0
The last seq the client has: everything after it is replayed.
Welcome
pydantic-model
¶
Bases: _Model
The server's answer to hello.
Fields:
-
type(Literal['welcome']) -
protocol(str) -
workspace_id(WorkspaceId) -
head_seq(int) -
reset(bool)
EventFrame
pydantic-model
¶
Bases: Envelope
An envelope from the log.
Config:
frozen:True
Fields:
-
seq(int) -
id(EventId) -
ts(AwareDatetime) -
workspace_id(WorkspaceId) -
actor(Actor) -
causation(Causation | None) -
correlation_id(str) -
traceparent(str | None) -
event(AnyEvent) -
type(Literal['event'])
ReplayComplete
pydantic-model
¶
CommandFrame
pydantic-model
¶
CommandResult
pydantic-model
¶
ErrorFrame
pydantic-model
¶
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 |
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
¶
Identifiers¶
Identifiers are plain strings. These aliases say what a string identifies, and the factories generate or derive ids with a short type prefix.
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.
ScopeKey
¶
ScopeKey = str
Identifies one scope of a rule: the canonical JSON of its scope fields' values.
new_id
¶
firing_id
¶
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 |
derived_event_id
¶
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.