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.

Design section
Built today; sections 4 to 6 (agreed 6 Oct 2026) for what changes
Main code
packages/flows, contract/.../flows, core/.../workflows/engine.py, core/.../scoring/pipeline.py
Main tables
timeline_item, scoreboard, stat_value, fixture_competitor, fixture_competitor_rank, domain_event
Read time
about 25 minutes

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.

6
files in every sport folder (structure gate)
513
lines: the whole football sport
50 ms
time budget for one validation (from the code)
200 ms
time budget for one action (from the code)

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:

  1. Find the command by its code.
  2. Check the input against the command's input model.
  3. Run its validations, in order. The first "no" is the answer.
  4. Write the log row.
  5. Run its actions, and write the rows they return.
  6. Answer the sender.
A sketch titled One command, three kinds of code. A dashed frame labelled one sport folder holds six file shapes: constants.py, commands.py, validations.py, actions.py, stats.py, sport.py. Below it a row of boxes runs left to right: Console or scorer app, an arrow labelled command to Command (input model), an arrow labelled in order to Validations (first no wins), an arrow labelled all said yes to Actions (return rows), an arrow labelled engine writes to a Postgres cylinder with log row and changed rows under it. Dotted lines link commands.py to Command, validations.py to Validations and actions.py to Actions. A red dashed arrow goes down from Validations to a red box Refused, the rule's sentence. A green arrow labelled answer goes from Postgres back to the console.
A command passes its validations, then its actions return rows, and the engine writes them. Each kind of code has its own file.

The words you will meet on this page:

WordIn one sentence
ComponentOne decorated function: a command, validation, action, branch or context.
KeyA component's permanent name: workflow code, a dot, the function name (football.goal).
RegistryThe list of every component the installed plugin has, built at start-up.
Ledger, or logThe timeline_item table: one row per accepted command, never edited.
Apply modeEach command's actions run once, on its own input. Used by Results and Games.
Fold modeEach command replays the whole match log to rebuild the state. Used by sport scoring.
Derived stateThe 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.

FileWhat it holdsExample from football
constants.pyWords the sport uses, and its format modelGOAL_KINDS, CARD_COLOURS, FootballFormat (halves)
commands.pyInput models, and one @command per thing a person can sendmatch_start, period_start, goal, card, period_end, match_end, abandon
validations.pyOne @validation per rule, each with its sentencenot_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 placesstate, tallies, placings
stats.pyThe sport's stat cards, declared as datafootball.scoring.node, football.top_scorers.node, football.table.node, football.scoring.career
sport.pyThe class that joins the other five, its code, its fact list, its contextclass 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_model is the shape of the fixture's format. Football reads only halves.
  • roles lists the roles a command may name a person in. A goal names a scorer and an assist.
  • Sport is 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 code plus 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_registry imports each module named in the group omnium.flows and 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 the OMNIUM_FLOWS_VERSION pin 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_camel means the app may send minute or a camelCase name, and both work.
  • kind: GoalKind is a fixed list of words from constants.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

ArgumentWhat it does
input=GoalInputThe 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.
descriptionThe 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' _Input adds extra="forbid" (results/commands.py:33-34). A field the model does not know is refused, not dropped. Football's _Input does not set it.
  • Every field but side is 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=True marks 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_FIELDS lists the columns a no-code Scout integration maps, with labels and hints for the screen (games/constants.py:97 onwards).

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: Any means 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.derived is the saved scoreboard as it was before this command (scoring/pipeline.py:1179-1188). That is why half_is_open can 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_fixture is shared by edit_side and remove_side, so it reads the raw input.
  • side_medal_is_known returns Verdict(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.impl is the component key, like results.side_medal_is_known. flows_host.function_for turns it into a callable.
  • flows_host turns a Verdict into True, 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

givesReturnsWritten toUsed by
STATEOne dictionary: the match nowscoreboard.snapshotSport scoring (fold)
LINESA list of Line: each subject's numbersstat_valueSport scoring (fold)
PLACESA list of Place: the finishing orderfixture_competitor_rank, and rank on the linesSport scoring (fold)
RECORDSA list of Record(kind, data)Whatever the writer for kind writesResults, Games (apply)
NOTHINGNothing; side work onlyNothing; it may only run laterNot 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_open stops 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 as ownGoals, never goals.
  • 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 the state action just returned, not the log.
  • keys lists every number this action may write. Each must be in the sport's facts list in sport.py. A typo is refused when the program is published (contract/.../flows/sport.py:80-88).
  • A Line is 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_end and abandon list 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=True keeps only the fields the person sent. A field not sent is not touched.
  • by_alias=True turns participant_id into participantId, 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

  • _side loads the fixture_competitor row and refuses one from another fixture with a RecordError (writers/fixtures.py:65-74). A RecordError rolls back the log row too, and the sentence becomes the refusal.
  • pin(side, "result") adds the field name to the row's pinned list (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.
  • _touched counts the update and puts fixtureId in 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.

A sketch titled How the engine finds your code. An HTTP request box (code + input) has an arrow to Bridge router. From Bridge router, an arrow labelled apply mode goes to Workflow engine (workflows/engine.py), and an arrow labelled fold mode goes to Scoring pipeline (scoring/pipeline.py). Both point, with arrows labelled component key, to flows_host. An arrow labelled look up goes from flows_host to Registry, and an arrow labelled entry points goes from Registry to the omnium_flows plugin, which lists football, results and games.
The bridge picks the engine by mode. Both engines ask flows_host for a component by key; the registry found it through entry points.
@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.install builds 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_sha is 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

  • sport is one instance of the workflow class, made once per process (flows_host.py:118-123). Your method receives it as self.
  • ctx_for builds a RecordCtx for apply mode, or a Ctx with the format parsed through format_model for fold mode (flows_host.py:154-181, 268-269).
  • Any exception inside your function becomes a SandboxError with 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:

  1. Mode and subject. _subject refuses a scoring workflow, or the wrong subject kind (engine.py:229-240).
  2. Command and sender. _command loads 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).
  3. Time guard. For a person's command, SET LOCAL limits the transaction to scoring_txn_timeout_ms, 10,000 ms by default (scoring/pipeline.py:824-842, packages/core/src/omnium_core/settings.py:290).
  4. 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).
  5. Same key again? If this sender already sent this key on this stream, the answer is "duplicate" with the old seq. Nothing runs.
  6. Schema. pydantic checks the input against the command's model.
  7. Rules. In order; the first "no" is the answer. Nothing has been written yet.
  8. Log row. Inside a savepoint (a point inside the transaction that can be rolled back alone), one timeline_item row with unit = "record".
  9. Actions and writers. Each action's records go to their writers. A RecordError rolls back the savepoint, log row included.
  10. Outbox. One domain_event row, command.accepted, with the areas it touched (engine.py:171-188). Medal commands also republish the medal document (engine.py:194-195).
  11. Dirty mark. The event is marked changed, so the client feeds rebuild after commit (engine.py:205-206).
  12. Answer. CommandOutcome, which the bridge turns into CommandOut. 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 modeFold mode
