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.
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.
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.
| # | Question | Why it matters | Answer in short |
|---|---|---|---|
| 1 | Does every change take the same road? | A side door skips the checks, the log and the priority rules | Yes, no side doors (C1) |
| 2 | How is a retry made safe? | Networks drop. A second tap must not make a second goal | The same key gives the same answer (C2) |
| 3 | What if two people change the same thing at once? | One change must not wipe out the other unseen | A lock for order, a version check for "you saw an old one" (C3, C6) |
| 4 | When may the server say "saved"? | If it says "saved" and then crashes, the scorer trusts data that is gone | Only after the commit returns (C4) |
| 5 | What 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) |
| 6 | One engine, or two? | Two engines means two sets of bugs | One 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).
| Part | Example | What it is |
|---|---|---|
code | results.edit_side | Which command of the workflow to run. The plugin defines it. |
input | side, score | The values. The command's own model checks their shape. |
key | 7f3a2c1e-… | The retry key. The sender makes it once and sends it again on every retry. |
dry_run | false | "Try": run everything, then undo it. |
source | desk-2 | Which desk sent it, when a match has several. |
saw (new) | changed_seq per row, or the last play event | The version the sender saw. Agreed, not built. |
sentAt (new) | 2026-10-06T10:00:00.120Z | When 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

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.
- Quick key check, before any lock. One index lookup in the
commandtable. 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. - 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. - Key check again, under the lock. Seen with the same body: duplicate, with the first answer. Seen with a different body: refused (
key_reused). - Version check. An edit carries the
changed_seqof 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. - 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).
- 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 newfixture.last_seq, and thecommandrow with the exact answer. - Answer. "Accepted" goes out only after
COMMITreturns.
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.
| Outcome | HTTP | What it means | What the sender does |
|---|---|---|---|
accepted | 200 | Saved on disk | Removes it from its unsent list |
duplicate | 200 | Already saved; here is the first answer, exactly | The same as accepted |
refused | 422 | A rule said no. The body has a code, a sentence with the real values, the field, the value now, and what to do | Shows the sentence. Does not send it again |
conflict | 409 | You saw an old version. The body has the rows as they are now, or the events you missed | Updates the screen from the answer, asks the person |
try_again | 503 + Retry-After | Busy, a limit was hit, or the database is switching over. Cause: lock_timeout, worker_timeout, code_missing or db_unavailable | Sends 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 42 | Waits. |
FOR UPDATE, FOR SHARE, DELETE of match 42, or a change to its id | Waits. |
| Anything on match 43 | Not 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:
| Time | Command X | Command Y |
|---|---|---|
| t0 | locks match 7 | locks match 3 |
| t1 | wants match 3: waits for Y | wants match 7: waits for X |
| t2 | waits forever | waits 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 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.
| Time | Scorer A, key 7f3a…, saw 41 | Scorer B, key c19e…, saw 41 | Database |
|---|---|---|---|
| 0 ms | Quick key check: not seen. Takes the lock on match 42 | fixture.last_seq = 41 | |
| 1 ms | Key check under the lock: not seen. Play version: saw 41, now 41, fine. Rules pass | Quick key check: not seen. Asks for the lock: waits in Postgres's queue | |
| 2 ms | Writes timeline_item seq 42, domain_event, last_seq = 42, and the command row | still waiting | |
| 3 ms | COMMIT. Answer: accepted, event 42 | Gets the lock. Key c19e not seen. Play version: saw 41, now 42. ROLLBACK. Answer: conflict, with event 42 inside | last_seq = 42; nothing of B written |
| 5 ms | Screen 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.
- The request reaches the server. The command is saved as seq 57 and committed. The answer is lost in the tunnel.
- 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.
- It sends the same command again, with the same key
7f3a. - Quick key check: the
commandrow for7f3ais there, and the body hash matches. Answer: duplicate, with the first answer (seq 57). This took about 0.44 ms on the laptop. - 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
commandrow. 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_idis 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.
| File | New or today | What it holds |
|---|---|---|
core/.../engine/road.py | new | run(): the one entry point for every command |
core/.../engine/locks.py | new | take(): 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.py | new | The command table: prior(), save(); Answer and its five kinds |
core/.../engine/versions.py | new | stale_rows(): the changed_seq check |
core/.../engine/failures.py | new | Counts failed tries in command_failure; refuses after 5 |
core/.../outbox/reader.py | new | read_after(): the xid8 read for the live stream and the stats queue |
core/.../db.py | today | Gains build_command_engine() and build_import_engine() |
admin/.../deps.py | today | Commands stop using get_write_session; road.run commits itself |
admin/.../routers/bridge.py | today | send_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.
| Limit | A person's command | An import batch | Fires when | Answer |
|---|---|---|---|---|
Lock wait (lock_timeout) | 2 s | 5 s | It waited too long for the match lock | Try again |
One query (statement_timeout) | 5 s | 30 s | One query ran too long | Try again |
Silence in a transaction (idle_in_transaction_session_timeout) | 5 s | 30 s | The server stopped talking mid-transaction | Postgres ends the session and frees the lock |
| The sport's code | 50 ms of CPU | 2 s of CPU | A rule ran too long | Try again; after 5 tries, refused plus an alert |
| The screen waiting for an answer | 5 s | not used | No answer came | Sends 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:
- One statement that takes the lock and reads the key, the state and the row versions together.
- 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
| Piece | Where it runs (agreed, L15) | Why |
|---|---|---|
| Command service | Its own ECS service, several containers. Answers inside the HTTP request, no queue | A busy admin page or public API cannot slow a scorer. A queue would add a hop and a wait |
| The sport's code | A small pool of worker processes in each command container (L10) | A process can be killed at a limit; a Python thread cannot |
| Imports | Inside 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 pool | The scorer's process always gets the CPU first; a long import only takes longer |
| Writes and scoring reads | Postgres only | One place that is always right |
| Publishing to clients | Redis 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:
- Command A inserts its note and gets id 100. A is still running.
- Command B inserts its note, gets id 101, and commits first.
- The reader reads 101 and moves its mark to 101.
- 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.

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
| Metric | Kind | Alert |
|---|---|---|
omnium_command_seconds{workflow, outcome} | histogram | 95th percentile of a person's command over 30 ms for 5 minutes |
omnium_command_lock_wait_seconds{workflow} | histogram | any wait over 1 s |
omnium_command_answers_total{outcome, cause} | counter | try_again over 1% of commands for 5 minutes; any system failure |
omnium_db_oldest_transaction_seconds | gauge, from pg_stat_activity every 15 s | over 10 s (an import batch: 30 s) |
omnium_outbox_lag_seconds{reader} | gauge | over 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 server | What the scorer sees |
|---|---|
| Accepted in a few ms | The change, marked "sending", turns into a normal confirmed change |
| A slow moment: 300 ms instead of 3 ms | The same, a little later. Nothing is refused |
| Lock wait over 2 s, or the database switching over | Still "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 minutes | Still "sending"; nothing is dropped. The team gets an alert |
| Conflict | The screen updates from the answer, keeps what was typed, marks what changed, asks to confirm |
| Refused | One 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 wrong | What happens |
|---|---|
| The network drops after the commit, before the answer arrives | The screen sends again with the same key and gets "duplicate". One goal, not two |
| The server crashes before the commit | Postgres 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 answer | The change is saved. The retry gets "duplicate" |
The connection breaks during COMMIT itself | Answer "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 once | Both 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 match | Both take the same match lock. One waits a few ms. Neither overwrites the other unseen |
| The sport's code throws an error | Rolled 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 loops | The worker process is killed at its limit; the command rolls back; "try again" |
| The server hangs inside a transaction | Postgres 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 circle | Cannot happen with sorted locks. If it does, Postgres cancels one in about 1 s; it gets "try again" |
| An import batch hangs halfway | Only 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 match | The 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 body | Refused: key_reused |
| Bad input: wrong shape, missing field | Refused, and saved in the refusals log |
| A provider sends the same message twice | Its key comes from its message number, so the second is a duplicate |
| A scorer comes back online with 6 unsent changes | Sent in the order they were made, one at a time |
| The outbox reader runs while a slow transaction is open | It 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 today | What happens | Fixed by |
|---|---|---|
| An import holds the whole event's lock for its whole run | Imports 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 s | L6, L7 |
| The sport's code runs on the server's own thread | A looping rule: Postgres frees the lock after 10 s, but that process answers nothing, for any match, until restarted | L10 |
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.
| # | Decision | In plain words |
|---|---|---|
| C1 | One road: every change, from any source, is a command through one engine. No side doors | The admin panel's match form, which writes without a command today, moves onto the road |
| C2 | Every 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 refused | Phone resends key 7f3a, gets "duplicate". One goal |
| C3 | An 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 disk | Fixes today's answer-before-commit |
| C5 | Every refusal is saved, with who, what and why. A refused command writes nothing else | Support can say "refused at 14:02, a medal needs a finished match" |
| C6 | Every command that changes a match takes that match's lock, imports included. Big imports go in small batches, in a fixed order | An import and a hand fix can no longer overwrite each other |
| C7 | Five answers: accepted, duplicate, refused, conflict, try again | Each tells the screen what to do |
| C8 | Over the time budget means "try again", not "refused". Budgets are per workflow | On 1 Oct a correct fix was refused for taking 50 ms |
| C9 | One engine, two ways to work out a change: "apply once" and "fold" | One set of rules to test |
| C10 | One transaction saves the log row, claims and winner, the outbox note and the command row | A goal and its note to feeds can never disagree |
| C11 | Every command can be tried first ("Try"), with nothing saved | Exists today; kept |
| C12 | Each sender sends its unsent commands in order, one at a time | An undo never arrives before its ball |
| L1 | The match lock is the match's own row, FOR NO KEY UPDATE; it replaces the 8-byte advisory key | Every match gets its own lock |
| L2 | Lock in one fixed order: event, matches, people and teams, each sorted by id | No deadlock circles |
| L3 | Key checked before the lock and again under it; the command row commits with the change | Two copies can never both run; "was it saved?" is always answerable |
| L4 | The version check covers only the rows a command changes (changed_seq) | Provider goals do not make every operator edit conflict |
| L5 | A 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 wait | Bugs are reported, never retried forever |
| L6 | Time 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 workflow | A 300 ms command is still saved |
| L7 | Imports run in batches of at most 25 matches, each its own transaction and key | An operator waits at most one short batch |
| L8 | Commands and reads use separate pools; at most 2 commands per match wait per container | One busy match cannot use up every connection |
| L9 | Outbox readers read by transaction id, below the oldest running transaction | A late commit is never skipped |
| L10 | The sport's code runs in worker processes that can be killed at a limit | One bad rule costs one command, not a frozen server |
| L11 | The answer goes out only after COMMIT returns; a broken commit answers "try again" | "Accepted" always means saved |
| L12 | READ COMMITTED, plus: no write to a record without its lock. A two-writer test proves it | The standard setting, made safe by the lock |
| L13 | Command rows kept 30 days, deleted daily. The log is never deleted | Same answer for 30 days; the table does not grow forever |
| L14 | Locks are watched: alert on a lock wait over 1 s or a transaction open over 10 s (import: 30 s); a runbook | We hear from an alert, not a scorer |
| L15 | The 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 |
| L16 | Time limits are set once per connection pool, not with SET LOCAL. As few round trips as possible | Saves about 0.3 ms per command |
| L17 | The version check has three levels per command: none, rows, play. A conflict carries the missed events | Two scorers on one match never save on top of unseen events |
| L18 | Every 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 sentence | Never just "rejected" |
| M9 | A 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.
| Piece | Today | Agreed design | Status |
|---|---|---|---|
| One engine for apply-mode commands | engine.run, packages/core/src/omnium_core/workflows/engine.py:89 | road.run for every command, apply and fold | Partly built |
| Live scoring engine | A second engine, pipeline.command, scoring/pipeline.py:845; not on prod | Same road, "fold" mode (Section 5) | Partly built |
| Answer after commit | No: commit in deps.py:28, which FastAPI 0.139 runs after the response is sent | road.run commits before it returns (C4, L11) | Agreed, to build |
| Retry key | Optional (bridge.py:187); bridge makes a new one per call (app.ts:332); console passes none | Required, 8 to 128 characters, kept by the bridge outbox | Partly built |
| Key storage | In timeline_item.idempotency_key with a unique index (timeline.py:59); no body hash; refusals not kept | command table with payload_hash and the full answer, 30 days | Agreed, to build |
| Match lock | Advisory 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 lock | One 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 check | None | changed_seq per row; three levels (C3, L4, L17) | Agreed, to build |
| Refusals log | Not saved | Saved in command with code and sentence (C5, L18) | Agreed, to build |
| Five answers | Two 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 transaction | Yes, jobs_outbox.emit (jobs_outbox.py:44, engine.py:171) | Kept (C10) | Built today |
| Outbox readers | By id, can skip a late commit (bridge.py:924, triggers.py:68) | By xid8 below the oldest running transaction (L9) | Agreed, to build |
| Time limits | SET 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 limit | Pool-level server_settings, two pools, lock 2 s, query 5 s, silence 5 s (L6, L16) | Partly built |
| Sport code isolation | On 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 commands | Commands run inside the admin service; imports in the admin service and the scheduler | Command ECS service; imports in a process of their own inside it (L15, changed 8 Oct) | Agreed, to build |
| 2 database trips | About 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 metrics | None for commands | Five metrics, two lock alerts, a runbook (L14) | Agreed, to build |
Numbers
| Number | What | Where it came from |
|---|---|---|
| 0.09 ms / 0.16 ms | Key lookup, key not seen, 1 sender, median / 99th percentile | Measured, laptop, 5 Oct, 10 million command rows (design doc) |
| 0.24 ms / 0.34 ms | The same, 8 senders at once | Measured, as above |
| 0.44 ms / 1.42 ms | Key lookup, key seen (a real retry), 1 sender | Measured, as above |
| 0.23 ms / 0.51 ms | The same, 8 senders | Measured, as above |
| 1.48 ms / 3.98 ms (614 a second) | The whole command transaction, different matches, 1 sender | Measured, as above |
| 3.41 ms / 4.72 ms (2,276 a second) | The same, 8 senders | Measured, as above |
| 1.17 ms / 1.78 ms (812 a second) | Whole transaction without the three SET LOCAL lines, 1 sender | Measured, as above |
| 2.69 ms / 6.01 ms (2,853 a second) | The same, 8 senders | Measured, as above |
| 8.67 ms / 12.89 ms (904 a second) | 8 senders all on one match, all waiting for the same lock | Measured, as above |
| 3 levels, 737 MB | The command table's index at 10 million rows | Measured, as above |
| 0.0044 ms / 0.05 ms | sha256 / canonical JSON of a 10 KB state | Measured, laptop, 6 Oct (Section 5) |
| 3 to 8 ms | Lock held per command on prod | Estimate (design doc), to measure in Section 12 |
| 0.3 to 1 ms | One database round trip on AWS, same zone | Estimate (design doc) |
| 2 to 3 ms | A command with 2 trips (M9) on prod | Estimate (Section 5) |
| 20 to 60 ms | Network from India to an AWS region in India | Estimate (design doc) |
| about 1 a second | Commands from one live match | Estimate (design doc) |
| 75,000 a day | Commands at 50 matches a day × 1,500 | Estimate (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.
Read next
- The live scoring engine: what "fold" does with a command once it is on this road, and the sport worker processes.
- Bridge and live updates: how the outbox note reaches a second screen.
- Imports from Scout: how a run becomes batches of 25 matches on this road.
- Database and the read-only copy: the 2-trip write function, growth and retention.