Getting data in · Imports
Imports from Scout
How rows that Scout scraped from a results site become commands on our data, and what an import must never do.
In one minute
An import reads one run of a Scout job (a list of rows scraped from a results site), maps each row onto our fields, and sends the rows that changed as import commands through the workflow engine. At the Asian Games almost every change on prod came this way.
Built today, an import is one long database transaction. It fetches from Scout over HTTP inside it, compares each row with what this integration last read, sends all changed rows under one lock on the whole event, cancels units a full run no longer lists, and tells nobody when it fails.
The agreed design (Section 3, 24 decisions) keeps today's no-code mapping, pins, guards and held rows, and changes the rest. The fetch happens with no transaction open. The run is saved first. Each row is compared with what our database holds now, minus fields a person pinned. Changes are saved in batches of at most 25 matches, each batch a command that is safe to repeat. A missing row becomes a question for a person, never a delete. Finished matches are read again for 72 hours. Every failure raises an alert through one central alert service.
Guards judge what a run would do, not how big it is. After the Games a run had only 4 rows, and those 4 mattered.
What this part does
An import answers one question again and again: what did the site change, and should our data change with it? It is the main way data gets into omnium for a multi-sport event.
Words used on this page:
| Word | Meaning |
|---|---|
| Integration | One Scout job connected to one import command, set up with no code. A row in integration. |
| Run | One read of one Scout run. A row in integration_run. |
| Row | One scraped item, for example one match. Each row gets a key (which record it is about) and a hash (a fingerprint of its mapped value). |
| Import command | A workflow command marked imports=True, for example games.import_units. It writes the rows. |
| Pin | A field name a person set by hand, saved in fixture.pinned. An import never overwrites a pinned field. |
| Held row | A row the import cannot place on its own. It waits for a person in integration_answer. |
The design section was built around seven questions. Each one went wrong at least once at the Asian Games:
| # | Question | What went wrong at the Games |
|---|---|---|
| 1 | How does an import decide that a row changed? | It compared with its own last read, not with our data. Golf and sailing results stayed out of date for days (F17, F36). |
| 2 | What may an import never overwrite? | The Section 2 priority list and the command's own rules decide who wins, never the import alone. A hand edit is a source near the top by default. |
| 3 | What does a row missing from a run mean? | A partial list looked like "everything else was deleted". |
| 4 | What happens with an id we do not know? | A new player or team must not be dropped or invented silently. |
| 5 | How is a run saved so a failure halfway does no damage? | A run that writes 300 matches and dies at match 150 must be safe to run again. |
| 6 | How long does an import keep reading a finished match? | Results changed after protests a day or more later, after we stopped reading. |
| 7 | How do we see that an import is failing? | A job failed Scout's checks for days and nobody was told (N4). |
How it works

