Getting data in · Commands

Commands and the workflow engine

The one road every change takes: a key, a lock, a check, one transaction, then an answer you can trust.

Design section
Section 4, agreed 6 Oct 2026 (30 decisions)
Main code
core/workflows/engine.py, core/scoring/pipeline.py, admin/routers/bridge.py
Main tables
fixture, timeline_item, domain_event; new: command, command_failure
Read time
about 35 minutes

In one minute

Every change in omnium is a command: a scorer's tap, a staff fix in the console, a provider message or a batch from a Scout import. A command has a code (what to do), an input (the values) and a key (an id the sender makes once and sends again on every retry).

Every command takes the same road. Postgres locks the match's own row. Then the engine checks the key, the version the sender saw, the input and the rules, and works out the change. It saves everything in one transaction and only then sends the answer. There are five answers: accepted, duplicate, refused, conflict and try again. Each one tells the screen exactly what to do next.

The road exists today and its order is mostly right, but it has real gaps. The answer can go out before the commit. A retry can count twice. There is no version check. Two kinds of lock never meet. The outbox reader can skip a late commit. The agreed design closes each gap with plain Postgres: no queue, no Redis lock, no new tools.

It is fast. On a laptop with 10 million command rows, the whole command transaction took 1.48 ms at the median and 3.98 ms at the 99th percentile. Eight senders all on one match still got 904 commands a second.

0.09 ms
key lookup, 10M rows (laptop, 5 Oct)
1.48 ms
whole command, median (laptop, 5 Oct)
3.98 ms
whole command, 99th pct (laptop, 5 Oct)
12.89 ms
8 senders on one match, 99th pct

What this part does

The command road answers six questions. If any one is answered wrong, data is lost, counted twice, or saved without anyone knowing.

#QuestionWhy it mattersAnswer in short
1Does every change take the same road?A side door skips the checks, the log and the priority rulesYes, no side doors (C1)
2How is a retry made safe?Networks drop. A second tap must not make a second goalThe same key gives the same answer (C2)
3What if two people change the same thing at once?One change must not wipe out the other unseenA lock for order, a version check for "you saw an old one" (C3, C6)
4When may the server say "saved"?If it says "saved" and then crashes, the scorer trusts data that is goneOnly after the commit returns (C4)
5What do we keep when a command is refused?Support must answer "why did my change not go through?"Every refusal is saved, with the reason (C5)
6One engine, or two?Two engines means two sets of bugsOne engine, two ways to work out a change (C9)

What a command is

A command is one request to change something, sent over HTTP to the bridge. Its address today is POST /workflows/{workflow}/subjects/{kind}/{ref}/commands, and its body is CommandIn (packages/admin/src/omnium_admin/routers/bridge.py:183).

PartExampleWhat it is
coderesults.edit_sideWhich command of the workflow to run. The plugin defines it.
inputside, scoreThe values. The command's own model checks their shape.
key7f3a2c1e-…The retry key. The sender makes it once and sends it again on every retry.
dry_runfalse"Try": run everything, then undo it.
sourcedesk-2Which desk sent it, when a match has several.
saw (new)changed_seq per row, or the last play eventThe version the sender saw. Agreed, not built.
sentAt (new)2026-10-06T10:00:00.120ZWhen the device made it. Agreed, not built.

The subject is the thing the command changes: one match (a fixture) or one event (a competition, for medal tables, schedules and players). How a command, its validations and its actions are written in the flows plugin is on Writing a workflow in code. This page is the engine side: what happens to a command once it arrives.

How it works

A sketch of one command's road. On the left, a white box for the scorer, console or import, holding key 7f3a, sends POST command into a middle column of seven steps going down: 1 quick key check, 2 lock the match row with FOR NO KEY UPDATE, 3 key check again, 4 version check, 5 rules then apply, 6 commit in the database with the log, outbox and command row, and 7 the answer accepted in green. Four side exits lead right: from step 3 to a green box duplicate with the first answer; from step 2, after waiting 2 seconds, to a red box try again; from step 4 to a red box conflict, you saw an old version; and from step 5 to a red box refused, saved with the reason.
The agreed road. A command either reaches step 7 and is saved, or leaves by a side exit that tells the sender what to do.

The road, step by step, as agreed in Section 4. The numbers match the drawing. Steps 1 to 5 can each end the command early, by a side exit. Only steps 6 and 7 make it real.

Before the road: the sender keeps the command, with its key. A screen keeps the command on the device until an answer comes. It makes the key once, and every retry sends the same key. A provider message takes its key from the provider's own message number. An import batch uses the run's id plus the batch number.

  1. Quick key check, before any lock. One index lookup in the command table. A retry of a finished command is answered "duplicate" at once. This is only a shortcut; step 3 is the check that counts. Input of the wrong shape is also refused here, before any lock.
  2. Lock the match row. SELECT last_seq FROM fixture WHERE id = $match FOR NO KEY UPDATE. One command at a time per match, for a few milliseconds. The same read returns the match's log number. A command that waits more than 2 s gets try again.
  3. Key check again, under the lock. Seen with the same body: duplicate, with the first answer. Seen with a different body: refused (key_reused).
  4. Version check. An edit carries the changed_seq of each row it changes. A ball carries the last play event the sender saw. If anything moved on: conflict, with what changed. A pure add (a commentary line) skips this.
  5. Rules, then apply. The command's validations run. A "no" from a rule is saved as a refusal, and nothing else is written. Then the change is worked out: "apply once" for results, schedules and imports, "fold" for live scoring (Scoring engine).
  6. Commit it all in one transaction. The log row (timeline_item), the claims and the winner from Section 2, the note to feeds (domain_event), the new fixture.last_seq, and the command row with the exact answer.
  7. Answer. "Accepted" goes out only after COMMIT returns.

After the road: feeds and screens read the note, by transaction id, so nothing is skipped (see "The outbox, and the late commit" below).

The five answers

Every answer has the same body, with an outcome field saying which of the five it is.

OutcomeHTTPWhat it meansWhat the sender does
accepted200Saved on diskRemoves it from its unsent list
duplicate200Already saved; here is the first answer, exactlyThe same as accepted
refused422A rule said no. The body has a code, a sentence with the real values, the field, the value now, and what to doShows the sentence. Does not send it again
conflict409You saw an old version. The body has the rows as they are now, or the events you missedUpdates the screen from the answer, asks the person
try_again503 + Retry-AfterBusy, a limit was hit, or the database is switching over. Cause: lock_timeout, worker_timeout, code_missing or db_unavailableSends again with the same key after a short wait