Used byresults, gamesfootball, cricket and the other sports
What a rule readsRecordCtx: event, fixture and sides, senderCtx: fixture, log, saved scoreboard, masters
Log rowtimeline_item, unit = "record"timeline_item, unit = "event", plus timeline_actor rows
Actions getThe command's inputThe whole effective log (or one event, if incremental)
WritesEach record's writerscoreboard, stat_value, fixture_competitor_rank
"Try" (dry run)YesNo: the bridge refuses it (bridge.py:266-267)
AlsoMarks the event changed for feedsPublishes 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.

RuleReadsAnswer
side_is_in_fixtureinput["side"], ctx.fixture["sides"]ok: the side is in this fixture
side_keeps_someonewhich fields were sentok: neither participantId nor label was sent
side_status_is_knowninput.statusok: not sent
side_medal_is_knowninput.medalok: "gold" is in SIDE_MEDALS

7. Log row. One timeline_item row is added:

ColumnValue
stream_idthe fixture's ledger_stream row
fixture_id3f1c2a9e-...
event_typeresults.edit_side
unitrecord
seqthe stream's next number, say 12
statusCONFIRMED
payloadthe input exactly as sent
idempotency_keyc0a8f1d2-...
actor, source_codethe 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 differentWhere it stopsThe answer
"medal": "Golden"Step 6, side_medal_is_known200, accepted: false, refusal: "rejected", message: "A medal is Gold, Silver or Bronze"
"rank": 0Step 5, schema200, accepted: false, message: "rank: Input should be greater than or equal to 1" (pydantic's wording)
Same key sent againStep 4200, accepted: true, duplicate: true, seq: 12; nothing written
"code": "results.set_line"Step 3HTTP 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, ranking and standings are shared stat cards in omnium_flows/shared/stat_cards.py. Each builds a spec and calls declared_stat from omnium_contract.stats.
  • declared_stat compiles 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 the tallies action writes to stat_value. That is the link between actions and stats.
  • FootballTotals is mixed into Football like the other four classes. Results and Games have a stats.py too, 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.

  1. Input model in commands.py. Subclass the file's _Input. Give every field a type and a limit (max_length, ge).
  2. 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.
  3. Validations in validations.py. One @validation(message="...") per rule. Return ok() or reject("..."). Write the sentence for the person who will read it. List them in the command in the order the person should hear them.
  4. Action in actions.py.
    • Fold sport: return Lines or Places. Put every new number name in keys and in the sport's facts in sport.py.
    • Results or Games: return Record(kind=..., data=...) for a writer that exists. A new record kind needs a new @writer in packages/core/src/omnium_core/workflows/writers/, imported by workflows/catalogue.py. That is a core change, not a plugin change.
  5. Who may send it. Set imports=True only for a command an integration sends; then add fields for the integration screens.
  6. 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.

  1. 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.py replays 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.
  2. Run make check. Its structure step 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.py imports 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.environ and file opens in plugin code (packages/flows/pyproject.toml:63-73), because components must be pure.
  1. Publish. flows_publish.install_all regenerates 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 after make seed, from python -m omnium_worker install-programs, or from the admin route POST /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 the fanos-omnium-flows package (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.

TodayAgreedDecision
Two engines: workflows/engine.py and scoring/pipeline.pyOne 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 commandAn 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 logS1, S2, S3, M3
Football and cricket fold the whole log; only athletics stepsCricket and football must get a step function before they go live; state under 16 KBS2, M2
No command can set one player's line by handA new Results command, results.set_line, with a new writer writers/lines.py. Keys are checked against the sport's catalogue and pinnedS12
key is optional; refusals are not saved; two answersA command id is required; five answers (accepted, duplicate, refused, conflict, try again); each refusal saved, with a fixed code and a sentenceC2, C5, C7, L18
Last save winsEdits carry the version the sender saw; "conflict" if it movedC3, L4, L17
Components run on the web server's threadComponents run in small worker processes per sport and releaseL10, M6
make new-sport makes an old file shapeThe scaffold makes the six files, a test and a golden matchP11
A scorecard-only sport still needs a pluginA sport scored as a scorecard needs no code: settings plus the Results workflowP1

Read the details on the pages that own them: