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.

Design section
Section 3, agreed 6 Oct 2026
Main code
packages/core/.../integrations/ (service.py, jobs.py)
Main tables
integration, integration_run, integration_row, integration_answer, integration_job
Read time
about 25 minutes

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.

1.69 ms
fast path, 1,000 keys in 1 million rows (median, laptop, 6 Oct)
1.75 ms
one batch of 25 matches (median, laptop, 6 Oct)
25
matches per batch, at most (agreed, a setting)
72 h
settle window after a result is official (agreed, a setting)

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:

WordMeaning
IntegrationOne Scout job connected to one import command, set up with no code. A row in integration.
RunOne read of one Scout run. A row in integration_run.
RowOne 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 commandA workflow command marked imports=True, for example games.import_units. It writes the rows.
PinA field name a person set by hand, saved in fixture.pinned. An import never overwrites a pinned field.
Held rowA 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:

#QuestionWhat went wrong at the Games
1How 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).
2What 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.
3What does a row missing from a run mean?A partial list looked like "everything else was deleted".
4What happens with an id we do not know?A new player or team must not be dropped or invented silently.
5How 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.
6How long does an import keep reading a finished match?Results changed after protests a day or more later, after we stopped reading.
7How do we see that an import is failing?A job failed Scout's checks for days and nobody was told (N4).

How it works

A sketch of one import run in the agreed design. On the far left a white box Scout sends an arrow labelled doorbell or check to the main column. The main column has six yellow boxes, top to bottom: 1. Fetch rows, no transaction open; 2. Record run, status received; 3. Map rows, key plus hash; 4. Compare, database minus pins; 5. Save in batches, 25 matches each; 6. Mark missing, full runs only. Below them a green box Run done. A blue cylinder key state has a dashed arrow labelled skip unchanged into Compare. Save in batches has an arrow labelled one command per batch into a blue cylinder Postgres. Three red boxes on the right: Held rows, from Map rows, labelled unknown id; Failed batch, from Save in batches, labelled after 3 tries; Missing twice, from Mark missing. All three red boxes point to a white box Waiting for you. Failed batch also points to a yellow box Alert service.
One run in the agreed design. Good rows go down the middle; anything that cannot be saved safely goes right, to a person.

The agreed run has six steps. What runs today is described step by step in the low-level design below.

  1. 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.
  2. Record the run first. In its own short transaction, the run is saved with status received and its row count. From then on, no failure can erase the record that the run happened.
  3. Map. Today's no-code steps, unchanged. Each row gets a key and a hash.
  4. 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.
  5. 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.
  6. 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-000100Status 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).

  1. The job opens one session, and the fetch reads 500 rows over HTTP inside it (jobs.py:279, service.py:559-563).
  2. 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).
  3. "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.
  4. The 4 rows go as one games.import_units command, inside one savepoint, holding the lock on the whole event (engine.py:116). Every console edit on the event waits.
  5. The writer compares field by field and skips pinned fields (fixtures.py:1141-1150). If the football start is pinned, it is kept.
  6. All 500 rows are written to integration_row, the 496 unchanged ones too (service.py:513-545).
  7. One commit at the end. If anything failed, the rollback also removes the run record (jobs.py:294-306).

Agreed design.

  1. Fetch: 500 rows over HTTP, no transaction open.
  2. Record: integration_run row saved as received, rows_collected = 500. Committed.
  3. Map: 500 keys and hashes.
  4. 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_official on. If it is off, I9 opens a review item and the official result stays.
  5. Fast path: one query on integration_key_state for the 500 keys. 495 are skipped: same hash, and the record's changed_seq has 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 whose changed_seq moved when staff fixed it.
  6. 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: start is 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 result is pinned, so it is kept.
  7. 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.
  8. Mark missing: this is a full build, so keys not in the run get missing_runs + 1. None this time.
  9. 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.

TableOne row perKey columns
integrationScout job wired to one import commandcode, 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_versionsaved version of the settingsnumber, status (draft, live, archived), settings JSONB with seven stages: source, rows, keep, fields, match, safety, schedule
integration_runrunscout_run_id, scout_started_at, via (doorbell, check, run_now, replay, try), status, rows_collected, counts JSONB, guard, message
integration_rowrow of a runsubject_key, subject_kind, subject_id, row_hash (sha256 hex), raw, mapped, outcome, timeline_item_id
integration_answerquestion for a personkind (held_row, guard_pause, new_values, ...), key, payload, status (open, answered, dismissed); one open answer per key
integration_jobqueued piece of workkind (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-Signature header with a sha256= prefix.
  • An integration that is not on does not queue a pull.
  • If no doorbell arrives, a loop every 60 seconds (CHECK_SECONDS, scheduler/__main__.py:43 and admin/integration_queue.py:42) queues a check for any integration not heard from for check_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_run row. Only the job's error and the integration's last_error stay.
  • 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 usualWhat 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 moreNormal
More than 10% of rows unplaced or refusedWhole 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_hashes reads the newest row_hash per key from integration_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_rows set). Today that is games.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 kind INTEGRATION.
  • 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_ms in settings.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:batch is checked first. A repeat after a crash answers duplicate and 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).
Two time lines. Top, labelled Today in red: one long red bar labelled one transaction, event locked, containing fetch (HTTP), map, send all rows, row log, cancel missing and feed ids, ending in a single commit mark. Bottom, labelled Agreed in green: a white box fetch (HTTP) marked no transaction, a short yellow bar record run, a white box map plus fast path, four short yellow bars each labelled batch with the note 25 matches, own commit, a yellow bar mark missing, and a green box safe to repeat.
Today one transaction wraps the whole run. In the agreed design each yellow step commits on its own, and the HTTP fetch has no transaction open.

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:

  1. The fetch runs before any transaction, with 3 retries.
  2. Transaction 1 saves the run as received. A Scout-failed run or an older run ends here, recorded.
  3. Mapping is pure code with no database.
  4. A short read asks the key-state table which keys can be skipped.
  5. 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.
  6. 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 as duplicate.
  7. 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 with auto_cancel_missing on sends the same command itself. It is off by default.
A state sketch. A green box Listed has an arrow labelled not in a full run to a yellow box Missing once, marked, nothing changes. A curved arrow labelled seen again goes from Missing once back to Listed. Missing once has an arrow labelled not in next full run to a red box Missing twice, review item. Missing twice points to a white box A person decides, which has two arrows down: to a yellow box Cancel, a command, and to a green box Keep.
Missing is a question, not a delete. Only a person (or an opt-in setting) cancels.

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 officialRead again every
0 to 6 hours1 hour
6 to 24 hours6 hours
24 to 72 hours1 day
After 72 hoursnot 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).

A time line from official, through 6 h and 24 h, to 72 h, then a dashed part labelled stop reading. Above the line, yellow bands: read every hour from official to 6 h, read every 6 hours from 6 h to 24 h, read once a day from 24 h to 72 h, with dots for each read. Below, a red box protest result changes points up to a read between 24 h and 72 h, and a green box beside it says saved on next read. A note at the bottom: next_read_at on the key, 250 pages per Scout run.
A finished result is read on a falling schedule for 72 hours, so a late protest still reaches our data.

Settings Agreed, to build

Agreed design. Every number is a setting, so a new event or provider changes settings, not code.

SettingWhereDefault
batch_matchesintegration settings (version JSON)25
fetch_retriesintegration settings3
auto_cancel_missingintegration settingsoff
missing_runs_before_reviewintegration settings2
settle_hoursintegration settings72
may_change_officialintegration settings (Section 2)off
row_log_dayssettings.py30

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

FileNew or todayChange
core/.../integrations/jobs.pytodaywork() stops wrapping a whole run in one session; it calls pipeline.run()
core/.../integrations/pipeline.pynewThe stages: fetch, record, map, compare, batches, finish
core/.../integrations/service.pytodayKeeps guards, mapping glue, held-row questions. _changed and _last_hashes move to keystate.py; _retire_gone moves to missing.py
core/.../integrations/keystate.pynewread_skippable(), upsert() on integration_key_state
core/.../integrations/missing.pynewmark_missing() after a full run; opens review items (integration_answer, kind missing)
core/.../integrations/settle.pynewdue_keys() and the next_read_at rules; used by followups.py
core/.../integrations/retention.pynewDaily delete of runs older than 30 days, in batches of 10,000 rows
core/.../ingest/scout.pytodayFetch gains 3 retries; Scout's run status is kept, not swallowed
core/.../ingest/retire.pytodayIts coverage query moves into missing.py; the direct UPDATE to CANCELLED is removed
core/.../ingest/fixtures.pytodayapply_fixtures returns the field-level diff it wrote; skipped unnamed units are no longer reported as accepted
core/.../catalogue/feed_ids.pytodayassign_feed_ids becomes the writer of a command, games.assign_feed_ids

Migration order Agreed, to build

  1. Fix two bugs: _retire_gone gets a savepoint; skipped unnamed units stop being reported as accepted.
  2. Fetch outside the transaction, with retries; keep Scout's run status.
  3. Record the run in its own transaction first (I2).
  4. Create integration_key_state, filled once from today's integration_row (the newest row per key), then kept up to date by the writers.
  5. Add changed_seq on fixture and fixture_competitor (Section 4, C3); the fast path turns on.
  6. Batches through road.run with the import pool (Section 4, L7).
  7. Missing rows become review items; the direct cancel is removed.
  8. The settle window on the key; followups.py reads due keys.
  9. Row log trimming and the 30-day delete.
  10. 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):

