Getting data in · Workflow code
Writing a workflow in code
How a command, its validations and its actions are written in the flows plugin today, read line by line from the real code.
In one minute
A workflow is a Python class in the flows plugin (packages/flows/src/omnium_flows/). It is built from small functions, each marked with a decorator: @command (what a person or a feed may send), @validation (a yes or no, with the sentence a person reads) and @action (what is now true, returned as rows). The engine finds these functions by a key such as football.goal or results.edit_side. It never imports the plugin by name.
A function never touches the database. A validation returns a Verdict. An action returns rows: a scoreboard, Lines, Places or Records. The engine does every write, in one transaction, next to the log row.
There are two engines today. Apply mode (workflows/engine.py) runs Results and Games commands once against their own input. Fold mode (scoring/pipeline.py) runs a sport's live scoring by replaying the match log. Both call the same decorated functions, through omnium_core/flows_host.py.
Every code block on this page is real, copied from the repo on 7 Oct 2026, with its file and line numbers. "..." marks lines left out.
What a workflow is
A workflow is one class that says what can change, when it may change, and what changes. A sport's live scoring is one workflow. Results (any fixture in any sport) is another. Games operations (medal table, schedule, players, imports) is a third.
Every change in omnium is a command. A command is a named request with an input, like "goal, for united, scored by rashford". The engine takes it through the same steps every time:
- Find the command by its code.
- Check the input against the command's input model.
- Run its validations, in order. The first "no" is the answer.
- Write the log row.
- Run its actions, and write the rows they return.
- Answer the sender.