The agreed run has six steps. What runs today is described step by step in the low-level design below.
- Fetch outside the database. A job reads the Scout run over HTTP with no database transaction open. A failed call is tried 3 times, with longer waits each time. A Scout run that failed its own checks is saved as "Scout failed" and alerted, not skipped quietly.
- Record the run first. In its own short transaction, the run is saved with status
receivedand its row count. From then on, no failure can erase the record that the run happened. - Map. Today's no-code steps, unchanged. Each row gets a key and a hash.
- Compare with our data, not with the last read. For each row, the engine reads the current values of that match (or entry, or medal row), leaves out every pinned field, and keeps only fields that really differ. A fast path skips the read when nothing has touched the record since this integration last wrote it.
- Save in batches of 25 matches. Each batch is one import command with the key
run:batch. It takes match locks, has import time limits, and is one transaction. A failed batch is tried again on its own; the others carry on. - Mark missing. After a full run, keys the run no longer listed are counted. Missing from two full runs in a row opens a review item. Nothing is cancelled by an import on its own.
Three more rules sit around these steps:
- Unknown ids are held. A row naming a team, person or match we do not know waits with a clear question. It is tried again by itself on every later run. A held row is never reported as saved.
- Finished matches are read again. After the source marks a result official, its page is read again for a settle window: hourly, then every 6 hours, then daily, up to 72 hours.
- A person is told. A failed run, 2 Scout runs in a row that failed their checks, a batch that failed 3 times, a guard hold, or no new run for three intervals: each raises an alert and shows on the integration's card.
A worked example
Example · A schedule run of 500 rows, where 4 changed
The setting. Day 9 of a Games. The schedule integration (games.import_units, build all) gets a doorbell from Scout: run r-8812 passed, 500 rows. Since the last run, the site changed 4 things. The values below are an example, not prod data.
| Key (example) | What the site changed |
|---|---|
ATH:M100METRES---FNL-000100 | Status went from unofficial to official |
SAL:race 7 (sailing) | Result changed after a protest, 30 hours after it was first official |
FBL:match 41 (football) | Start time moved by 30 minutes |
BDM:match 12 (badminton) | Names a new player we do not have |
One more thing happened outside the import: a staff member fixed a golf round in the console and pinned its result. The site still shows the old value, so its row hash is the same as last time.
Today (built).
- The job opens one session, and the fetch reads 500 rows over HTTP inside it (
jobs.py:279,service.py:559-563). - Size check: the usual size is the median of the last 5 written runs, say 500. 500 is not under half of that, so the run goes on (
service.py:805-833). - "Only changes": each row's hash is compared with the newest hash this integration saved for that key (
service.py:835-851). The 4 changed rows differ; 496 rows do not. The golf row has the same hash, so it is skipped. That is fine here, because the staff fix is pinned. But a fix that was not pinned, or a value a writer refused, is never tried again while the site's value stays the same. That is how golf and sailing stayed wrong for days. - The 4 rows go as one
games.import_unitscommand, inside one savepoint, holding the lock on the whole event (engine.py:116). Every console edit on the event waits. - The writer compares field by field and skips pinned fields (
fixtures.py:1141-1150). If the football start is pinned, it is kept. - All 500 rows are written to
integration_row, the 496 unchanged ones too (service.py:513-545). - One commit at the end. If anything failed, the rollback also removes the run record (
jobs.py:294-306).
Agreed design.
- Fetch: 500 rows over HTTP, no transaction open.
- Record:
integration_runrow saved asreceived,rows_collected = 500. Committed. - Map: 500 keys and hashes.
- Guards (I12): any size of run writes. This run would not clear values on more than 10 records, nor reopen more than 5 finished matches. The sailing result is already official, so this integration must have
may_change_officialon. If it is off, I9 opens a review item and the official result stays. - Fast path: one query on
integration_key_statefor the 500 keys. 495 are skipped: same hash, and the record'schanged_seqhas not moved since this integration wrote it. 5 keys go on: the 4 with a new hash, and the golf round, whose hash is the same but whosechanged_seqmoved when staff fixed it. - One batch (5 matches, under the 25 limit), command key
r-8812-run-id:0:- locks the 5 matches, sorted by id;
- reads their current values and pins;
- athletics: status differs, written;
- sailing: result differs, written (the settle window is why this page was still being read 30 hours later);
- football:
startis pinned, so it is kept as set by hand, nothing written; - badminton: unknown player, so the row is held with a question;
- golf: site value differs from ours, but
resultis pinned, so it is kept.
- Saved in that same transaction: 2 field changes, 1 timeline item with the diff, 5 key-state rows, and only 3 rows in
integration_row(2 changed, 1 held). Not 500. - Mark missing: this is a full build, so keys not in the run get
missing_runs + 1. None this time. - Run screen: What changed (2), Kept as set by hand (2), Held (1), Missing (0), Problems (0).
Database time for this run: one fast-path query plus one batch. From the measured numbers below, that is a few milliseconds (estimate). Today's version held the event lock for the whole run.
Example · After the Games: a run of 4 rows
This is the case behind decision I12. During the Asian Games a schedule run usually changed 200 to 500 matches. After the Games ended, the final run sent only 4 rows. Those 4 mattered.
Today (built): the usual size is about 300. 4 is under 10% of 300, so the run counts as a "collapse": it writes nothing and pauses the integration (service.py:120, service.py:597-611). Even without that, if 1 of the 4 rows could not be placed, that is 25%, over the 10% limit, and the whole run is taken back (service.py:860-883).
Agreed design: the 4 rows are saved. An unplaced row is held on its own and never blocks the rest. A run is only held for a person if it would clear values on more than 10 records, reopen more than 5 finished matches, or change an official result it may not change. Size only decides whether a full run may mark rows missing.
Low-level design
This section has two parts: the code as it runs today, quoted from the repo, and the agreed design, which is not in the code yet. Each code block says which one it is. How commands themselves are written in the flows plugin is on Writing a workflow in code.
The tables today Built today
All six tables live in packages/contract/src/omnium_contract/models/integrations.py.
| Table | One row per | Key columns |
|---|---|---|
integration | Scout job wired to one import command | code, event_id, import_command, scout_job_id, secret, state (setting_up, on, paused, paused_by_guard, off), check_seconds (default 1800), build (all, day, pages, ...), last_outcome, last_error |
integration_version | saved version of the settings | number, status (draft, live, archived), settings JSONB with seven stages: source, rows, keep, fields, match, safety, schedule |
integration_run | run | scout_run_id, scout_started_at, via (doorbell, check, run_now, replay, try), status, rows_collected, counts JSONB, guard, message |
integration_row | row of a run | subject_key, subject_kind, subject_id, row_hash (sha256 hex), raw, mapped, outcome, timeline_item_id |
integration_answer | question for a person | kind (held_row, guard_pause, new_values, ...), key, payload, status (open, answered, dismissed); one open answer per key |
integration_job | queued piece of work | kind (run, pull, check, replay, try, ...), params, status (queued, running, done, failed) |
Today's run statuses (integrations.py:119-120): queued, running, written, nothing_new, older, paused_by_guard, failed, tried, cancelled. Row outcomes (integrations.py:175-176): accepted, unchanged, kept_by_hand, refused, not_placed, held, skipped, older, problem.
How a run starts today Built today
Scout's webhook is only a doorbell: it says a run passed and carries ten sample rows, never the full list (scout.py:10-23). omnium checks the signature and queues a pull job. It never fetches inside the web request.
packages/admin/src/omnium_admin/routers/integrations.py:2252-2284 (trimmed)
@doorbell_router.post("/doorbell/{code}")
async def doorbell(code, request, session, x_scout_signature=Header(default=None)):
integration = (await session.execute(
sa.select(Integration).where(Integration.code == code))).scalar_one_or_none()
body = await request.body()
if integration is None or not scout.verify_signature(
body, x_scout_signature, integration.secret
):
raise HTTPException(status.HTTP_401_UNAUTHORIZED, "the signature did not match")
...
if integration.state != "on":
return {"queued": False, "reason": f"the integration is {integration.state}"}
job = await jobs.enqueue(session, integration, "pull", {"via": "doorbell", "run": run.run_id})
return {"queued": True, "job": job.id}
- The signature is HMAC-SHA256 over the exact body, in the
X-Scout-Signatureheader with asha256=prefix. - An integration that is not
ondoes not queue a pull. - If no doorbell arrives, a loop every 60 seconds (
CHECK_SECONDS,scheduler/__main__.py:43andadmin/integration_queue.py:42) queues acheckfor any integration not heard from forcheck_seconds(jobs.py:339-368).
Where the queue is worked. Today two processes work this queue: admin-api, on a thread of its own (admin/integration_queue.py; the integration_queue_in_admin setting is on by default, settings.py:257), and the scheduler (scheduler/__main__.py:131-132). A job is taken under a row lock, so each job runs in only one of them (jobs.py:65-78). In the plan the command service works the queue, in an import process of its own beside the process that answers scorers, at lower CPU priority; admin-api only receives the doorbell (decided 8 Oct 2026; Section 4, L15). A big pull then never shares a process with the console, and the scorer's process always gets the CPU first.
A restart must not block imports for an hour Proposed. Today a job cut off by a restart blocks its integration until it is an hour old (jobs.py:40, 65, 372-377). With imports in the command service, every deploy can cut a run. Proposed: the running job renews a timestamp every 10 s, and a job not renewed for 60 s is freed. The next run continues from the last batch done; batches carry the key {run_id}:{batch_no}, so nothing is written twice.
Fetching from Scout today Built today
packages/core/src/omnium_core/ingest/scout.py:152-181 (trimmed)
async def fetch_items(base_url: str, run_id: str, *, limit: int = 50_000) -> list[Any]:
out: list[Any] = []
offset = 0
async with httpx.AsyncClient(timeout=_TIMEOUT) as client: # 20 s, connect 8 s
while True:
response = await client.get(
f"{_base(base_url)}/v1/runs/{run_id}/items",
params={"offset": str(offset), "limit": str(_PAGE)}, # _PAGE = 500
)
...
out.extend(items)
if len(items) < _PAGE:
return out
offset += len(items)
if len(out) >= limit:
raise ScoutError(f"run {run_id} has more than {limit} rows; refusing to read on")
- Rows come 500 a page, up to 50,000 rows.
- There is no retry. One failed HTTP call fails the job.
- A Scout run that did not pass is refused with an error when it is named (
service.py:306-312). The scheduled check just picks the newest passed run (scout.py:104-150), so failing runs are skipped quietly.
One transaction for the whole job today Built today
packages/core/src/omnium_core/integrations/jobs.py:277-306 (trimmed)
async def work(factory, job_id: int) -> str:
async with factory() as session: # one session, one transaction
job = await session.get(IntegrationJob, job_id)
integration = await session.get(Integration, job.integration_id) if job else None
...
try:
outcome = await perform(session, integration, kind, params) # fetch + map + send + log
status, error = "done", None
except IntegrationError as exc:
outcome, status, error = "failed", "failed", str(exc)
except Exception as exc:
log.exception("integration job %s failed", job_id)
outcome, status, error = "failed", "failed", f"something went wrong: {exc}"
if error is not None:
await session.rollback() # the run record goes too
integration = await session.get(Integration, integration_id)
if integration is not None:
integration.last_outcome = "failed"
integration.last_error = error[:2000]
...
await session.commit()
- The session's first read opens a transaction. The HTTP fetch then happens inside it.
- On any error the rollback removes the
integration_runrow. Only the job's error and the integration'slast_errorstay. - There is no alert: only a log line and a screen update.
Map: key and hash today Built today
packages/core/src/omnium_core/integrations/mapping.py:342-346
parts = [read_path(row.value, key) for key in settings.match.key_fields]
row.key = ":".join(str(part) for part in parts) if all(p is not None for p in parts) else None
row.row_hash = hashlib.sha256(
json.dumps(row.value, sort_keys=True, default=str).encode()
).hexdigest()
The key is built from the version's match.key_fields (default ["externalId"]). The hash is sha256 over the mapped value with sorted keys, so the same value always gives the same hash. The agreed design keeps this step as it is.
Guards today Built today
Two kinds of guard run today. The first looks at size, the second at how many rows failed to place.
packages/core/src/omnium_core/integrations/service.py:805-832 (trimmed)
_COLLAPSE_SHARE = 0.1 # service.py:120
async def _too_few(session, integration, fetched, settings) -> tuple[bool, str | None]:
usual = await _usual_rows(session, integration) # median of the last 5 written runs
if not usual or fetched.rows_collected >= settings.safety.min_share * usual: # min_share 0.5
return False, None
why = f"Scout sent {fetched.rows_collected} rows; it usually sends about {round(usual)}."
return fetched.rows_collected < _COLLAPSE_SHARE * usual, why
packages/core/src/omnium_core/integrations/service.py:860-883 (trimmed)
def _guard(mapped, to_send, sent, settings):
held = sum(1 for row in mapped.rows if row.settled == "held")
considered = len(to_send) + held
...
not_placed = words.count("not_placed") + words.count("held") + held
refused = words.count("refused")
if not_placed / considered > settings.safety.max_not_placed: # default 0.1
return ("not_placed", f"{not_placed} of {considered} rows were not placed or held. "
"Nothing from this run was written.")
if refused / considered > settings.safety.max_refused: # default 0.1
return ("refused", ...)
| Size of run, against the usual | What happens today |
|---|---|
| Under 10% | "Collapse": nothing written, integration paused (paused_by_guard) |
| 10% to 50% | Written, but it may not cancel missing units |
| 50% or more | Normal |
| More than 10% of rows unplaced or refused | Whole run taken back (savepoint rollback), integration paused |
Defaults are in integrations/settings.py:173-181. Section 3 replaces the size rule with I12 (below).
Change detection today Built today
packages/core/src/omnium_core/integrations/service.py:835-851 (trimmed)
async def _changed(session, integration, target, ready, replay_of=None) -> list[Row]:
if not target.max_rows or replay_of is not None:
return ready
last = await _last_hashes(session, integration)
return [row for row in ready if not (row.key and last.get(row.key) == row.row_hash)]
_last_hashesreads the newestrow_hashper key fromintegration_row(service.py:342-365).- It compares with what this integration last held, not with the database. That is the gap Section 3 closes (I3).
- It applies only to imports whose rows stand alone (
max_rowsset). Today that isgames.import_units, up to 1,000 rows per command (UNIT_ROWS_PER_COMMAND,flows/games/commands.py:52). The medal table and entries go in whole, every run. - This query hit the 15 s statement limit on 20 Sep and killed ten checks in a row, before an index was added (
service.py:347-352).
How rows become commands today Built today
packages/core/src/omnium_core/integrations/service.py:422-478 (trimmed)
size = target.max_rows or max(len(rows), 1)
actor = Actor(kind=INTEGRATION, name=f"integration:{integration.code}", source_code=integration.code)
for index, start in enumerate(range(0, len(rows), size)):
chunk = rows[start : start + size]
payload = {
"rows": [{**row.value, "_position": row.position, ...} for row in chunk],
"idSource": settings.match.id_source or integration.code,
}
outcome = await engine.run(
session, workflow=target.workflow, event=target.event, fixture=None,
code=integration.import_command, payload=payload, actor=actor,
idempotency_key=f"{key}:{index}", # key = "<scout run id>:v<version>"
dry_run=dry_run,
)
if outcome.duplicate: # sent before: every row "unchanged"
...
continue
if not outcome.accepted: # refused: every row "refused"
...
continue
for row in chunk:
sent.outcomes[row.position] = outcome.rows.get(row.position) or ("accepted", None, None)
In plain words:
- Rows are cut into chunks of the command's
maxItems(1,000 for units). A medal table or entry list goes as one chunk. - Each chunk is one import command on the event (
fixture=None), sent by an actor of kindINTEGRATION. - The idempotency key is the Scout run id, the version number and the chunk number. Sending the same run twice returns
duplicate, and nothing is written twice. - Every chunk of the run sits inside one savepoint (
service.py:615). A guard can take back the whole run with one rollback. - Because the subject is the event, the engine takes an advisory lock on the whole event (
engine.py:116), held until the job commits. Imports also skip the 10 s time limit a person's command gets (engine.py:110-115,scoring_txn_timeout_msinsettings.py:290); only the 15 s statement limit applies (settings.py:275).
The writer and pins today Built today
The import command's writer compares each field with the database and skips pinned ones.
packages/core/src/omnium_core/ingest/fixtures.py:1141-1150
pinned = set(row.pinned or [])
changed = False
def allowed(name: str, same: bool) -> bool:
if same:
return False
if name in pinned:
report.kept_by_hand += 1
return False
return True
fixture.pinned and fixture_competitor.pinned are JSONB lists of field names a person set by hand: name, start, venue, status, result, competitors on a fixture (models/fixtures.py:127-132). So the writer already compares with the database. Only the step before it, _changed, decides on the last read, and drops rows the writer would have fixed.
Missing units today Built today
After a written run with build all or day (a day build only up to its own Games day), units the run no longer lists are set to CANCELLED, unless they hold a result, a medal or a pin. The guard on which day and discipline a run covered is careful (retire.py:1-60), but the action is a direct write, outside any command.
packages/core/src/omnium_core/ingest/retire.py:279-283
await session.execute(
sa.update(Fixture)
.where(Fixture.id.in_([unit.fixture_id for unit in report.cancelled]))
.values(status=FixtureStatus.CANCELLED)
)
The setting retire_missing_units is on by default (settings.py:162). In one test, 132 of 270 units it called gone were still on the site an hour later (retire.py:30-34). The cancel step catches its own error but has no savepoint, so a database error there aborts the transaction, and the whole good run rolls back (service.py:701-712).
The new table Agreed, to build
Agreed design, not in the code yet. One row per integration and key: what this integration last saw and wrote. It replaces the _last_hashes query.
-- NEW (K1)
CREATE TABLE integration_key_state (
integration_id uuid NOT NULL REFERENCES integration(id) ON DELETE CASCADE,
subject_key text NOT NULL, -- e.g. "ATH:M100METRES---FNL-000100"
subject_kind text NOT NULL, -- fixture, entry, medal_row ...
subject_id uuid, -- our record, once known
last_hash bytea, -- sha256 of the mapped value last seen
applied_seq bigint, -- the record's changed_seq right after this integration wrote it
last_seen_run uuid, -- the last run that listed this key
missing_runs int NOT NULL DEFAULT 0,
official_at timestamptz, -- when the source first marked it official (I8)
next_read_at timestamptz, -- when to read it again in the settle window (I8)
PRIMARY KEY (integration_id, subject_key)
);
CREATE INDEX ix_key_state_next_read ON integration_key_state (next_read_at)
WHERE next_read_at IS NOT NULL;
-- integration_run gains the batch count and Scout's own status.
ALTER TABLE integration_run
ADD COLUMN batch_total int,
ADD COLUMN batch_done int NOT NULL DEFAULT 0,
ADD COLUMN scout_status text; -- passed / failed, as Scout reported it
-- status values: received, applying, done, partly_done, failed, scout_failed, older, paused_by_guard
The fast path Agreed, to build
Agreed design, not in the code yet. Keys that can be skipped: same value as last time, and nobody touched the record since.
SELECT k.subject_key
FROM integration_key_state k
JOIN fixture f ON f.id = k.subject_id
WHERE k.integration_id = $integration
AND k.subject_key = ANY($keys) -- the run's keys, 1,000 at a time
AND k.last_hash = ANY($hashes) -- matched to the same key in code
AND f.changed_seq = k.applied_seq;
Both conditions must hold (K2). Same hash alone is not enough: that is today's bug. Every key the fast path does not skip is compared with the database inside its batch. Measured on a laptop on 6 Oct with 1 million key rows: 1.69 ms median and 3.85 ms at the 99th percentile for 1,000 keys (from the design doc).
One batch, statement by statement Agreed, to build
Agreed design, not in the code yet. Each batch is one import command, sent through the Section 4 road (see Commands).
BEGIN; -- import pool: lock 5 s, query 30 s, idle 30 s
SELECT id, changed_seq FROM fixture WHERE id = ANY($ids) ORDER BY id FOR NO KEY UPDATE; -- up to 25 match locks
SELECT outcome, answer FROM command WHERE command_id = $run || ':' || $batch; -- a repeat? -> duplicate
SELECT ... current values of these matches and their sides ...;
-- in code: diff = site values - current values - pinned fields; empty diff = nothing to write
UPDATE fixture SET <changed fields>, changed_seq = $seq WHERE id = ...; -- only real changes
UPDATE fixture_competitor SET <changed fields>, changed_seq = $seq WHERE id = ...;
INSERT INTO timeline_item (...) VALUES (...); -- one log row for the batch, with the diff
INSERT INTO integration_key_state (...) VALUES (...)
ON CONFLICT (integration_id, subject_key)
DO UPDATE SET last_hash = EXCLUDED.last_hash, applied_seq = EXCLUDED.applied_seq,
last_seen_run = EXCLUDED.last_seen_run, missing_runs = 0;
INSERT INTO integration_row (...) VALUES (...); -- only changed, held and problem rows
INSERT INTO domain_event (...); INSERT INTO command (...);
UPDATE integration_run SET batch_done = batch_done + 1 WHERE id = $run;
COMMIT;
- Locks sorted by id. Two writers on the same matches always lock in the same order, so they cannot deadlock.
- The batch key
run:batchis checked first. A repeat after a crash answersduplicateand writes nothing. - Key state is written in the same transaction as the data. So the fast path can never be ahead of the data.
- Measured with 25 matches and their sides: 1.75 ms median, 3.95 ms at the 99th percentile (laptop, 6 Oct, from the design doc). A 1,240-row schedule run with 50 changed matches is about 2 batches plus the fast path: under 10 ms of database time (estimate in the design doc).