ScreenWhat the operator seesWhat they can doBackend
Integration cardTwo 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 runReads 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 detailWhen, 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 unitReads 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 youOne 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 placeNew 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 wrongWhat 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 valueThe 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 slowThe 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 failedRecorded 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 12Batches 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 writerTried 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 runThe 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 rowsAny size of run writes its rows. Size only decides whether a full run may mark rows missing (I6, I12).
A run would wipe resultsHeld 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 scheduleMissing 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 haveThe 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 nameHeld as "no name yet", not reported as saved (I7).
An import tries to change an official resultOnly 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 savingBoth wait for the match lock, a few milliseconds. The batch respects the person's pin (I4).
An HTTP call to Scout hangs for 30 sNo 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.

#DecisionIn plain words
I1The 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.
I2A run is recorded as "received" in its own short transaction before anything is applied.A failure never erases the record that a run happened.
I3Change 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.
I4Imports 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.
I5A 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.
I6A 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.
I7A 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.
I8A 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.
I9An 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.
I10Alerts 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.
I11The 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.
I12Guards 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.
I13One 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.
U1The 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.
U2The 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.
U3One "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.
K1A 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.
K2The 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.
K3Each 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.
K4The 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.
K5Missing 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.
K6Cancelling a unit and assigning feed ids become commands (no direct writes).Every change to our data goes by one road (Section 4, C1).
K7The 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.
K8integration_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).

PieceTodayAgreed designStatus
Doorbell and checkSigned doorbell queues a pull (routers/integrations.py:2252-2284); a check every 60 s for quiet integrations (jobs.py:339-368)SameBuilt today
Where the queue is workedadmin-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 restartBlocks its integration until the job is an hour old (jobs.py:40)A 60 s lease, renewed every 10 sProposed
No-code mapping, key and hashmapping.py:342-346Same, unchangedBuilt today
Pins respected by writersingest/fixtures.py:1141-1150Same; the compare also leaves pins outBuilt today
Held rows and questionsintegration_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 savedPartly built
Fetch outside a transaction, with retriesHTTP fetch inside the job's transaction (jobs.py:279); no retry (scout.py:152-181)I1, K4Agreed, to build
Scout-failed runs recordedNamed failed run refused with an error (service.py:306-312); check skips failed runs quietlyI1: saved as "Scout failed", alert after 2Agreed, to build
Run recorded firstRun row saved at the end; a rollback removes it (jobs.py:294-306)I2Agreed, to build
Compare with database minus pinsCompares with this integration's last hash (service.py:835-851)I3, K1, K2Agreed, to build
integration_key_state and changed_seqDo not exist (no match in packages/)K1, Section 4 C3Agreed, to build
Batches of 25 with key run:batchChunks 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, K3Agreed, to build
GuardsSize 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 sizeAgreed, to build
Missing rowsDirect 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 commandAgreed, to build
Settle windowPages 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 keyAgreed, to build
Feed idsDirect writes after a run (jobs.py:138-152, catalogue/feed_ids.py)K6: a command, games.assign_feed_idsAgreed, to build
Row logEvery row of every run saved, never deleted (service.py:513-545)I11, K8Agreed, to build
AlertsNone for imports; a shared Slack sender exists (notify/slack.py) but integrations do not call itI10, I13: central alert service (Section 11)Agreed, to build
ConsoleOne health word per card; runs list (omnium-console/src/lib/api.ts:461-491)U1 to U3Agreed, to build

Numbers

NumberWhatSource
1.69 ms median, 3.85 ms p99Fast-path query, 1,000 keys, 1 million key rowsMeasured on a laptop, 6 Oct (design doc)
1.75 ms median, 3.95 ms p99One batch of 25 matches and their sides: lock, read, update, key state, outboxMeasured on a laptop, 6 Oct (design doc)
under 10 msDatabase time for a 1,240-row run with 50 changed matchesEstimate from the two numbers above (design doc)
500 rows a page, 50,000 rows maxScout fetch pagingFrom the code, scout.py:77, scout.py:152
20 s, connect 8 sScout HTTP timeoutFrom the code, scout.py:68
60 sHow often the queue checks for quiet integrationsFrom the code, scheduler/__main__.py:43
1,800 sDefault check_seconds per integrationFrom the code, models/integrations.py:55-57
1,000 rowsRows per games.import_units command todayFrom the code, flows/games/commands.py:52
10% and 50%Today's collapse share and min_shareFrom the code, service.py:120, integrations/settings.py:177
250 pages, 900 sPages per Scout run, and Scout's run limitFrom the code, followups.py:106-107
132 of 270Units a day pull called gone that were still on the siteFrom a code comment, retire.py:30-34
200 to 500, then 4Matches changed per run during the Games, and in the final runFrom the user's comment on the design doc
  • 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.