The words you will meet on this page:
| Word | In one sentence |
|---|---|
| Component | One decorated function: a command, validation, action, branch or context. |
| Key | A component's permanent name: workflow code, a dot, the function name (football.goal). |
| Registry | The list of every component the installed plugin has, built at start-up. |
| Ledger, or log | The timeline_item table: one row per accepted command, never edited. |
| Apply mode | Each command's actions run once, on its own input. Used by Results and Games. |
| Fold mode | Each command replays the whole match log to rebuild the state. Used by sport scoring. |
| Derived state | The match as it stands, saved in the scoreboard table's snapshot column. |
The six files of a sport
Every sport folder holds the same six files, and a check in make check fails the build if one is missing. Someone who has read one sport can find their way in any other.
| File | What it holds | Example from football |
|---|---|---|
constants.py | Words the sport uses, and its format model | GOAL_KINDS, CARD_COLOURS, FootballFormat (halves) |
commands.py | Input models, and one @command per thing a person can send | match_start, period_start, goal, card, period_end, match_end, abandon |
validations.py | One @validation per rule, each with its sentence | not_started, two_sides, match_is_live, no_half_open, half_is_open, played_in_full |
actions.py | @actions that return the state, the lines and the places | state, tallies, placings |
stats.py | The sport's stat cards, declared as data | football.scoring.node, football.top_scorers.node, football.table.node, football.scoring.career |
sport.py | The class that joins the other five, its code, its fact list, its context | class Football(...), code = "football", facts, scoreboard |
The football files have 30, 103, 61, 187, 48 and 79 lines (counted on 7 Oct). Results and Games use the same six files, even though they are not sports.
The class in sport.py is built from the other files by mixing classes: each file holds one class, and sport.py inherits from all of them.
class Football(
FootballCommands,
FootballValidations,
FootballActions,
FootballTotals,
Sport,
):
code = "football"
format_model = FootballFormat
...
roles: ClassVar[frozenset[str]] = frozenset({"scorer", "assist", "carded"})
packages/flows/src/omnium_flows/football/sport.py:40-48, 64
code = "football"is permanent. It becomes the first part of every key.format_modelis the shape of the fixture's format. Football reads onlyhalves.roleslists the roles a command may name a person in. A goal names ascorerand anassist.Sportis the base class from the contract package. It collects the components when Python creates the class.
How the class becomes keys
The decorators do no work when they run; the base class turns each decorated function into a key. This code runs once, when Python loads the class.
for name in dir(cls):
if name.startswith("_"):
continue
member = getattr(cls, name, None)
spec = spec_of(member) if callable(member) else None
if spec is None:
continue
owned_here = name in cls.__dict__
inherited_from = None
if parent is not None and not owned_here:
inherited_from = f"{parent.code}.{name}"
else:
overrode_anything = True
components[name] = RegisteredComponent(
key=f"{code}.{name}",
sport_code=code,
spec=spec,
fn=member,
inherited_from=inherited_from,
)
packages/contract/src/omnium_contract/flows/sport.py:141-160
- It walks every member of the class and keeps those that carry a decorator's spec.
- The key is
codeplus the function name:football.goal,football.half_is_open. - A format (a subclass, like cricket's The Hundred) inherits the rest and overrides a few. The same file refuses a format of a format, a format that changes nothing, and a format with no
scope(sport.py:130-176).
How omnium finds the plugin
omnium discovers workflows through Python entry points, so adding a sport is one line in the plugin's pyproject.toml. An entry point is a name a package publishes so other code can find it without importing it by name.
# How omnium discovers sports. New sport = one line here, zero omnium changes.
[project.entry-points."omnium.flows"]
cricket = "omnium_flows.cricket"
swimming = "omnium_flows.swimming"
beach_volleyball = "omnium_flows.beach_volleyball"
volleyball = "omnium_flows.volleyball"
boxing = "omnium_flows.boxing"
football = "omnium_flows.football"
athletics = "omnium_flows.athletics"
# Workflows that are not one sport's scoring: an event's own operations, and
# the results of any fixture in any sport.
games = "omnium_flows.games"
results = "omnium_flows.results"
packages/flows/pyproject.toml:25-37
def load_registry(group: str = ENTRY_POINT_GROUP) -> Registry:
...
registry = Registry()
for entry_point in importlib_metadata.entry_points(group=group):
module = entry_point.load()
found = [
member
for member in vars(module).values()
if isinstance(member, type)
and issubclass(member, Workflow)
and not member.__dict__.get("__abstract__", False)
]
if not found:
raise RegistryError(
f"entry point {entry_point.name!r} ({entry_point.value}) defines no Workflow"
)
for workflow in found:
registry.add_workflow(workflow)
return registry
packages/contract/src/omnium_contract/flows/registry.py:94-117
load_registryimports each module named in the groupomnium.flowsand registers every workflow class in it.add_workflow(registry.py:35-43) refuses a workflow code or a component key declared twice.- omnium calls this once per process, through
flows_host.registry()(packages/core/src/omnium_core/flows_host.py:90-96). That function first checks theOMNIUM_FLOWS_VERSIONpin against the installed plugin version, if the pin is set.
A command, line by line Built today
A command is an input model plus a decorated, empty method. The method body is ...: a command does nothing itself. Its decorator says what input it takes, which rules guard it and which actions it runs.
The input model first. Football's goal:
class _Input(BaseModel):
model_config = ConfigDict(alias_generator=to_camel, populate_by_name=True)
...
class GoalInput(_Input):
entry: str
kind: GoalKind
scorer: str | None = None
assist: str | None = None
minute: int | None = Field(default=None, ge=1)
packages/flows/src/omnium_flows/football/commands.py:14-15, 27-32
- The input is a pydantic model: a Python class that checks data against typed fields. This model is the command's schema; there is no second copy.
alias_generator=to_camelmeans the app may sendminuteor a camelCase name, and both work.kind: GoalKindis a fixed list of words fromconstants.py:open_play,penalty,free_kick,own_goal. Any other word is refused before a rule runs.Field(default=None, ge=1)means a minute, if sent, is 1 or more.
Then the command itself:
class FootballCommands:
...
@command(
input=GoalInput,
validations=("half_is_open",),
actions=("tallies",),
actors=(
{"path": "input.scorer", "role": "scorer", "optional": True},
{"path": "input.assist", "role": "assist", "optional": True},
),
description="One goal, to one side, and how it came about. An own goal "
"names the side it counts for and the player who conceded it.",
)
def goal(self) -> None: ...
packages/flows/src/omnium_flows/football/commands.py:42, 57-68
| Argument | What it does |
|---|---|
input=GoalInput | The model the input must match. The engine checks it first. |
validations=("half_is_open",) | The rules to run, by function name, in this order. |
actions=("tallies",) | The row-making actions to run after the log row is written. The state action runs on every command and is not listed. |
actors=(...) | Which input fields name a person, and in what role. Each one becomes a timeline_actor row next to the log row. optional means the field may be empty. |
description | The sentence screens show for this command. |
The decorator also takes branches (routing that may add actions), stats (stat keys this command sets off), applies_to (only some fixture types), imports and fields. Here is its full signature:
def command(
*,
input: type[BaseModel], # mirrors the design docs
validations: tuple[str, ...] | list[str] = (),
actions: tuple[str, ...] | list[str] = (),
branches: tuple[str, ...] | list[str] = (),
stats: tuple[str, ...] | list[str] = (),
actors: tuple[dict[str, Any], ...] | list[dict[str, Any]] = (),
applies_to: tuple[str, ...] | list[str] = (),
description: str = "",
imports: bool = False,
fields: tuple[dict[str, Any], ...] | list[dict[str, Any]] = (),
) -> _Decorator:
packages/contract/src/omnium_contract/flows/components.py:131-143
The decorator only attaches a spec to the function and returns it unchanged (components.py:124-128, 156-176). That is why a test can call any component like a plain function.
A Results command
Results commands look the same. The one the console uses to fix a side:
class EditSideInput(_Input):
side: uuid.UUID
participant_id: uuid.UUID | None = None
label: str | None = Field(default=None, max_length=128)
score: str | None = Field(default=None, max_length=32)
rank: int | None = Field(default=None, ge=1)
status: str | None = None
#: The medal this side won: Gold, Silver, Bronze. "" takes it away.
medal: str | None = Field(default=None, max_length=32)
...
world_record: bool | None = None
...
event_record: bool | None = None
#: Whether this side went through, as the feeds print it (``qualified``).
qualified: bool | None = None
...
@command(
input=EditSideInput,
validations=(
"side_is_in_fixture",
"side_keeps_someone",
"side_status_is_known",
"side_medal_is_known",
),
actions=("side_changes",),
description=(
"Correct one side: who it is, its score, its place, its medal, its status, "
"whether its mark is a world record and whether it qualified."
),
)
def edit_side(self) -> None: ...
packages/flows/src/omnium_flows/results/commands.py:81-100, 168-182
- Results'
_Inputaddsextra="forbid"(results/commands.py:33-34). A field the model does not know is refused, not dropped. Football's_Inputdoes not set it. - Every field but
sideis optional. Only the fields the person sent are changed; that matters in the action below. - Four rules, in order. The order is the message order: if two rules would fail, the person sees the first.
A Games command, and an import
Games commands act on an event (a whole Games), not a fixture. One is a person's command; the other may only be sent by an integration.
@command(
input=SetMedalsInput,
validations=("noc_is_a_code", "sets_a_number"),
actions=("medals_changed",),
description="Type a country's counts or place. Each number typed is kept from the feed.",
)
def set_medals(self) -> None: ...
...
@command(
input=ImportMedalTableInput,
actions=("medal_table_imported",),
imports=True,
fields=MEDAL_FIELDS,
description="One country per row, with its golds, silvers and bronzes.",
)
def import_medal_table(self) -> None: ...
packages/flows/src/omnium_flows/games/commands.py:243-249, 412-419
imports=Truemarks a command only an integration may send. The engine refuses it from a person, and refuses every other command from an integration (packages/core/src/omnium_core/workflows/engine.py:259-263).fields=MEDAL_FIELDSlists the columns a no-code Scout integration maps, with labels and hints for the screen (games/constants.py:97onwards).
Validations, line by line Built today
A validation reads the input and the context, and returns ok() or reject("a sentence"). It never writes. It is pure: the same input and context always give the same answer.
The contract's three small types:
class Verdict(BaseModel):
"""What a validation says: yes, or no with the message the scorer reads."""
model_config = ConfigDict(frozen=True)
ok: bool
message: str | None = None
def ok() -> Verdict:
return Verdict(ok=True)
def reject(message: str) -> Verdict:
"""Refuse, with the sentence a scorer will read. Write it for them."""
return Verdict(ok=False, message=message)
packages/contract/src/omnium_contract/flows/types.py:179-194
A football rule:
class FootballValidations:
@validation(message="The match has already begun")
def not_started(self, input: Any, ctx: Ctx[FootballFormat]) -> Verdict:
"""Refuses a second kick-off — a match starts once.
...
"""
if not ctx.derived.get("started"):
return ok()
return reject("The match has already begun")
...
@validation(message="No half is open")
def half_is_open(self, input: Any, ctx: Ctx[FootballFormat]) -> Verdict:
"""Goals, cards and the end of a half all need a half in progress."""
if ctx.derived.get("halfOpen") is True:
return ok()
return reject("No half is open")
packages/flows/src/omnium_flows/football/validations.py:15-26, 49-54
@validation(message=...)stores a fallback sentence. It is used when the function says no without its own sentence.input: Anymeans the rule gets the raw input dictionary. If the type is an input model, the engine builds that model first (flows_host.py:184-195).ctx: Ctx[FootballFormat]is everything a sport rule may read: the fixture, the log, the derived state, stats, masters and child fixtures (types.py:112-126).- In fold mode,
ctx.derivedis the saved scoreboard as it was before this command (scoring/pipeline.py:1179-1188). That is whyhalf_is_opencan answer by reading one flag.
Results and Games rules get a different context, RecordCtx: the event, the fixture with its sides, who sent it (person or integration), and the sender's name (types.py:144-161).
@validation(message="That side is not part of this fixture")
def side_is_in_fixture(self, input: Any, ctx: RecordCtx) -> Verdict:
"""A side from another fixture would change the wrong match.
Read from the raw input: the edit and the removal both carry ``side``.
"""
if str(input.get("side")) in _side_ids(ctx):
return ok()
return reject("That side is not part of this fixture")
...
@validation(message="A medal is Gold, Silver or Bronze")
def side_medal_is_known(self, input: EditSideInput, ctx: RecordCtx) -> Verdict:
"""A medal is one of three words, or nothing at all.
...
"""
if input.medal is None or not input.medal.strip():
return Verdict(ok=True)
return Verdict(ok=input.medal.strip().casefold() in SIDE_MEDALS)
packages/flows/src/omnium_flows/results/validations.py:67-75, 93-103
side_is_in_fixtureis shared byedit_sideandremove_side, so it reads the raw input.side_medal_is_knownreturnsVerdict(ok=False)with no message. The engine then uses the decorator's sentence: "A medal is Gold, Silver or Bronze".
What happens when one fails
The engine runs the rules in the order the command lists them, and the first "no" is the answer. Nothing is written. Here is the apply-mode loop:
def _rules(
program: Any, refs: tuple[str, ...], payload: dict[str, Any], scope: dict[str, Any]
) -> str | None:
"""Each rule in the order the command lists it; the first no is the answer."""
for ref in refs:
rule = next((v for v in program.validations if v.id == ref), None)
if rule is None:
return f"rule {ref!r} is missing from this workflow"
message = rule.message or rule.code
try:
if rule.checks:
...
value = flows_host.function_for(rule.impl)(
payload, scope, determinism=Determinism(clock_nanos=0), limits=SYNC_LIMITS
).value
except (ExprError, SandboxError) as exc:
return f"{rule.code}: {exc}"
if value is True or (isinstance(value, dict) and value.get("ok") is True):
continue
if isinstance(value, dict) and value.get("message"):
return str(value["message"])
return str(message)
return None
packages/core/src/omnium_core/workflows/engine.py:335-361
rule.implis the component key, likeresults.side_medal_is_known.flows_host.function_forturns it into a callable.flows_hostturns aVerdictintoTrue, or into{"ok": False, "message": ...}(flows_host.py:271-278).- The sentence comes from
reject(...)if given, else from@validation(message=...), else the rule's code. - A rule that crashes, or runs over its time budget, is also a "no", with the error as the sentence.
- The fold pipeline does the same in
_validate(scoring/pipeline.py:1037-1078) and also records each step in a trace.
What the person sees: the HTTP answer is 200 with accepted: false, refusal: "rejected" and the sentence in message (packages/admin/src/omnium_admin/routers/bridge.py:246-259). A rule's "no" is not an HTTP error. The console's send helper turns it into a Refused error that the screen shows (omnium-console/src/lib/bridge.ts:113-127, the separate console repo).
Actions, line by line Built today
An action answers "what is now true because of this command?" by returning rows. It never writes them. What it returns is set by gives, and the engine decides where each kind lands.
GIVES_TABLE: dict[Gives, str] = {
Gives.STATE: "scoreboard",
Gives.LINES: "stat_value",
Gives.PLACES: "fixture_competitor_rank",
Gives.NOTHING: "",
#: Not one table: each record kind has its own writer in ``omnium_core.workflows``.
Gives.RECORDS: "",
}
packages/core/src/omnium_core/scoring/program.py:106-113
gives | Returns | Written to | Used by |
|---|---|---|---|
STATE | One dictionary: the match now | scoreboard.snapshot | Sport scoring (fold) |
LINES | A list of Line: each subject's numbers | stat_value | Sport scoring (fold) |
PLACES | A list of Place: the finishing order | fixture_competitor_rank, and rank on the lines | Sport scoring (fold) |
RECORDS | A list of Record(kind, data) | Whatever the writer for kind writes | Results, Games (apply) |
NOTHING | Nothing; side work only | Nothing; it may only run later | Not used by any workflow today (searched 7 Oct) |
Football's state action
The state action replays the log, event by event, and returns the match as one dictionary. This is the "fold".
class FootballActions:
@action(gives=Gives.STATE, reads=["log"])
def state(self, log: tuple[Event, ...], ctx: Ctx[FootballFormat]) -> dict[str, Any]:
...
halves_to_play = ctx.fixture.format.halves or 2
...
for event in log:
payload = event.input
if event.code == "football.match_start":
...
elif event.code == "football.goal":
if current is None or current["closed"] or done:
continue
entry = payload["entry"]
goals[entry] = goals.get(entry, 0) + 1
current["goals"][entry] = current["goals"].get(entry, 0) + 1
scorer = payload.get("scorer")
if scorer:
if payload.get("kind") == "own_goal":
_line(by_player, scorer)["ownGoals"] += 1
else:
_line(by_player, scorer)["goals"] += 1
assist = payload.get("assist")
if assist and payload.get("kind") != "own_goal":
_line(by_player, assist)["assists"] += 1
...
half_open = current is not None and not current["closed"]
return {
"started": started,
"done": done,
"abandoned": abandoned,
"sides": sides,
"goals": goals,
"halves": halves,
"halvesPlayed": len(halves),
"halfOpen": half_open,
"playedInFull": len(halves) >= halves_to_play and not half_open,
"byPlayer": by_player,
"winner": winner,
"drawn": drawn,
}
packages/flows/src/omnium_flows/football/actions.py:33-38, 49-51, 64-78, 105-119
reads=["log"]says it needs the whole log. The engine hands it every effective event of the match.- A goal outside an open half is skipped, not counted. The validation
half_is_openstops most of those; the fold is the second guard. - An own goal counts for the side named in
entry, and goes on the scorer's line asownGoals, nevergoals. - The returned flags (
halfOpen,playedInFull,started,done) are exactly what the validations read on the next command.
Football's lines and places
@action(
gives=Gives.LINES,
reads=["derived"],
keys=(
"goals",
"goals_for",
"goals_against",
"goals_scored",
"assists",
"own_goals",
"yellow_cards",
"red_cards",
),
)
def tallies(self, log: tuple[Event, ...], ctx: Ctx[FootballFormat]) -> list[Line]:
"""Goals a side per half, then per player: goals, assists, own goals, cards."""
derived = ctx.derived
if not derived.get("started") or derived.get("abandoned"):
return []
lines = [
Line(
entry=side,
segment=f"half-{half['number']}",
values={"goals": half["goals"].get(side, 0)},
)
for half in derived.get("halves") or []
for side in derived.get("sides") or []
]
...
packages/flows/src/omnium_flows/football/actions.py:121-148
reads=["derived"]means it reads the state thestateaction just returned, not the log.keyslists every number this action may write. Each must be in the sport'sfactslist insport.py. A typo is refused when the program is published (contract/.../flows/sport.py:80-88).- A
Lineis the whole truth for one subject in one segment. The engine replaces the set the action returned and deletes rows it returned last time but not now (scoring/projections.py:1-19). That is why a correction needs no undo code. segment="half-1"puts a side's goals in a part of the match. A line with no segment is about the whole match.
@action(gives=Gives.PLACES, reads=["derived"])
def placings(self, log: tuple[Event, ...], ctx: Ctx[FootballFormat]) -> list[Place]:
"""The finishing order. A draw is both sides on rank 1; an abandoned
match places nobody, which a table reads as a no-result."""
derived = ctx.derived
if not derived.get("done") or derived.get("abandoned"):
return []
sides = derived["sides"]
goals = derived["goals"]
def place(side: str, ordinal: int, *, shared: bool = False) -> Place:
scored = goals.get(side, 0)
return Place(entry=side, rank=ordinal, shared=shared, display=str(scored), value=scored)
if derived.get("drawn"):
return [place(sides[0], 1, shared=True), place(sides[1], 1, shared=True)]
winner = derived["winner"]
loser = sides[1] if sides[0] == winner else sides[0]
return [place(winner, 1), place(loser, 2)]
packages/flows/src/omnium_flows/football/actions.py:169-187
- A match that is not over places nobody: an empty list.
- A draw is both sides at rank 1 with
shared=True. match_endandabandonlist this action; a goal does not, so a goal never writes places.
Results' actions return records
Apply-mode actions take the input and a RecordCtx, and return Records. A Record names a writer in the engine and gives it data.
def _sent(model: BaseModel) -> dict[str, Any]:
return model.model_dump(mode="json", by_alias=True, exclude_unset=True)
...
@action(gives=Gives.RECORDS)
def side_changes(self, input: EditSideInput, ctx: RecordCtx) -> list[Record]:
"""The side fields that were sent."""
return [Record(kind="fixture.sides.set", data=_sent(input))]
packages/flows/src/omnium_flows/results/actions.py:25-26, 44-47
exclude_unset=Truekeeps only the fields the person sent. A field not sent is not touched.by_alias=Trueturnsparticipant_idintoparticipantId, the name the writer reads.- The action decides what changes. The writer does the database work.
The writer and the hand-edit rule
Writers live in omnium core, not in the plugin. Each is registered for one record kind:
@writer("fixture.sides.set")
async def set_side(ctx: WriteContext, data: dict[str, Any]) -> None:
"""Correct one side: who it is, its score, place, medal, status, record, qualified.
Each is pinned, so the next pull leaves what a person set alone.
"""
side = await _side(ctx, data.get("side"))
...
if "score" in data:
...
rest = {k: v for k, v in (side.result or {}).items() if k != "score"}
side.result = {**rest, "score": data["score"] or ""}
pin(side, "result")
if "medal" in data:
...
medal = str(data["medal"] or "").strip()
rest = {k: v for k, v in (side.result or {}).items() if k != "medal"}
side.result = {**rest, "medal": medal} if medal else rest
pin(side, "result")
...
if "rank" in data:
side.rank = data["rank"]
pin(side, "rank")
...
_touched(ctx)
packages/core/src/omnium_core/workflows/writers/fixtures.py:295-301, 312-325, 329-331, 341
_sideloads thefixture_competitorrow and refuses one from another fixture with aRecordError(writers/fixtures.py:65-74). ARecordErrorrolls back the log row too, and the sentence becomes the refusal.pin(side, "result")adds the field name to the row'spinnedlist (workflows/records.py:164-166). A later import will not overwrite a pinned field.may_write(records.py:173-180) is the same rule for writers that also serve imports._touchedcounts the update and putsfixtureIdin the answer (writers/fixtures.py:101-104).
How the engine runs them Built today
There are two engines today, and the bridge picks one by the workflow's mode. Both find components the same way: by key, through flows_host and the registry.