Locks, explained from the start

A lock is a "one at a time" rule. While one command holds the lock on match 42, a second command for match 42 waits its turn. Commands for other matches do not wait.

The agreed lock is a row lock that Postgres keeps on the match's own row in the fixture table. Our code never writes a "locked" flag anywhere. Postgres marks the row with the id of the transaction that holds the lock, and the mark counts only while that transaction is running. When the transaction ends, by commit, by rollback, or because the connection died, the lock is gone at once. Nothing is ever left behind to clean up.

Postgres has four strengths of row lock. We use FOR NO KEY UPDATE, the second strongest. It is the lock a normal UPDATE of non-key columns takes anyway, so we take it a little earlier, at the start of the command.

Someone else tries to…While we hold FOR NO KEY UPDATE on match 42
Read match 42 with a plain SELECT (feeds, public API, screens)Not blocked. Readers see the last saved version.
Insert a row that points at match 42, like a timeline_item or fixture_competitor (the foreign key check takes FOR KEY SHARE)Not blocked. This is why we use NO KEY UPDATE and not UPDATE.
Take FOR NO KEY UPDATE on match 42 (another command)Waits until we commit or roll back.
UPDATE a normal column of match 42Waits.
FOR UPDATE, FOR SHARE, DELETE of match 42, or a change to its idWaits.
Anything on match 43Not blocked. Each match has its own lock.

Why sorted order stops deadlocks. A deadlock is two commands each waiting for a lock the other holds, in a circle, forever. It can only happen when two commands take the same locks in a different order:

TimeCommand XCommand Y
t0locks match 7locks match 3
t1wants match 3: waits for Ywants match 7: waits for X
t2waits foreverwaits forever

If everyone locks in the same fixed order, the circle cannot form. Both X and Y try match 3 first. One gets it. The other waits at the very first lock, holding nothing, so the first one always finishes. The agreed order (L2) is: the event first, then matches, then people and teams, each kind sorted by id, all in one statement per kind. If a deadlock happens anyway, Postgres finds it within about 1 second (its default deadlock_timeout), cancels one command, and that command gets "try again".

What the lock does not decide. The lock decides the order: one command at a time per match. It does not decide what clients see when sources disagree. That is the Section 2 priority list (Truth and ownership). Every source's value is saved, and the winner of each field group is shown.

A timeline sketch with three lanes, Scorer A at the top, the Match 42 row in the middle, and Scorer B at the bottom, and a time axis marked 0, 1, 3 and 5 milliseconds. Scorer A sends a goal saying it saw event 41, takes the lock, and later gets accepted as event 42. The match row is locked by A from 0 to 3 milliseconds, then by B from 3 to 5. Scorer B sends a goal that also saw 41, waits alongside A's lock until 3 milliseconds, then gets a red conflict saying it missed event 42, and its screen shows A's goal.
Two writers on one match. The lock makes B wait; the version check stops B from saving on top of a state it never saw.

A worked example

Example · Two scorers send a goal at the same moment

Football, match 42, 63rd minute. Scorer A and scorer B both cover the match (Section 7 lets a match allow two scorers or one with a hand-over). The last play event both screens show is 41. Both see the same striker score and both tap "goal" within 1 ms.

TimeScorer A, key 7f3a…, saw 41Scorer B, key c19e…, saw 41Database
0 msQuick key check: not seen. Takes the lock on match 42fixture.last_seq = 41
1 msKey check under the lock: not seen. Play version: saw 41, now 41, fine. Rules passQuick key check: not seen. Asks for the lock: waits in Postgres's queue
2 msWrites timeline_item seq 42, domain_event, last_seq = 42, and the command rowstill waiting
3 msCOMMIT. Answer: accepted, event 42Gets the lock. Key c19e not seen. Play version: saw 41, now 42. ROLLBACK. Answer: conflict, with event 42 insidelast_seq = 42; nothing of B written
5 msScreen adds A's goal from the answer and says "1 new event from A: check, then tap again"

B looks: it is the same goal. B does nothing. The match shows 1–0, one goal. If B had seen a second goal, B taps again with "saw 42", and it is saved as event 43.

The times are an illustration. The laptop measured 1.2 to 1.5 ms median for the lock to be held; prod is expected at 3 to 8 ms (estimate in the design doc).

Today: both commands take the advisory lock (pipeline.py:869), so they run one after the other. But nothing carries a version, so both goals are saved: the log holds two goals. Whether a sport's own rule catches a double goal is up to that sport (not checked for football). Live scoring is not on prod today.

Example · A phone resends after losing network

Scorer A taps "goal" in a stadium tunnel. The command gets key 7f3a, and the bridge outbox saves it on the phone (IndexedDB) before the first send.

  1. The request reaches the server. The command is saved as seq 57 and committed. The answer is lost in the tunnel.
  2. The screen gets no answer in 5 s. The goal stays on screen, marked "sending". The phone waits about 0.2 s, plus a little randomness.
  3. It sends the same command again, with the same key 7f3a.
  4. Quick key check: the command row for 7f3a is there, and the body hash matches. Answer: duplicate, with the first answer (seq 57). This took about 0.44 ms on the laptop.
  5. The screen treats duplicate like accepted and removes it from the unsent list. One goal, not two.

Two other ways this can go, both safe:

  • The first copy is still running (a slow moment). The second copy passes the quick check, then waits for the match lock. When it gets the lock, the first copy has committed, so the check under the lock finds row 7f3a: duplicate.
  • The server crashed before the commit. Postgres rolled everything back, including the command row. The retry finds no row and runs normally: accepted. The goal is saved once.

Today: the bridge makes a new key on every call unless the caller passes one (packages-ts/omnium-bridge/src/app.ts:332), and the console never passes one (omnium-console src/lib/bridge.ts:122). Nothing resends by itself (the bridge has no retry loop, and the console's mutations use retry: 0, omnium-console src/main.tsx:33). But if the person taps again after a timeout, that is a new key, so it is a second command. Also, the answer can go out before the commit (see "Built today" below).

Low-level design

Today's road, in the code Built today

Today there are two engines. Games and Results commands (apply mode) go through engine.run. Live scoring desks go through pipeline.command. Both are called from send_command in the bridge router.

# packages/core/src/omnium_core/workflows/engine.py:107-126 (today, trimmed)
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)            # SET LOCAL timeouts, a person's command only
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))