The pipeline as code Agreed, to build
Agreed design, not in the code yet. A new file, integrations/pipeline.py, cuts the run into short transactions around today's mapping and writers.
# core/src/omnium_core/integrations/pipeline.py (new)
async def run(sessions: Sessions, integration_id: uuid.UUID, via: str) -> RunResult:
fetched = await scout.fetch_with_retries(integration) # no transaction open (K4)
async with sessions.short() as s: # transaction 1
run = await records.received(s, integration, fetched) # status "received" (I2)
if fetched.scout_status == "failed":
return await records.scout_failed(s, run, fetched.reason) # alert after 2 in a row (I10)
if await records.is_older(s, integration, fetched):
return await records.older(s, run)
rows = mapping.map_items(fetched.items, version.steps, lookups) # pure, today's code
if (stop := guards.too_few(rows, history)) is not None:
return await records.paused_by_guard(sessions, run, stop)
async with sessions.short() as s:
skip = await keystate.read_skippable(s, integration.id, rows) # K1, K2
todo = [r for r in rows if r.key not in skip]
batches = chunk_by_match(todo, size=25) # sorted by our id
await records.set_batches(sessions, run, len(batches))
for n, batch in enumerate(batches):
answer = await road.run(sessions.imports(), subject=integration.event_subject,
cmd=Command(key=f"{run.id}:{n}", code=integration.import_command,
input={"rows": batch, "runId": str(run.id)}))
if answer.outcome == "try_again":
answer = await retry_batch(...) # 3 tries, growing waits (I5)
if answer.outcome not in ("accepted", "duplicate"):
return await records.partly_done(sessions, run, n, answer) # alert (I10)
async with sessions.short() as s: # last transaction
await missing.mark_missing(s, integration, run, full=fetched.is_full_build) # K5
await records.done(s, run)
await followups.queue(sessions, integration, run) # pages, draws, medallists, settle window (K7)
Reading it top to bottom:
- The fetch runs before any transaction, with 3 retries.
- Transaction 1 saves the run as
received. A Scout-failed run or an older run ends here, recorded. - Mapping is pure code with no database.
- A short read asks the key-state table which keys can be skipped.
- The rest are cut into batches of up to 25 matches, sorted by our id. Each batch is a command keyed
run id:batch number. - A batch that says "try again" is retried on its own. A batch that still fails ends the run as
partly_done, with an alert. The next attempt of the same run starts from the first unsaved batch; saved batches come back asduplicate. - The last short transaction marks missing keys and closes the run.
Missing rows Agreed, to build
Agreed design, not in the code yet. After the last batch, in its own short transaction:
UPDATE integration_key_state SET missing_runs = missing_runs + 1
WHERE integration_id = $integration AND last_seen_run <> $run
AND covered_by_run(subject_key, $run); -- today's day x discipline box (retire.py), kept
-- keys with missing_runs >= 2 and no open item -> insert a "missing" review item (integration_answer)
- Only full runs count. A day build or a page run never marks anything missing.
- Only runs that passed the size check count (a full run must list at least half its usual keys).
- A key seen again resets to 0 (the batch upsert sets
missing_runs = 0). - Cancelling is a person's command,
results.set_status CANCELLED. An integration withauto_cancel_missingon sends the same command itself. It is off by default.

The settle window Agreed, to build
Agreed design, not in the code yet. When a row first says official, next_read_at is set to now plus 1 hour. Each read moves it on:
| Time since official | Read again every |
|---|---|
| 0 to 6 hours | 1 hour |
| 6 to 24 hours | 6 hours |
| 24 to 72 hours | 1 day |
| After 72 hours | not read; next_read_at empty |
The follow-up step picks keys with next_read_at in the past, oldest first, at most 250 per Scout run. That limit comes from Scout's 900 s run limit (PAGES_PER_RUN = 250, followups.py:106-107, built today).
Today a result page is read again only within lookahead_days (3 by default, integrations/settings.py:46) of the unit's start, and at most every 6 hours (PAGE_AGAIN_HOURS, followups.py:110). A protest decided later than that was missed (F17).

Settings Agreed, to build
Agreed design. Every number is a setting, so a new event or provider changes settings, not code.
| Setting | Where | Default |
|---|---|---|
batch_matches | integration settings (version JSON) | 25 |
fetch_retries | integration settings | 3 |
auto_cancel_missing | integration settings | off |
missing_runs_before_review | integration settings | 2 |
settle_hours | integration settings | 72 |
may_change_official | integration settings (Section 2) | off |
row_log_days | settings.py | 30 |
The I12 guard limits (more than 10 records cleared, more than 5 finished matches reopened) are also settings per integration. The design doc does not give them setting names.
Files that change Agreed, to build
| File | New or today | Change |
|---|---|---|
core/.../integrations/jobs.py | today | work() stops wrapping a whole run in one session; it calls pipeline.run() |
core/.../integrations/pipeline.py | new | The stages: fetch, record, map, compare, batches, finish |
core/.../integrations/service.py | today | Keeps guards, mapping glue, held-row questions. _changed and _last_hashes move to keystate.py; _retire_gone moves to missing.py |
core/.../integrations/keystate.py | new | read_skippable(), upsert() on integration_key_state |
core/.../integrations/missing.py | new | mark_missing() after a full run; opens review items (integration_answer, kind missing) |
core/.../integrations/settle.py | new | due_keys() and the next_read_at rules; used by followups.py |
core/.../integrations/retention.py | new | Daily delete of runs older than 30 days, in batches of 10,000 rows |
core/.../ingest/scout.py | today | Fetch gains 3 retries; Scout's run status is kept, not swallowed |
core/.../ingest/retire.py | today | Its coverage query moves into missing.py; the direct UPDATE to CANCELLED is removed |
core/.../ingest/fixtures.py | today | apply_fixtures returns the field-level diff it wrote; skipped unnamed units are no longer reported as accepted |
core/.../catalogue/feed_ids.py | today | assign_feed_ids becomes the writer of a command, games.assign_feed_ids |
Migration order Agreed, to build
- Fix two bugs:
_retire_gonegets a savepoint; skipped unnamed units stop being reported as accepted. - Fetch outside the transaction, with retries; keep Scout's run status.
- Record the run in its own transaction first (I2).
- Create
integration_key_state, filled once from today'sintegration_row(the newest row per key), then kept up to date by the writers. - Add
changed_seqonfixtureandfixture_competitor(Section 4, C3); the fast path turns on. - Batches through
road.runwith the import pool (Section 4, L7). - Missing rows become review items; the direct cancel is removed.
- The settle window on the key;
followups.pyreads due keys. - Row log trimming and the 30-day delete.
- The console: two-line health, the run screen, Waiting for you.
Screens
An operator should answer three questions in five seconds: is each source working, what did the last run change, and is anything waiting for me. Today's console answers the first only partly.
Today Built today, in omnium-console:
- The Integrations page shows one card per integration with a health word: On time, Late, Very late, Failing, Paused, Stopped by a check (
integrationHealth,omnium-console/src/lib/api.ts:461-491). "On time" means only that omnium ran recently: a Scout run that failed its own checks still shows On time. - The Runs page shows each run, its row count, time taken and its log. It says "passed with 1,240 rows", but not what changed in our data, what was kept as set by hand, or what was held.
- Held rows and guard pauses have no single place to act on them.
Agreed Agreed, to build (U1 to U3; built in Section 7):
| Screen | What the operator sees | What they can do | Backend |
|---|---|---|---|
| Integration card | Two lines of health. Scout: "run passed 4 min ago" or "Scout's check failed twice". omnium: "saved 12 changes" or "partly done 7 of 12". One sentence why, when not green. | Run now, Pause, Start, Open last run | Reads games.integrations (gains Scout's status and the batch count). Commands games.run_integration, games.stop_integration, games.start_integration. Live: the INTEGRATION_CHANGED note through the outbox. |
| Run detail | When, which Scout run, rows read, time taken, a bar of batches saved. Five tabs: What changed, Kept as set by hand, Held, Missing from source, Problems. | Open a match; Hand back to the feed; answer a held row; Cancel or Keep a missing unit | Reads games.run (gains the change list: match, field, old, new). Commands results.hand_back, the held-row answer, results.set_status CANCELLED. Live: batch progress through the outbox. |
| Waiting for you | One list across all integrations: held rows, missing units, official results a source wants to change, guard holds. A count on the sidebar. | Answer each in place | New read games.reviews. Each answer is a command, so it is checked and logged like any change. |
Every screen handles: loading, empty ("no runs yet"), failed to read (with Try again), running (bar moving), partly done, paused by a check (with the reason and Start), and off. Nothing shows green unless both Scout and omnium are fine.
When things go wrong
Every case ends with no good data overwritten, nothing saved twice, and a person told when they need to act. This is the agreed design.
| What goes wrong | What happens |
|---|---|
| The site corrects a result two days later (a protest) | The finished unit is still in its settle window, so it is read again. The compare with our data sees the new value and saves it (I3, I8). |
| Another path fixed a match; the site still shows the old value | The fast path sees changed_seq moved and sends the key to the compare. A pinned value is kept. Otherwise the Section 2 priority list decides (I3). |
| Scout is down or slow | The fetch tries 3 times, then the run is recorded as failed and an alert fires. Nothing in our data changes (I1, I10). |
| Scout marks its own run failed | Recorded as "Scout failed" with Scout's reason, shown on the card, alerted after 2 in a row (I1, I10). |
| The server dies at batch 7 of 12 | Batches 1 to 6 are saved; 7 was rolled back. The next attempt sends 7 to 12; 1 to 6 come back as duplicate by their keys (I4, I5). |
| A batch fails because of a bug in a writer | Tried alone 3 times, then the run ends partly_done with the batch's error, and an alert fires. The other batches stay saved (I5, I10). |
| A deploy during a run | The job stops after its current batch. The next job picks up from the first unsaved batch (I5). |
| The site sends half its schedule, or a run after the Games has 4 rows | Any size of run writes its rows. Size only decides whether a full run may mark rows missing (I6, I12). |
| A run would wipe results | Held for a person, nothing written, if it would clear values on more than 10 records, reopen more than 5 finished matches, or change an official result it may not change (I12). |
| A unit really is removed from the schedule | Missing from two full runs in a row: a review item "no longer on the site" with one button to cancel it, as a command (I6). |
| A run names a team we do not have | The row is held with a question; the rest of the run is saved. When the team is added or linked, the row is applied on the next run (I7). |
| A new match with no name | Held as "no name yet", not reported as saved (I7). |
| An import tries to change an official result | Only if this integration may (Section 2 setting). Otherwise a review item, and the official result stays (I9). |
| Two integrations write the same match (schedule and draws) | Both take the match lock. Each compares with what is in the database at that moment, so neither undoes the other unseen (I3, I4). |
| A person edits a match while its batch is saving | Both wait for the match lock, a few milliseconds. The batch respects the person's pin (I4). |
| An HTTP call to Scout hangs for 30 s | No database transaction is open during it, so no lock is held (K4). |
Decisions
All 24 decisions were agreed on 6 Oct 2026. I12 and I13 were added that day after the user's comments.
| # | Decision | In plain words |
|---|---|---|
| I1 | The fetch from Scout happens with no database transaction open, retries 3 times with growing waits, and records a Scout run that failed its checks as "Scout failed". | A slow Scout no longer keeps a database session open, and a Scout failure is no longer skipped quietly. |
| I2 | A run is recorded as "received" in its own short transaction before anything is applied. | A failure never erases the record that a run happened. |
| I3 | Change detection compares each row with what our database holds now, minus pinned fields. A fast path skips the read only when nothing has touched that record since this integration last wrote it. | The golf and sailing fixes that stayed wrong for days are saved on the next run. |
| I4 | Imports save in batches of at most 25 matches; each batch is one import command with key run:batch, through the Section 4 road. | An operator waits at most one short batch, not a whole run. |
| I5 | A failed batch is retried alone, 3 times with growing waits; then the run ends "partly done", and the next attempt starts from the first unsaved batch. | A crash at batch 7 of 12 loses nothing and repeats nothing. |
| I6 | A missing row never cancels anything by itself. Missing from one full run: marked. From two in a row: a review item with Cancel and Keep. Auto-cancel is a setting, off by default. | In one test, 132 of 270 units the code called gone were still on the site. |
| I7 | A row with an unknown id, or a new match with no name, is held with a clear question and retried by itself on every later run. A held row is never reported as saved. | A new team appears: the rest is saved, and its rows land once someone links it. |
| I8 | A finished unit is read again until the source marks it official, then for a settle window (72 hours, a setting): hourly, then every 6 hours, then daily. | A protest decided the next day still reaches our data and our clients. |
| I9 | An import changes an official result only if its integration is allowed to (Section 2 setting); otherwise it opens a review item. | A late site edit never silently overwrites a result we published as official. |
| I10 | Alerts on: a failed run, 2 Scout runs in a row that failed their checks, a batch failed 3 times, a guard hold, no new run for three intervals. All through the central alert service (I13). | We hear about a broken import from an alert, not from a client. |
| I11 | The row log keeps only rows that changed, were held or had a problem, plus one latest row per key; run logs are kept 30 days. | Today every row of every run is saved forever. |
| I12 | Guards judge what a run would do, not how big it is. Any size of run writes; unplaced rows are held one by one. A run is held only if it would clear values on more than 10 records, reopen more than 5 finished matches, or change an official result it may not. Size only gates missing-marking. | After the Games a run had 4 rows, and they mattered. Today that run is paused as a collapse. |
| I13 | One central alert service in our repo. Every part raises an alert the same way, in the same transaction as its cause; the service groups repeats, routes them, shows them in the console, and closes them when the cause clears. | One set of rules and one place to see alerts. Designed in Section 11. |
| U1 | The integration card shows Scout's health and omnium's health separately; green only when both are fine. | "On time" can no longer hide a failing Scout run. |
| U2 | The run screen shows what changed (match, field, old, new), what was kept as set by hand, held, missing, and problems, with batch progress. | An operator sees what a run did in five seconds. |
| U3 | One "Waiting for you" list across all integrations, with a count in the sidebar; every answer is a command. | Nothing that needs a person hides on one integration's page. |
| K1 | A new integration_key_state table: one row per integration and key. | "Did this change?" in about 2 ms for 1,000 keys (measured), not a query that hit 15 s. |
| K2 | The fast path skips a key only when its hash is the same and the record's changed_seq has not moved since this integration wrote it. | Skip only when we can prove nothing changed; otherwise look. |
| K3 | Each batch: lock up to 25 matches sorted by id, check the batch key, read current values, write only real differences, update key state, keep only changed or held rows, count the batch done, commit. | A batch takes about 2 ms (measured) and is safe to repeat. |
| K4 | The fetch runs before any transaction; the run row, each batch, and the missing-row marking are each their own transaction. | No HTTP call ever runs with a transaction open, and no step can erase another's record. |
| K5 | Missing rows are counted per key, only on full runs that passed the size check, using today's day x discipline box; 2 in a row opens a review item. | The proven coverage guard from retire.py stays; only the action changes from cancel to ask. |
| K6 | Cancelling a unit and assigning feed ids become commands (no direct writes). | Every change to our data goes by one road (Section 4, C1). |
| K7 | The settle window lives on the key (next_read_at: 1 h, then 6 h, then 24 h, ending at 72 h). Follow-ups pick due keys, at most 250 per Scout run. | Late corrections are read on a schedule we can see and tune. |
| K8 | integration_row keeps only changed, held and problem rows; a daily job deletes runs older than 30 days in batches of 10,000 rows. | The row log stops growing without limit, and the delete never holds a long lock. |
Built today, or still to build
Checked in the code on 7 Oct 2026, on branch feat/result-feeds (commit 76ef1e2).
| Piece | Today | Agreed design | Status |
|---|---|---|---|
| Doorbell and check | Signed doorbell queues a pull (routers/integrations.py:2252-2284); a check every 60 s for quiet integrations (jobs.py:339-368) | Same | Built today |
| Where the queue is worked | admin-api on a thread of its own, and the scheduler; a row lock gives each job to one of them (admin/integration_queue.py, scheduler/__main__.py:131-132) | The command service, in an import process of its own at lower CPU priority; admin-api only receives the doorbell (L15, changed 8 Oct) | Agreed, to build |
| A run cut off by a restart | Blocks its integration until the job is an hour old (jobs.py:40) | A 60 s lease, renewed every 10 s | Proposed |
| No-code mapping, key and hash | mapping.py:342-346 | Same, unchanged | Built today |
| Pins respected by writers | ingest/fixtures.py:1141-1150 | Same; the compare also leaves pins out | Built today |
| Held rows and questions | integration_answer kind held_row (service.py:886-921); answers like "create a new one" re-read (service.py:947) | Retried by itself every run; never reported as saved | Partly built |
| Fetch outside a transaction, with retries | HTTP fetch inside the job's transaction (jobs.py:279); no retry (scout.py:152-181) | I1, K4 | Agreed, to build |
| Scout-failed runs recorded | Named failed run refused with an error (service.py:306-312); check skips failed runs quietly | I1: saved as "Scout failed", alert after 2 | Agreed, to build |
| Run recorded first | Run row saved at the end; a rollback removes it (jobs.py:294-306) | I2 | Agreed, to build |
| Compare with database minus pins | Compares with this integration's last hash (service.py:835-851) | I3, K1, K2 | Agreed, to build |
integration_key_state and changed_seq | Do not exist (no match in packages/) | K1, Section 4 C3 | Agreed, to build |
Batches of 25 with key run:batch | Chunks of up to 1,000 rows, key scout run:version:chunk, one savepoint, event lock (service.py:422-478, 615; engine.py:116) | I4, I5, K3 | Agreed, to build |
| Guards | Size bands 10% and 50%; 10% unplaced or refused takes the run back (service.py:120, 805-833, 860-883) | I12: judge the effect, not the size | Agreed, to build |
| Missing rows | Direct UPDATE ... CANCELLED (retire.py:279-283), on by default (settings.py:162), no savepoint (service.py:701-712) | I6, K5, K6: review item, cancel is a command | Agreed, to build |
| Settle window | Pages re-read within 3 days of start, at most every 6 h (followups.py:106-110, settings.py:46) | I8, K7: 72 h after official, on the key | Agreed, to build |
| Feed ids | Direct writes after a run (jobs.py:138-152, catalogue/feed_ids.py) | K6: a command, games.assign_feed_ids | Agreed, to build |
| Row log | Every row of every run saved, never deleted (service.py:513-545) | I11, K8 | Agreed, to build |
| Alerts | None for imports; a shared Slack sender exists (notify/slack.py) but integrations do not call it | I10, I13: central alert service (Section 11) | Agreed, to build |
| Console | One health word per card; runs list (omnium-console/src/lib/api.ts:461-491) | U1 to U3 | Agreed, to build |
Numbers
| Number | What | Source |
|---|---|---|
| 1.69 ms median, 3.85 ms p99 | Fast-path query, 1,000 keys, 1 million key rows | Measured on a laptop, 6 Oct (design doc) |
| 1.75 ms median, 3.95 ms p99 | One batch of 25 matches and their sides: lock, read, update, key state, outbox | Measured on a laptop, 6 Oct (design doc) |
| under 10 ms | Database time for a 1,240-row run with 50 changed matches | Estimate from the two numbers above (design doc) |
| 500 rows a page, 50,000 rows max | Scout fetch paging | From the code, scout.py:77, scout.py:152 |
| 20 s, connect 8 s | Scout HTTP timeout | From the code, scout.py:68 |
| 60 s | How often the queue checks for quiet integrations | From the code, scheduler/__main__.py:43 |
| 1,800 s | Default check_seconds per integration | From the code, models/integrations.py:55-57 |
| 1,000 rows | Rows per games.import_units command today | From the code, flows/games/commands.py:52 |
| 10% and 50% | Today's collapse share and min_share | From the code, service.py:120, integrations/settings.py:177 |
| 250 pages, 900 s | Pages per Scout run, and Scout's run limit | From the code, followups.py:106-107 |
| 132 of 270 | Units a day pull called gone that were still on the site | From a code comment, retire.py:30-34 |
| 200 to 500, then 4 | Matches changed per run during the Games, and in the final run | From the user's comment on the design doc |
Read next
- Truth and ownership: the Section 2 priority list that decides who wins when an import and a person disagree.
- Commands: the one road every batch takes, with match locks and time limits.
- Monitoring: the central alert service that import alerts go through.
- Database: the tables, the pools and the limits an import runs under.