@router.post("/workflows/{workflow}/subjects/{kind}/{ref}/commands")
async def send_command(
workflow: str,
kind: Kind,
ref: str,
body: CommandIn,
session: WriteSessionDep,
caller: CallerDep,
event: str | None = None,
) -> CommandOut:
subject = await _resolve(session, workflow, kind, ref, event)
if subject.workflow.mode == "apply":
return await _apply(session, subject, body, caller)
return await _score(session, subject, body, caller)
packages/admin/src/omnium_admin/routers/bridge.py:213-226
From key to function
The engines never hold a Python function. They hold a program: a document generated from the registry and saved in scoring_program_version. Each node in it carries impl, the component key.
commands = tuple(
CommandNode(
code=f"{root}.{component.spec.name}",
...
validations=component.spec.validations,
actions=component.spec.actions,
branches=component.spec.branches,
stats=component.spec.stats,
applies_to=component.spec.applies_to,
impl=component.key,
impl_sha=component.spec.source_sha,
)
for component in of(ComponentKind.COMMAND)
)
packages/core/src/omnium_core/flows_publish.py:93-113
flows_publish.installbuilds this document and publishes it when the plugin's code changed (flows_publish.py:377-444). A fixture is pinned to the version it was scored with.impl_shais a hash of the function's own source. It does not cover helpers the function calls (contract/.../flows/components.py:73-76).- At run time,
flows_host.function_for(impl)looks the key up in the registry (flows_host.py:107-128). A key the installed plugin does not have is an error that names the mismatch.
FlowsFunction._invoke is the one place that calls your function. It picks the right arguments by kind and mode:
if kind is ComponentKind.VALIDATION:
payload, scope = args
verdict: Verdict = component.fn(
sport, _typed_input(component, dict(payload or {})), ctx_for(scope)
)
if verdict.ok:
return True
return {"ok": False, "message": verdict.message}
if kind is ComponentKind.ACTION:
if component.spec.gives is Gives.RECORDS:
payload, scope = args
records = component.fn(
sport, _typed_input(component, dict(payload or {})), _record_ctx(scope)
)
return [
(r if isinstance(r, Record) else Record.model_validate(r)).model_dump()
for r in records
]
return self._action(sport, args)
packages/core/src/omnium_core/flows_host.py:271-290
sportis one instance of the workflow class, made once per process (flows_host.py:118-123). Your method receives it asself.ctx_forbuilds aRecordCtxfor apply mode, or aCtxwith the format parsed throughformat_modelfor fold mode (flows_host.py:154-181, 268-269).- Any exception inside your function becomes a
SandboxErrorwith the key in front (flows_host.py:233-238).
Apply mode: workflows/engine.py
kind, subject_id = _subject(workflow, event, fixture)
program, node = await _command(session, workflow, code, actor)
...
if not flows_host.spec_for(node.impl).imports:
await guard_scoring_txn(session)
await session.execute(sa.select(sa.func.pg_advisory_xact_lock(ledger.lock_key(subject_id))))
stream_id = await streams.stream_for(session, kind, subject_id)
if idempotency_key:
prior = await _prior_seq(session, stream_id, actor, idempotency_key)
if prior is not None:
return CommandOutcome(accepted=True, code=code, duplicate=True, seq=prior)
problems = flows_host.command_input_problems(node.impl, payload)
if problems:
return _refusal(code, "; ".join(problems))
scope: dict[str, Any] = {
"event": subjects.event_facts(event),
"fixture": await subjects.fixture_facts(session, fixture) if fixture is not None else None,
"current": {},
"source": actor.kind,
"actor": actor.name,
}
rejection = _rules(program, node.validations, payload, scope)
if rejection is not None:
return _refusal(code, rejection)
savepoint = await session.begin_nested()
report = WriteReport()
try:
item = TimelineItem(
stream_id=stream_id,
fixture_id=fixture.id if fixture is not None else None,
source_code=actor.source_code,
event_type=code,
unit=UNIT,
seq=await streams.next_seq(session, kind, subject_id),
status=TimelineItemStatus.CONFIRMED,
occurred_at=datetime.now(UTC),
payload=_logged(node.impl, payload),
idempotency_key=idempotency_key,
actor=actor.name,
)
session.add(item)
await session.flush()
await _write(program, node, payload, scope, WriteContext(...))
except RecordError as exc:
await savepoint.rollback()
return _refusal(code, str(exc), report, reason=exc.reason)
...
packages/core/src/omnium_core/workflows/engine.py:107-166 (the WriteContext arguments are shortened)
And _write, which runs each action and hands each record to its writer:
limits = None if flows_host.spec_for(node.impl).imports else SYNC_LIMITS
for ref in node.actions:
action = next((a for a in program.actions if a.id == ref), None)
if action is None or not action.impl:
continue
records = flows_host.function_for(action.impl)(
payload, scope, determinism=Determinism(clock_nanos=0), limits=limits
).value
for record in records or []:
kind = str(record["kind"])
await writer_for(kind)(context, dict(record.get("data") or {}))
context.report.kinds.add(kind)
await context.session.flush()
packages/core/src/omnium_core/workflows/engine.py:306-318
The numbered trace of one apply-mode call:
- Mode and subject.
_subjectrefuses a scoring workflow, or the wrong subject kind (engine.py:229-240). - Command and sender.
_commandloads the published program and finds the node by code. A person may not send an import; an integration may only send imports (engine.py:243-264). - Time guard. For a person's command,
SET LOCALlimits the transaction toscoring_txn_timeout_ms, 10,000 ms by default (scoring/pipeline.py:824-842,packages/core/src/omnium_core/settings.py:290). - Lock. A Postgres advisory lock on the subject, held until the transaction ends. Its key is the first 8 bytes of the subject's id (
scoring/ledger.py:167-176). - Same key again? If this sender already sent this
keyon this stream, the answer is "duplicate" with the oldseq. Nothing runs. - Schema. pydantic checks the input against the command's model.
- Rules. In order; the first "no" is the answer. Nothing has been written yet.
- Log row. Inside a savepoint (a point inside the transaction that can be rolled back alone), one
timeline_itemrow withunit = "record". - Actions and writers. Each action's records go to their writers. A
RecordErrorrolls back the savepoint, log row included. - Outbox. One
domain_eventrow,command.accepted, with the areas it touched (engine.py:171-188). Medal commands also republish the medal document (engine.py:194-195). - Dirty mark. The event is marked changed, so the client feeds rebuild after commit (
engine.py:205-206). - Answer.
CommandOutcome, which the bridge turns intoCommandOut. A dry run ("Try") rolls the savepoint back here instead of keeping it (engine.py:220-225).
The commit itself happens after the route returns, in the request's session teardown (packages/admin/src/omnium_admin/deps.py:25-31).
Fold mode: scoring/pipeline.py
The fold pipeline takes the same rings, with a replay after the log row:
await session.execute(sa.select(sa.func.pg_advisory_xact_lock(ledger.lock_key(fixture.id))))
# 1 — resolve. A command a fixture's type does not offer is not a command.
node = _command_for(program, fixture, code)
...
# 2 — idempotency. A retry is a no-op, not a second ball.
if idempotency_key:
prior = await ledger.seen(session, fixture.id, source_code, idempotency_key)
...
# 3 — schema. A flows-backed command's schema IS its pydantic model.
if node.impl:
from omnium_core import flows_host # noqa: PLC0415
problems = flows_host.command_input_problems(node.impl, payload)
...
context.derived = await _current_derived(session, fixture)
...
# 4 — validations, in the order the command lists them: the first failure is
# the message the scorer reads.
rejection = _validate(compiled, node, payload, context, result)
if rejection is not None:
return result.reject("validate", rejection)
# 5 — the op. The only write of truth, and always an insert.
try:
item, fires = await _apply(
...
)
...
# 6+7 — the actions this command runs, in dependency order, with its
# branches free to add to that set once the state is known.
context = await refold(
...
)
packages/core/src/omnium_core/scoring/pipeline.py:869-958
Inside refold, the state action always runs; the row-making actions run only if the command named them (or a branch added them):
ordered = action_order(program)
folds = tuple(node for node in ordered if node.gives is Gives.STATE)
resume = await _resume_points(session, fixture, folds)
if fold_state(compiled, context, determinism, folds, result, resume=resume):
await projections.write_scoreboard(session, fixture.id, context.derived)
await _write_checkpoints(session, fixture, context)
...
# Routing, out loud, on the state the command just produced.
if only is not None and branches:
only = _route(compiled, branches, payload or {}, context, only, result)
stamp = await fixture_stamp(session, fixture)
for node in ordered:
if not writes_rows_now(node):
continue
if only is not None and node.id not in only:
continue
rows = rows_of(node, run_action(compiled, node, context, determinism))
remember_rows(context, node, rows)
written = await projections.write(
...
)
packages/core/src/omnium_core/scoring/pipeline.py:704-758
The differences from apply mode, in short:
| Apply mode | Fold mode | |
|---|---|---|
| Used by | results, games | football, cricket and the other sports |
| What a rule reads | RecordCtx: event, fixture and sides, sender | Ctx: fixture, log, saved scoreboard, masters |
| Log row | timeline_item, unit = "record" | timeline_item, unit = "event", plus timeline_actor rows |
| Actions get | The command's input | The whole effective log (or one event, if incremental) |
| Writes | Each record's writer | scoreboard, stat_value, fixture_competitor_rank |
| "Try" (dry run) | Yes | No: the bridge refuses it (bridge.py:266-267) |
| Also | Marks the event changed for feeds | Publishes the fixture document (pipeline.py:965-967) |
A full example, end to end Built today
Example · A console operator gives a side its gold medal and rank 1
The operator opens a finished fixture in the console, picks a side, sets the medal to Gold and the rank to 1, and saves. The ids below are examples, in the right shape.
1. The console sends the command. editSide calls send("results.edit_side", { side: sideId, ...body }, fixtureRef(id)) (omnium-console/src/lib/api.ts:803-806). The bridge app adds a fresh key from crypto.randomUUID() (packages-ts/omnium-bridge/src/app.ts:55-59, 323-334). The host page then calls omnium:
POST /workflows/results/subjects/fixture/3f1c2a9e-5b7d-4c11-9e2a-7a0d8c4b6e21/commands?event=ag-2026
Authorization: Bearer <token>
Content-Type: application/json
{
"code": "results.edit_side",
"input": {
"side": "8d2e4f60-1a3b-4c5d-8e9f-0a1b2c3d4e5f",
"medal": "Gold",
"rank": 1
},
"key": "c0a8f1d2-6b3e-4f7a-9c2d-1e5f8a7b6c4d"
}
The URL shape is built in packages-ts/omnium-bridge/src/http.ts:181-186; the body shape is the one its own test checks (http.test.ts:41-58).
2. The bridge route. send_command resolves the workflow results and the fixture, checks that the fixture belongs to event ag-2026, sees mode == "apply" and calls engine.run (bridge.py:213-259).
3. The engine finds the command. The pinned program for results has a command node results.edit_side with impl = "results.edit_side". The sender is a person and the command is not an import, so it may send it.
4. Lock and key. It locks the fixture and looks for this key from this account on the fixture's stream. None, so it goes on.
5. Schema. EditSideInput.model_validate(...) passes: side is a UUID, rank is 1 or more, medal is under 32 characters.
6. Rules, in order.
| Rule | Reads | Answer |
|---|---|---|
side_is_in_fixture | input["side"], ctx.fixture["sides"] | ok: the side is in this fixture |
side_keeps_someone | which fields were sent | ok: neither participantId nor label was sent |
side_status_is_known | input.status | ok: not sent |
side_medal_is_known | input.medal | ok: "gold" is in SIDE_MEDALS |
7. Log row. One timeline_item row is added:
| Column | Value |
|---|---|
stream_id | the fixture's ledger_stream row |
fixture_id | 3f1c2a9e-... |
event_type | results.edit_side |
unit | record |
seq | the stream's next number, say 12 |
status | CONFIRMED |
payload | the input exactly as sent |
idempotency_key | c0a8f1d2-... |
actor, source_code | the account's name and code |
8. Action. side_changes returns one record:
[Record(kind="fixture.sides.set", data={"side": "8d2e4f60-...", "medal": "Gold", "rank": 1})]
9. Writer. set_side loads the fixture_competitor row, sets result["medal"] = "Gold", sets rank = 1, and adds result and rank to pinned. A later Scout import will not overwrite either.
10. Outbox and dirty mark. One domain_event row: kind command.accepted, payload with command: "results.edit_side", workflow: "results" and touches: ["fixture"]. The event is marked changed so the feeds rebuild after commit.
11. The answer. HTTP 200:
{
"accepted": true,
"duplicate": false,
"seq": 12,
"message": null,
"refusal": null,
"reason": null,
"touches": ["fixture"],
"value": { "fixtureId": "3f1c2a9e-5b7d-4c11-9e2a-7a0d8c4b6e21" },
"counts": { "updated": 1 },
"notes": [],
"dryRun": false,
"fixtureStatus": "COMPLETED"
}
Field names are camelCase because CommandOut uses AdminModel, whose alias generator is to_camel (packages/admin/src/omnium_admin/schemas.py:11-12). The host tells other open screens that fixture changed (packages-ts/omnium-bridge/src/host.ts:319), and the console reads the fixture again (api.ts:805).
The same request, three other ways:
| What is different | Where it stops | The answer |
|---|---|---|
"medal": "Golden" | Step 6, side_medal_is_known | 200, accepted: false, refusal: "rejected", message: "A medal is Gold, Silver or Bronze" |
"rank": 0 | Step 5, schema | 200, accepted: false, message: "rank: Input should be greater than or equal to 1" (pydantic's wording) |
Same key sent again | Step 4 | 200, accepted: true, duplicate: true, seq: 12; nothing written |
"code": "results.set_line" | Step 3 | HTTP 404, detail.code: "not_offered"; this command does not exist yet |
For a football goal the road is the same up to step 6, then it takes the fold path: the rule half_is_open reads halfOpen from the saved scoreboard; one timeline_item row with unit = "event" (scoring/ledger.py:36) is written, with a timeline_actor row for each named scorer and assist; state replays the log and rewrites scoreboard; tallies returns lines that replace the match's rows in stat_value. The answer's value is the desk's fresh board (bridge.py:286-295).
Stats in the same package Built today
A sport's stats are declared as data in its stats.py, and they compile to SQL inside the plugin when the sport is imported. The engine runs what it is handed and never knows any sport's numbers.
SCORING = {"goals": "goals", "assists": "assists", "own_goals": "own_goals"}
class FootballTotals:
"""Mixed into ``Football``."""
goals_node = totals(
key="football.scoring.node",
scope="node",
name="Goals and assists",
description="A player's goals, assists and own goals in this competition.",
numbers=SCORING,
appearances="appearances",
)
...
table = standings(
key="football.table.node",
name="League table",
description="Played, won, drawn, lost, goals and points for this competition.",
scored="goals_for",
conceded="goals_against",
)
packages/flows/src/omnium_flows/football/stats.py:10-23, 33-39
totals,rankingandstandingsare shared stat cards inomnium_flows/shared/stat_cards.py. Each builds a spec and callsdeclared_statfromomnium_contract.stats.declared_statcompiles the spec into a stat function at import time, with the compiler in the contract package (packages/contract/src/omnium_contract/stats/spec.py:142-184). A spec that does not hold together fails the import, while it is still a draft.- The numbers they read (
goals,goals_for,goals_against) are the keys thetalliesaction writes tostat_value. That is the link between actions and stats. FootballTotalsis mixed intoFootballlike the other four classes. Results and Games have astats.pytoo, with an empty class, because the structure gate asks for the file.- Stats run in the stats queue, after the command commits, never inside it (
scoring/pipeline.py:993-998).
How to add a new command
Add the input model, the command, its rules and its action, then one test that names each new function. The structure gate fails the build without the test.
- Input model in
commands.py. Subclass the file's_Input. Give every field a type and a limit (max_length,ge). - Command in
commands.py. An empty method with@command(input=..., validations=(...), actions=(...)). The method name becomes the key's last part, and it is permanent once used: it is written in every log row. - Validations in
validations.py. One@validation(message="...")per rule. Returnok()orreject("..."). Write the sentence for the person who will read it. List them in the command in the order the person should hear them. - Action in
actions.py.- Fold sport: return
Lines orPlaces. Put every new number name inkeysand in the sport'sfactsinsport.py. - Results or Games: return
Record(kind=..., data=...)for a writer that exists. A new record kind needs a new@writerinpackages/core/src/omnium_core/workflows/writers/, imported byworkflows/catalogue.py. That is a core change, not a plugin change.
- Fold sport: return
- Who may send it. Set
imports=Trueonly for a command an integration sends; then addfieldsfor the integration screens. - Test in
packages/flows/tests/<workflow>/test_<workflow>.py. Call the functions directly, as plain methods:
def test_two_sides_refuses_a_team_playing_itself() -> None:
football = Football()
same = MatchStartInput.model_validate({"home": "united", "away": "united"})
assert not football.two_sides(same, ctx()).ok
different = MatchStartInput.model_validate({"home": "united", "away": "city"})
assert football.two_sides(different, ctx()).ok
packages/flows/tests/football/test_football.py:201-206
Football's tests also assert the full list of registered keys (test_football.py:41-61), so a new component must be added there too.
- Golden match, for a fold sport whose numbers change. A golden is a recorded log plus the expected output of each action, in
packages/flows/tests/goldens/<sport>/*.json.test_goldens.pyreplays every golden on every run and fails on any difference (tests/test_goldens.py:45-113). Football has four:two_one_home_win,own_goal_draw,abandoned,empty_pitch. - Run
make check. Itsstructurestep runs two scripts:
REQUIRED_MODULES = (
"constants.py",
"commands.py",
"validations.py",
"actions.py",
"stats.py",
"sport.py",
)
...
for component in registry.components():
name = component.spec.name
if component.key.rsplit(".", 1)[-1] != name:
problems.append(f"{component.key}: key does not end with member name {name!r}")
if component.inherited_from is not None:
continue # the parent's tests cover inherited components
if name not in test_text:
problems.append(f"{component.key}: no test names it")
packages/flows/scripts/check_structure.py:26-33, 54-61
- The first part fails a workflow folder without the six files, or without a
tests/<folder>/directory. - The second fails any component whose name appears in no test file.
check_clean_import.pyimports each entry-point module in a fresh process with network sockets blocked, and fails on any output or error (packages/flows/scripts/check_clean_import.py:18-42).- The plugin's own ruff settings ban clocks,
random,os.environand file opens in plugin code (packages/flows/pyproject.toml:63-73), because components must be pure.
- Publish.
flows_publish.install_allregenerates each workflow's program from the registry and publishes the ones that changed (flows_publish.py:447-457). It does not run on its own at start-up. It runs aftermake seed, frompython -m omnium_worker install-programs, or from the admin routePOST /workflows/install(packages/worker/src/omnium_worker/__main__.py:94-96, 407-410;packages/admin/src/omnium_admin/routers/workflows.py:408-422). The plugin itself ships as thefanos-omnium-flowspackage (make publish-flows).
How it changes in the agreed design Agreed, to build
Sections 4, 5 and 6 (all closed 6 Oct 2026) keep the decorators and the six files, and change the engine around them. None of this is in the code yet.
| Today | Agreed | Decision |
|---|---|---|
Two engines: workflows/engine.py and scoring/pipeline.py | One package, omnium_core/engine/, with two modes inside: scorecard (apply once) and events (step) | Section 5, module layout; C9 |
| A sport's state action replays the whole log on every command | An events sport provides a step function: initial_state, check, step, summary. The saved state plus one event gives the new state. check replaces rules that read the whole log | S1, S2, S3, M3 |
| Football and cricket fold the whole log; only athletics steps | Cricket and football must get a step function before they go live; state under 16 KB | S2, M2 |
| No command can set one player's line by hand | A new Results command, results.set_line, with a new writer writers/lines.py. Keys are checked against the sport's catalogue and pinned | S12 |
key is optional; refusals are not saved; two answers | A command id is required; five answers (accepted, duplicate, refused, conflict, try again); each refusal saved, with a fixed code and a sentence | C2, C5, C7, L18 |
| Last save wins | Edits carry the version the sender saw; "conflict" if it moved | C3, L4, L17 |
| Components run on the web server's thread | Components run in small worker processes per sport and release | L10, M6 |
make new-sport makes an old file shape | The scaffold makes the six files, a test and a golden match | P11 |
| A scorecard-only sport still needs a plugin | A sport scored as a scorecard needs no code: settings plus the Results workflow | P1 |
Read the details on the pages that own them:
- Commands and the workflow engine: the one road, answers, keys and locks.
- The live scoring engine: scorecard and events modes, the step function,
results.set_line. - Sport plugins and rule versions: releases, pins, and the fixed scaffold.
Read next
- Commands and the workflow engine: the road every command takes, and what changes in it.
- The live scoring engine: how a match state is kept, corrected and checked.
- Sport plugins and rule versions: how a plugin release reaches a match safely.
- Truth and ownership: the pin rule that writers apply.