In plain words: it finds the subject and the command, sets the time limits for a person's command, then takes a lock on the subject. The key check runs under the lock, which is right. But the lock is an advisory lock on the subject id: for a Results edit the match, for a Games command or an import the whole event. The key is looked up in timeline_item, so a refused command leaves no key behind. Nothing compares the body, so the same key with a different body is answered "duplicate".

# packages/core/src/omnium_core/scoring/ledger.py:167-177 (today)
def lock_key(ref_id: uuid.UUID) -> int:
    """A signed 64-bit advisory-lock key for anything identified by a uuid."""
    return int.from_bytes(ref_id.bytes[:8], "big", signed=True)

In plain words: the lock number is the first 8 bytes of the id. Our ids are uuid7 (uuid_utils, packages/contract/src/omnium_contract/models/base.py:13), which start with the time they were made. Ids made in the same moment share their first 8 bytes: five made in a row on 7 Oct all began 01a1165c-f545-7e91. So matches made together get the same number and share a lock. The two speed-climbing matches fixed on 1 Oct (01a0a101-2736-79a2-…) share one. No data is lost, but an edit on one waits for the other. L1 replaces this with the row itself.

# packages/core/src/omnium_core/workflows/streams.py:60-70 (today, trimmed)
async def next_seq(session: AsyncSession, kind: str, subject_id: uuid.UUID) -> int:
    """The next position in this subject's log. One row lock; no scan, no race."""
    if kind == FIXTURE:
        seq = (await session.execute(
            sa.update(Fixture).where(Fixture.id == subject_id)
            .values(last_seq=Fixture.last_seq + 1).returning(Fixture.last_seq)
        )).scalar_one_or_none()

In plain words: the match row already counts its log in fixture.last_seq, and an event counts in ledger_stream.last_seq. This UPDATE already takes a row lock on the match, but only late, after the rules have run. The agreed design takes the same row lock at the start instead, and drops the advisory lock. No new lock table is needed.

# packages/admin/src/omnium_admin/deps.py:12-32 (today, trimmed)
async def get_write_session(request: Request) -> AsyncIterator[AsyncSession]:
    factory: async_sessionmaker[AsyncSession] = request.app.state.session_factory
    async with factory() as session:
        try:
            yield session
            await session.commit()
        except Exception:
            await session.rollback()
            raise
        await dirty.flush(request.app.state.cache, session)

In plain words: the commit runs in the teardown of a FastAPI dependency. In FastAPI 0.139.0 (the installed version), a dependency's teardown runs on the "request" exit stack, which closes after await response(scope, receive, send) has sent the answer (fastapi/routing.py:132-137, fastapi/dependencies/utils.py:669-671; checked 7 Oct). So send_command (bridge.py:213) can answer "accepted" and then fail to commit. This is the C4 bug. It hits every console edit. The fix: the command road commits inside itself, before it returns.

The tables Agreed, to build

Two new tables and two new columns. The match lock needs no new table: it is the fixture row (and ledger_stream for an event).

-- NEW. Every command we finished, with the answer we gave. A retry reads it.
CREATE TABLE command (
  command_id    text PRIMARY KEY,         -- the key: made by the sender, sent again on every retry
  sender        text NOT NULL,            -- device, provider or integration
  subject_id    uuid NOT NULL,            -- the match (or event) it changed
  payload_hash  bytea NOT NULL,           -- same key + different content = refused
  outcome       text NOT NULL,            -- 'accepted' or 'refused'
  answer        jsonb NOT NULL,           -- exactly what we answered
  seq           bigint,                   -- the log row it made, if accepted
  received_at   timestamptz NOT NULL,
  finished_at   timestamptz NOT NULL
);                                        -- kept 30 days, then deleted by a daily job

-- NEW. Commands that failed with an error (not a rule's no). Written in its own
-- short transaction, because the failed one was rolled back.
CREATE TABLE command_failure (
  command_id  text PRIMARY KEY,
  attempts    int NOT NULL,
  last_error  text NOT NULL,
  last_at     timestamptz NOT NULL
);

-- NEW COLUMN. Every row a command can edit gets the log number of its last change.
ALTER TABLE fixture_competitor ADD COLUMN changed_seq bigint;   -- and the other editable rows

-- NEW COLUMN. The outbox row records its transaction, so readers never skip a late commit.
ALTER TABLE domain_event ADD COLUMN xid xid8 NOT NULL DEFAULT pg_current_xact_id();

The log, timeline_item, stays as it is: one row per accepted command, numbered per match, never updated or deleted. Today it also holds the retry key, with a unique index on stream, source and key (uq_timeline_item_idempotency, packages/contract/src/omnium_contract/models/timeline.py:59).

Why a separate command table, and why 30 days (C2, L13). The command row holds the answer, including a refusal, which the log never holds. A retry within 30 days always gets the same answer. The daily job deletes old rows in batches of 10,000, so it never holds a long lock. If the table ever passes about 50 million rows, it switches to one partition per day and drops a whole day at once.

Why it stays in the same database. The user asked whether the command table could move to another database later. The design advises against it. "Was it saved?" has a true answer only because the command row and the change commit together (L3, L11). With two databases, a crash between two commits makes that answer a guess. Rows older than 30 days can be copied elsewhere for history.

Size. Expected rows: 50 matches a day (Section 1) × about 1,500 commands a match (estimate) = 75,000 a day, about 2.3 million in 30 days. The speed test used 10 million rows, four times that. The index has 3 levels at 10 million rows (737 MB). Ten times more rows adds about one level: one more page read, a few microseconds. bigint lasts about 29 million years at 10,000 commands a second, and the log number counts per match.

One command, statement by statement Agreed, to build

This is the whole transaction for one command on one match, from the design doc. The doc's first draft had three SET LOCAL lines after BEGIN; they are left out here, because L16 moved the limits onto the connection pool. Everything between BEGIN and COMMIT is saved together or not at all.

-- 0. Fast path, before any lock: a retry of a finished command is answered at once.
SELECT outcome, answer, payload_hash FROM command WHERE command_id = $id;

BEGIN;                                              -- READ COMMITTED
-- (time limits come from the connection pool, not SET LOCAL: see L16)

-- 1. Take the match lock and read its log number. NO KEY UPDATE, not UPDATE,
--    so other tables can still insert rows that point at this match.
SELECT last_seq FROM fixture WHERE id = $match FOR NO KEY UPDATE;

-- 2. Seen this key? Checked again under the lock, so two copies of one
--    command arriving together can never both run.
SELECT outcome, answer, payload_hash FROM command WHERE command_id = $id;
--    found, same hash  -> ROLLBACK; answer "duplicate" with the saved answer
--    found, other hash -> ROLLBACK; answer "refused: key used for another command"

-- 3. Version check, only on the rows this command edits.
SELECT changed_seq FROM fixture_competitor WHERE id = ANY($rows);
--    any changed_seq different from the one the sender saw -> ROLLBACK; "conflict"

SAVEPOINT work;
-- 4. The sport's rules and the change, worked out in plugin code (no SQL inside).
-- 5. The writes:
INSERT INTO timeline_item (stream_id, seq, ...) VALUES (..., $last_seq + 1, ...);
--    the claims and the winner (Section 2), and the edited rows with
--    changed_seq = $last_seq + 1
INSERT INTO domain_event (...) VALUES (...);         -- the note to feeds and screens
UPDATE fixture SET last_seq = $last_seq + 1 WHERE id = $match;
INSERT INTO command (command_id, ..., outcome, answer, seq)
     VALUES ($id, ..., 'accepted', $answer, $last_seq + 1);

--    If a rule said no:
--      ROLLBACK TO SAVEPOINT work;
--      INSERT INTO command (..., outcome, answer) VALUES (..., 'refused', $why);

COMMIT;                                              -- synchronous_commit = on
-- 6. Only now is the answer sent.

In plain words:

  • The key is checked twice. The quick check answers most retries without waiting for the lock. The check under the lock is the one that counts. If two copies arrive together, both pass the quick check, but only one gets the lock first; the second then finds the first's row. The primary key on command_id is the last guard: a third copy can never insert a second row.
  • The version check looks only at the rows this command edits (L4). If it checked the whole match, every new goal from a provider would make every operator edit fail with "conflict" during a busy match.
  • A rule's "no" rolls back to the savepoint and saves only the refusal row. So a refused command writes nothing else, and support can still read why (C5).
  • The version has three levels set per command (L17): none (a commentary line), rows (an edit, using changed_seq), and play (every ball, goal or point carries the last play event the sender saw).
  • Isolation is READ COMMITTED, the Postgres default, plus one rule: no write to a record without holding its lock (L12). SERIALIZABLE was not chosen because under normal load it would turn many commands into retries.

The road as code Agreed, to build

These files are new; the code is from the design doc's low-level part and is not in the repo yet. Section 5 builds its engine on the same road.run.

FileNew or todayWhat it holds
core/.../engine/road.pynewrun(): the one entry point for every command
core/.../engine/locks.pynewtake(): locks the event, matches, then people and teams, sorted. Replaces ledger.lock_key and every pg_advisory_xact_lock call (engine.py:116, pipeline.py:869)
core/.../engine/answers.pynewThe command table: prior(), save(); Answer and its five kinds
core/.../engine/versions.pynewstale_rows(): the changed_seq check
core/.../engine/failures.pynewCounts failed tries in command_failure; refuses after 5
core/.../outbox/reader.pynewread_after(): the xid8 read for the live stream and the stats queue
core/.../db.pytodayGains build_command_engine() and build_import_engine()
admin/.../deps.pytodayCommands stop using get_write_session; road.run commits itself
admin/.../routers/bridge.pytodaysend_command() builds a Command, calls road.run, maps the five answers to HTTP
# core/src/omnium_core/engine/locks.py   (new, agreed design)
ORDER = ("event", "fixture", "participant")          # L2: always this order, ids sorted inside each

async def take(session: AsyncSession, want: Locks) -> dict[uuid.UUID, int]:
    """Lock every record the command changes. Returns each fixture's last_seq."""
    seqs: dict[uuid.UUID, int] = {}
    if want.event:
        await session.execute(sa.select(Competition.id).where(Competition.id == want.event)
                              .with_for_update(key_share=False, read=False, nowait=False, of=Competition))
    if want.fixtures:
        rows = await session.execute(sa.select(Fixture.id, Fixture.last_seq)
                                     .where(Fixture.id.in_(sorted(want.fixtures)))
                                     .order_by(Fixture.id)
                                     .with_for_update(key_share=True))   # = FOR NO KEY UPDATE
        seqs = dict(rows.tuples())
    if want.participants:
        await session.execute(sa.select(Participant.id).where(Participant.id.in_(sorted(want.participants)))
                              .order_by(Participant.id).with_for_update(key_share=True))
    return seqs
# lock_timeout comes from the connection (L16); a timeout raises 55P03 -> try_again(lock_timeout)

In plain words: one statement per kind, in the fixed order, with the ids sorted. In SQLAlchemy, with_for_update(key_share=True) is how you write FOR NO KEY UPDATE. The same read returns each match's last_seq, so the lock and the log number come in one trip. Postgres error 55P03 means "lock not available", and the road turns it into "try again".

# core/src/omnium_core/engine/answers.py   (new, agreed design)
def payload_hash(cmd: Command) -> bytes:
    # canonical JSON: sorted keys, no spaces, so the same command always hashes the same
    body = json.dumps({"code": cmd.code, "input": cmd.input, "subject": str(cmd.subject)},
                      sort_keys=True, separators=(",", ":"), default=str)
    return hashlib.sha256(body.encode()).digest()

async def prior(session, cmd) -> Answer | None:
    row = (await session.execute(sa.select(command.c.payload_hash, command.c.answer)
                                 .where(command.c.command_id == cmd.key))).first()
    if row is None:                         return None
    if row.payload_hash != payload_hash(cmd): return Answer.refused("key_reused", "This key was used for a different command.")
    return Answer.from_saved(row.answer, duplicate=True)

async def save(session, cmd, answer) -> None:
    await session.execute(sa.insert(command).values(
        command_id=cmd.key, sender=cmd.source, subject_id=cmd.subject, payload_hash=payload_hash(cmd),
        outcome=answer.outcome, answer=answer.as_json(), seq=answer.seq,
        received_at=cmd.received_at, finished_at=sa.func.now()))
    # a second copy racing past prior() fails here on the primary key -> answered as duplicate

In plain words: the hash is of the code, the input and the subject, in a fixed JSON form. prior() is used both for the quick check and the check under the lock. A screen bug that sends a used key with a different body is refused, so a bug cannot hide a change. Section 5 measured the cost: sha256 of a 10 KB state takes 0.0044 ms, canonical JSON of it 0.05 ms (laptop, 6 Oct).

# core/src/omnium_core/engine/versions.py   (new, agreed design)
async def stale_rows(session, cmd) -> list[dict] | None:
    saw = cmd.saw.rows if cmd.saw else {}
    if not saw: return None                                     # an add: no check (C3)
    now = dict((await session.execute(sa.select(FixtureCompetitor.id, FixtureCompetitor.changed_seq)
                .where(FixtureCompetitor.id.in_(saw)))).tuples())
    stale = [rid for rid, seq in saw.items() if now.get(uuid.UUID(rid)) != seq]
    return await read_rows(session, stale) if stale else None   # the rows as they are now, for the screen

In plain words: compare each row's changed_seq with what the sender saw. Any row that moved on goes back in the conflict answer, as it is now, so the screen can update without a reload.

Pools and time limits Agreed, to build

Every wait has a limit, and crossing a limit gives an answer the sender knows how to handle. The limits are safety nets, not speed targets (L6): set about 1,000 times above today's measured median, so a slow moment never refuses a correct command.

LimitA person's commandAn import batchFires whenAnswer
Lock wait (lock_timeout)2 s5 sIt waited too long for the match lockTry again
One query (statement_timeout)5 s30 sOne query ran too longTry again
Silence in a transaction (idle_in_transaction_session_timeout)5 s30 sThe server stopped talking mid-transactionPostgres ends the session and frees the lock
The sport's code50 ms of CPU2 s of CPUA rule ran too longTry again; after 5 tries, refused plus an alert
The screen waiting for an answer5 snot usedNo answer cameSends again with the same key

The speed target is separate: a person's command under 30 ms of server time at the 95th percentile. Crossing it refuses nobody. An alert fires if it stays over for 5 minutes.

Where the limits live (L8, L16). Today a person's command sets two limits with SET LOCAL on every transaction:

# packages/core/src/omnium_core/scoring/pipeline.py:838-842 (today)
ms = int(get_settings().scoring_txn_timeout_ms)          # 10_000 (settings.py:290)
if ms <= 0:
    return
await session.execute(sa.text(f"SET LOCAL statement_timeout = {ms}"))
await session.execute(sa.text(f"SET LOCAL idle_in_transaction_session_timeout = {ms}"))

Each SET LOCAL line is one more round trip. Measured 5 Oct: the three SET LOCAL lines in the design's first draft added about 0.3 ms to every command (1.48 ms with them, 1.17 ms without, one sender). The fix costs nothing: db.py already sends statement_timeout on every new connection through asyncpg's server_settings, so the other limits ride the same way.

# packages/core/src/omnium_core/db.py:24-37 (today, trimmed)
return create_async_engine(
    url or settings.database_url,
    pool_size=settings.db_pool_size,              # 10 (settings.py:270)
    max_overflow=settings.db_max_overflow,        # 5
    connect_args={
        "server_settings": {"statement_timeout": str(settings.db_statement_timeout_ms)}  # 15_000
    },
)
# core/src/omnium_core/db.py   (agreed: two new callers of build_engine)
def build_command_engine(s: Settings) -> AsyncEngine:      # people and scorers
    return build_engine(s, server_settings={"lock_timeout": "2000", "statement_timeout": "5000",
                                           "idle_in_transaction_session_timeout": "5000"})
def build_import_engine(s: Settings) -> AsyncEngine:       # integration batches
    return build_engine(s, server_settings={"lock_timeout": "5000", "statement_timeout": "30000",
                                           "idle_in_transaction_session_timeout": "30000"})

In plain words: two pools, one for people's commands and one for imports, each with its own limits fixed on the connection when it opens. A command sends no SET lines. Commands and reads also use separate pools, so one busy match can never use up the connections that screens read with.

At most 2 waiting per match (L8). A command waiting for a lock holds a database connection. So each container lets at most 2 commands per match wait in Postgres at once (an asyncio.Semaphore(2) per fixture id). The rest wait in that container's memory, without a connection. With 4 containers, at most 8 connections wait on one match. Nothing correct depends on this line; it only saves connections. Postgres still decides the real order, first come, first served.

Two database trips per command (M9) Agreed, to build

The statement-by-statement transaction above is about 9 round trips. On AWS one trip in the same zone costs about 0.3 to 1 ms (estimate, to measure in Section 12). Nine trips would cost more than the work itself. Section 5 agreed M9: a command makes 2 trips:

  1. One statement that takes the lock and reads the key, the state and the row versions together.
  2. One call to a database function that does all the writes and the commit.

That brings a command to about 2 to 3 ms on prod (estimate). It does not grow with the number of matches until the database itself is busy. The function's details belong to Section 10 (Database).

Where the command road runs Agreed, to build

PieceWhere it runs (agreed, L15)Why
Command serviceIts own ECS service, several containers. Answers inside the HTTP request, no queueA busy admin page or public API cannot slow a scorer. A queue would add a hop and a wait
The sport's codeA small pool of worker processes in each command container (L10)A process can be killed at a limit; a Python thread cannot
ImportsInside the command service, in a process of their own at lower CPU priority (changed 8 Oct 2026). Batches go through the same engine code, on their own connection poolThe scorer's process always gets the CPU first; a long import only takes longer
Writes and scoring readsPostgres onlyOne place that is always right
Publishing to clientsRedis cache and the CDN (Section 9)Millions of reads never touch the write path

Today commands and imports both run inside the admin service, and the sport's code runs on the request's own thread. flows_host.py:240-241 says so itself: "a Python component cannot be [cut off], so the budget is enforced after the fact". Section 5 later narrowed L10: only sports scored event by event start workers; apply mode (operations and scorecards) never does, and a worker's wall-clock limit is 500 ms per ball (Scoring engine).

Senders Agreed, to build

The bridge outbox. The bridge gains a queue that keeps the key, saved in IndexedDB per device and subject, so a reload or a crash does not lose an unsent change.

// packages-ts/omnium-bridge/src/outbox.ts   (new, agreed design)
async push(cmd: Omit<Pending, "key" | "tries" | "sentAt">): Promise<Answer> {
  const p = { ...cmd, key: crypto.randomUUID(), sentAt: new Date().toISOString(), tries: 0 };
  await this.store.add(p);              // saved BEFORE the first send
  return this.drain(p.key);             // one at a time, oldest first (C12)
}
private async sendOne(p: Pending): Promise<Answer> {
  for (;;) {
    const a = await this.http.command(p);                 // same key every time
    if (a.outcome !== "try_again" && !a.networkError) {
      await this.store.delete(p.key);                      // accepted, duplicate, refused, conflict
      return a;                                            // conflict stops the queue: the app decides
    }
    p.tries += 1;
    await sleep(Math.min(30_000, 200 * 2 ** p.tries) * (0.5 + Math.random()));   // backoff + jitter
  }
}

In plain words: save first, then send. Send one at a time, oldest first, so an undo never arrives before the ball it undoes (C12). Only "try again" and a network error are retried, with the same key, after a growing wait plus randomness so all screens do not come back at the same moment. A conflict stops the queue until the person decides.

Imports (L7). A run's rows are grouped by match, sorted by fixture id, and cut into batches of at most 25 matches. Each batch is its own command and its own transaction, with key {run_id}:{batch_no}, on the import pool. A retried batch is a duplicate. Today the per-chunk key already exists (f"{key}:{index}", packages/core/src/omnium_core/integrations/service.py:452), but every chunk runs in one transaction owned by the caller (service.py:17), under one event-wide lock. See Imports.

The outbox, and the late commit Partly built

The outbox is the domain_event table. Each command inserts one row in the same transaction as its change, so a change and its note can never disagree. That part is built (packages/core/src/omnium_core/jobs_outbox.py:44, called at engine.py:171).

The bug is in how readers follow it. Today two readers read "every row with an id greater than the last one I saw":

# packages/admin/src/omnium_admin/routers/bridge.py:915-928 (today, trimmed): the live stream to screens
rows = (await session.execute(
    sa.select(DomainEvent.id, DomainEvent.kind, DomainEvent.fixture_id,
              DomainEvent.competition_id, DomainEvent.payload)
    .where(DomainEvent.id > cursor)
    .order_by(DomainEvent.id)
    .limit(500)
)).all()

The stats queue does the same (packages/stats/src/omnium_stats/triggers.py:68). A row gets its id when it is inserted, but readers see it only after its transaction commits:

  1. Command A inserts its note and gets id 100. A is still running.
  2. Command B inserts its note, gets id 101, and commits first.
  3. The reader reads 101 and moves its mark to 101.
  4. A commits. Row 100 is now visible, but the reader is already past it. It never reads it.

A console screen misses a change, and a stat is never worked out. It is rare, but silent, and it gets more likely as more commands run at once.

A sketch in two panels over one database cylinder labelled domain_event. The left panel, Today: read id after last, shows A writes note 100 and commits late, B writes note 101 and commits first, the reader reads 101 and sets its mark to 101, and a red box: note 100 skipped forever. The right panel, Agreed: read by transaction, shows A still running, the reader waits below A, A commits, and a green box: the reader reads 100, then 101.
Reading by id can skip a late commit forever. Reading by transaction, below the oldest one still running, cannot.

The fix (L9). Each note stores the id of its transaction, as xid8. A reader reads only notes from transactions older than the oldest one still running.

# core/src/omnium_core/outbox/reader.py   (new, agreed design)
async def read_after(session, cursor: tuple[int, int], limit: int = 500) -> tuple[list[Row], tuple[int, int]]:
    rows = (await session.execute(sa.text("""
        SELECT id, xid::text::bigint AS x, kind, fixture_id, seq, payload FROM domain_event
         WHERE (xid::text::bigint, id) > (:x, :id)
           AND xid < pg_snapshot_xmin(pg_current_snapshot())
         ORDER BY xid, id LIMIT :limit"""), {"x": cursor[0], "id": cursor[1], "limit": limit})).all()
    return rows, ((rows[-1].x, rows[-1].id) if rows else cursor)

In plain words: the cursor is now a pair (transaction id, row id). No transaction older than the "oldest still running" point can add a note any more, so nothing behind the reader's mark can appear later. Notes wait a few milliseconds longer than today. A long transaction makes every reader wait, which is one more reason imports run in short batches. Old rows get xid = 0 in the migration and are read first. Each reader keeps its own cursor: the live stream in memory per container, the stats queue in its existing cursor row.

Pruning: jobs_outbox.prune deletes notes older than 30 days (jobs_outbox.py:41, :75). I found no service that calls it outside tests (grep, 7 Oct).

Metrics and alerts Agreed, to build

MetricKindAlert
omnium_command_seconds{workflow, outcome}histogram95th percentile of a person's command over 30 ms for 5 minutes
omnium_command_lock_wait_seconds{workflow}histogramany wait over 1 s
omnium_command_answers_total{outcome, cause}countertry_again over 1% of commands for 5 minutes; any system failure
omnium_db_oldest_transaction_secondsgauge, from pg_stat_activity every 15 sover 10 s (an import batch: 30 s)
omnium_outbox_lag_seconds{reader}gaugeover 10 s

Alerts are grouped by match and rule, so one problem sends one message per 15 minutes, not hundreds. The runbook has one query that shows who blocks whom (pg_blocking_pids) and one that ends that session (pg_terminate_backend). Ending a session rolls back only that one command, and its sender sends it again. Routing is Monitoring.

What a scorer sees

What happens on the serverWhat the scorer sees
Accepted in a few msThe change, marked "sending", turns into a normal confirmed change
A slow moment: 300 ms instead of 3 msThe same, a little later. Nothing is refused
Lock wait over 2 s, or the database switching overStill "sending". The screen retries by itself with the same key after 0.2, 0.4, 0.8 s and so on. No error box
Database down for minutesStill "sending"; nothing is dropped. The team gets an alert
ConflictThe screen updates from the answer, keeps what was typed, marks what changed, asks to confirm
RefusedOne sentence with the real values and the way out, for example "A medal needs a finished match. Men's 100m Final is still LIVE. Set the match to Finished first, or leave the medal empty."
Our own bug"The system failed on this command. Our team has been alerted. Reference 7f3a." Never a stack trace

The screens are Section 7 (Console and apps).

When things go wrong

Every case ends with nothing lost and nothing counted twice. This is the agreed design; the "Built today" section says what is different now.

What goes wrongWhat happens
The network drops after the commit, before the answer arrivesThe screen sends again with the same key and gets "duplicate". One goal, not two
The server crashes before the commitPostgres rolls everything back, and the lock goes with the connection. No "accepted" was sent. The screen sends again
The server crashes after the commit, before the answerThe change is saved. The retry gets "duplicate"
The connection breaks during COMMIT itselfAnswer "try again". The retry looks for the command row: there means saved (duplicate), missing means not saved (it runs). Always right, because the row commits with the change
Two people change the same value at onceBoth wait for the match lock. The first saves. The second carries an old version and gets "conflict", with the new value
An import batch and a hand edit hit the same matchBoth take the same match lock. One waits a few ms. Neither overwrites the other unseen
The sport's code throws an errorRolled back, lock freed at once. "Try again". A failure count is saved in its own small transaction. After 5 failures: "refused: the system failed", an alert, and that scorer's later commands wait until a person decides
The sport's code loopsThe worker process is killed at its limit; the command rolls back; "try again"
The server hangs inside a transactionPostgres ends the session after 5 s of silence (an import batch: 30 s) and frees the lock
Two commands wait for each other in a circleCannot happen with sorted locks. If it does, Postgres cancels one in about 1 s; it gets "try again"
An import batch hangs halfwayOnly its 25 matches are locked. Edits to the other matches and the medal table run as normal. After 30 s of silence Postgres ends it; the batch is retried alone; the run shows "partly done"
A deploy during a matchThe old server stops taking new commands and finishes the ones it has (ECS gives 30 s by default). A command cut off before its commit is rolled back and sent again
The database switches to its standby (up to 2 minutes)"Try again" until it is back. Screens keep their commands. Whether prod copies every commit to the standby before confirming: not checked (Section 10)
A screen bug sends a used key with a different bodyRefused: key_reused
Bad input: wrong shape, missing fieldRefused, and saved in the refusals log
A provider sends the same message twiceIts key comes from its message number, so the second is a duplicate
A scorer comes back online with 6 unsent changesSent in the order they were made, one at a time
The outbox reader runs while a slow transaction is openIt waits below that transaction, then reads its note. Nothing skipped

Can a lock freeze the system? The user asked this directly. In the agreed design, no. A lock is held for milliseconds on one match, and Postgres always lets it go: at commit, at rollback, when the connection dies, after 5 s of silence, or when a statement passes its limit. Reads never wait. The worst case: one match is slow for up to 5 s (an import batch: 30 s), its commands get "try again", and every other match carries on.

Today has two real ways to freeze (checked in code on 5 Oct, per the design doc):

Risk todayWhat happensFixed by
An import holds the whole event's lock for its whole runImports skip the person's guard (engine.py:114), so only the 15 s per-query limit applies (settings.py:275). If an import hangs between queries, the event lock stays until its process restarts, and every Games command for that event fails after 10 sL6, L7
The sport's code runs on the server's own threadA looping rule: Postgres frees the lock after 10 s, but that process answers nothing, for any match, until restartedL10

Decisions

All 30 decisions in Section 4 are Agreed (closed 6 Oct 2026). M9 is from Section 5 and is listed because it shapes this road.

#DecisionIn plain words
C1One road: every change, from any source, is a command through one engine. No side doorsThe admin panel's match form, which writes without a command today, moves onto the road
C2Every command has a key the sender makes once and sends on every retry. The first answer is kept 30 days. A used key with different content is refusedPhone resends key 7f3a, gets "duplicate". One goal
C3An edit, removal or correction carries the version the sender saw. An old version gets "conflict". Adds may skip it (per command)Two operators at version 14: the second gets "the score changed to 2–1"
C4"Accepted" goes out only after the save is on diskFixes today's answer-before-commit
C5Every refusal is saved, with who, what and why. A refused command writes nothing elseSupport can say "refused at 14:02, a medal needs a finished match"
C6Every command that changes a match takes that match's lock, imports included. Big imports go in small batches, in a fixed orderAn import and a hand fix can no longer overwrite each other
C7Five answers: accepted, duplicate, refused, conflict, try againEach tells the screen what to do
C8Over the time budget means "try again", not "refused". Budgets are per workflowOn 1 Oct a correct fix was refused for taking 50 ms
C9One engine, two ways to work out a change: "apply once" and "fold"One set of rules to test
C10One transaction saves the log row, claims and winner, the outbox note and the command rowA goal and its note to feeds can never disagree
C11Every command can be tried first ("Try"), with nothing savedExists today; kept
C12Each sender sends its unsent commands in order, one at a timeAn undo never arrives before its ball
L1The match lock is the match's own row, FOR NO KEY UPDATE; it replaces the 8-byte advisory keyEvery match gets its own lock
L2Lock in one fixed order: event, matches, people and teams, each sorted by idNo deadlock circles
L3Key checked before the lock and again under it; the command row commits with the changeTwo copies can never both run; "was it saved?" is always answerable
L4The version check covers only the rows a command changes (changed_seq)Provider goals do not make every operator edit conflict
L5A rule's "no" is saved after rolling back to a savepoint. An error is counted in its own transaction; after 5, refused plus an alert, and later commands waitBugs are reported, never retried forever
L6Time limits are safety nets far above normal; crossing one means "try again". Sport code is limited by CPU time. Speed targets are watched by alerts. All values are settings per workflowA 300 ms command is still saved
L7Imports run in batches of at most 25 matches, each its own transaction and keyAn operator waits at most one short batch
L8Commands and reads use separate pools; at most 2 commands per match wait per containerOne busy match cannot use up every connection
L9Outbox readers read by transaction id, below the oldest running transactionA late commit is never skipped
L10The sport's code runs in worker processes that can be killed at a limitOne bad rule costs one command, not a frozen server
L11The answer goes out only after COMMIT returns; a broken commit answers "try again""Accepted" always means saved
L12READ COMMITTED, plus: no write to a record without its lock. A two-writer test proves itThe standard setting, made safe by the lock
L13Command rows kept 30 days, deleted daily. The log is never deletedSame answer for 30 days; the table does not grow forever
L14Locks are watched: alert on a lock wait over 1 s or a transaction open over 10 s (import: 30 s); a runbookWe hear from an alert, not a scorer
L15The command road runs in its own ECS service, answers inside the request. Imports run inside it, in a process of their own at lower CPU priority (changed 8 Oct 2026; before: the stats-worker, and a separate worker)A scorer is never slowed by an admin page, and an import cannot take the scorer's CPU
L16Time limits are set once per connection pool, not with SET LOCAL. As few round trips as possibleSaves about 0.3 ms per command
L17The version check has three levels per command: none, rows, play. A conflict carries the missed eventsTwo scorers on one match never save on top of unseen events
L18Every refusal has a code, a sentence with the real values, the field, the value now, and what to do. A test fails the build if a rule has no sentenceNever just "rejected"
M9A command makes 2 trips to the database: one read statement, one write function (Section 5)About 2 to 3 ms per command on prod (estimate)

Built today, or still to build

Checked in the code on 7 Oct 2026 (repo at commit 76ef1e2), unless a row says otherwise.

PieceTodayAgreed designStatus
One engine for apply-mode commandsengine.run, packages/core/src/omnium_core/workflows/engine.py:89road.run for every command, apply and foldPartly built
Live scoring engineA second engine, pipeline.command, scoring/pipeline.py:845; not on prodSame road, "fold" mode (Section 5)Partly built
Answer after commitNo: commit in deps.py:28, which FastAPI 0.139 runs after the response is sentroad.run commits before it returns (C4, L11)Agreed, to build
Retry keyOptional (bridge.py:187); bridge makes a new one per call (app.ts:332); console passes noneRequired, 8 to 128 characters, kept by the bridge outboxPartly built
Key storageIn timeline_item.idempotency_key with a unique index (timeline.py:59); no body hash; refusals not keptcommand table with payload_hash and the full answer, 30 daysAgreed, to build
Match lockAdvisory lock on the first 8 bytes of the id (engine.py:116, pipeline.py:869, ledger.py:167)Row lock FOR NO KEY UPDATE, sorted (L1, L2)Agreed, to build
Import lockOne event-wide advisory lock for the whole run, one transaction (integrations/service.py:17, :444)Batches of 25 matches, each its own transaction, locking its own matches (L7)Agreed, to build
Version checkNonechanged_seq per row; three levels (C3, L4, L17)Agreed, to build
Refusals logNot savedSaved in command with code and sentence (C5, L18)Agreed, to build
Five answersTwo shapes: 200 with accepted: false, or an HTTP error (bridge.py:77-91)One body, outcome of five (C7)Agreed, to build
Outbox in the same transactionYes, jobs_outbox.emit (jobs_outbox.py:44, engine.py:171)Kept (C10)Built today
Outbox readersBy id, can skip a late commit (bridge.py:924, triggers.py:68)By xid8 below the oldest running transaction (L9)Agreed, to build
Time limitsSET LOCAL 10 s for a person's command (pipeline.py:841-842); 15 s per query on every connection (db.py:34-36); no lock limit; imports get no silence limitPool-level server_settings, two pools, lock 2 s, query 5 s, silence 5 s (L6, L16)Partly built
Sport code isolationOn the request thread; budget checked after the fact (flows_host.py:240)Killable worker processes (L10, refined in Section 5)Agreed, to build
"Try" (dry run)Yes: savepoint rolled back (engine.py:139, :220-222)Kept (C11)Built today
Own service for commandsCommands run inside the admin service; imports in the admin service and the schedulerCommand ECS service; imports in a process of their own inside it (L15, changed 8 Oct)Agreed, to build
2 database tripsAbout 9 trips as designed; more today (research counted 35 to 45 per ball in live scoring)One read statement, one write function (M9)Agreed, to build
Lock alerts and metricsNone for commandsFive metrics, two lock alerts, a runbook (L14)Agreed, to build

Numbers

NumberWhatWhere it came from
0.09 ms / 0.16 msKey lookup, key not seen, 1 sender, median / 99th percentileMeasured, laptop, 5 Oct, 10 million command rows (design doc)
0.24 ms / 0.34 msThe same, 8 senders at onceMeasured, as above
0.44 ms / 1.42 msKey lookup, key seen (a real retry), 1 senderMeasured, as above
0.23 ms / 0.51 msThe same, 8 sendersMeasured, as above
1.48 ms / 3.98 ms (614 a second)The whole command transaction, different matches, 1 senderMeasured, as above
3.41 ms / 4.72 ms (2,276 a second)The same, 8 sendersMeasured, as above
1.17 ms / 1.78 ms (812 a second)Whole transaction without the three SET LOCAL lines, 1 senderMeasured, as above
2.69 ms / 6.01 ms (2,853 a second)The same, 8 sendersMeasured, as above
8.67 ms / 12.89 ms (904 a second)8 senders all on one match, all waiting for the same lockMeasured, as above
3 levels, 737 MBThe command table's index at 10 million rowsMeasured, as above
0.0044 ms / 0.05 mssha256 / canonical JSON of a 10 KB stateMeasured, laptop, 6 Oct (Section 5)
3 to 8 msLock held per command on prodEstimate (design doc), to measure in Section 12
0.3 to 1 msOne database round trip on AWS, same zoneEstimate (design doc)
2 to 3 msA command with 2 trips (M9) on prodEstimate (Section 5)
20 to 60 msNetwork from India to an AWS region in IndiaEstimate (design doc)
about 1 a secondCommands from one live matchEstimate (design doc)
75,000 a dayCommands at 50 matches a day × 1,500Estimate (design doc)

The test setup. A laptop with 8 cores and 16 GB, Postgres 17 with only 128 MB of its own memory. Tables: 10 million command rows (2.7 GB), 5 million log rows, 100,000 matches with 2 million side rows. Each run: 20 seconds with pgbench. No network between app and database; prod adds a round trip per statement. The load tests on AWS (Section 12) fail if the lock is held over 20 ms at the 99th percentile.