Getting data out · Feeds and delivery

Stats, feeds and delivery

How a saved change becomes a client file and a pull answer, and reaches every client on time, even when one client's server is broken.

Design section
Section 9, agreed 7 Oct 2026
Main code
packages/core/.../publishing, api routers/getfeeds.py
Main tables
feed_document, delivery_assignment, delivery_state, feed_dependency (new)
Read time
about 22 minutes

In one minute

Everything a client sees is built from our data: client files sent by S3, FTP, SFTP and webhook, and answers to clients that pull from our API. This is the part clients pay for.

Agreed design. The delivery-worker builds and sends every file. Its build step (Section 9 calls it "the publisher") is woken by the same wake-up as the live screens. It builds files from the main database (never a read copy), and it rebuilds only the files a change touches. A file gets a new version only when its content really changed. Each client destination then has its own queue, its own time limits and retries, and its own proof of delivery. One broken client server delays only itself. Pull clients, such as Google during a cricket match, read from Redis, never straight from Postgres, and need a secret key.

Built today. It carried the Asian Games to paying clients. But every rebuild makes a new version of all 7 files and resends all of them. All clients share one send round, so one hanging SFTP host stalls everyone. No alert reaches a person. Every pull hit reads Postgres four times, and /getfeeds needs no key.

The promise: a saved change reaches the client's file within 10 seconds. An alert fires when a destination is 60 seconds behind or fails 3 times in a row.

10 s
change to client file, target (design doc, D6)
~24 MB
sent per rebuild today, changed or not (design doc, from file sizes)
0.30 ms
Redis write-if-newer, 3.3 KB file (measured, laptop, 7 Oct)
5,000/s
pull hits to design for (estimate, D15)

What this part does

This part takes saved data and gets it to clients. It answers seven questions. Each one went wrong at least once with a paying client at the Asian Games.

#QuestionWhy it matters
1How fresh must each client file and API answer be, and how do we prove it?Clients were promised updates "within seconds"
2When is a file rebuilt, and what is sent when it changes?Every rebuild resent all 7 files to every client, changed or not
3How does one slow or broken client server not delay everyone else?One hanging SFTP host stalled every client's send round
4How do we know each client really received each file?"Did NDTV get the final medal table?" must be answered from a screen
5Where are stats worked out, and what happens when one fails?A failed stat job told no one
6How are new clients and file shapes added?A new client must be settings, not code
7What does pull traffic hit, and how is the database protected?Google sends 2 to 3 million hits during one cricket match

A few words used on this page:

  • Feed file (or file): one built document, such as the calendar or the medal table of one Games. It is a row in feed_document.
  • Destination: one place we send files for one client, such as "NDTV over S3". It is a row in delivery_assignment.
  • Push: we send the file to the client (S3, FTP, SFTP, webhook).
  • Pull: the client asks our API for the file (/getfeeds and the /v1 documents).

How it works

A sketch titled Build once, send per destination. Postgres (main), with its outbox, wakes the delivery-worker's build step. It finds which files to rebuild using feed_dependency, then builds and hashes them. If the hash is the same, it stops. If the content changed, it saves to feed_document. From feed_document, the delivery-worker's send step fans out to four queues: Redis, NDTV S3, News18 FTP and DailyHunt SFTP, each sending to its own client: pull clients, NDTV, News18 and DailyHunt. The DailyHunt arrow is dashed, with a red note: hangs, only this queue waits.
Agreed design: the delivery-worker builds each changed file once; every destination sends on its own.

This is the agreed design. The steps, in order:

  1. A change is saved. A workflow command saves to Postgres and writes outbox rows. (The outbox is a table of "what changed" rows, written in the same transaction as the change.)
  2. The delivery-worker wakes. It uses the same wake-up as the live screens: a Postgres NOTIFY, plus an outbox reader that never skips a late commit. See Live updates.
  3. It finds the touched files. It reads which records changed (a match, a medal standing, an entry), and looks up the files built from them in feed_dependency. Only those files are rebuilt.
  4. It waits a moment. A 1 second timer per file gathers a burst of changes. The timer never runs more than 3 seconds from the first change, so a steady stream cannot delay a file.
  5. It builds each file from the main database. It calls the feed resolvers directly in its own process, not over HTTP through the public API.
  6. It hashes the result. Same hash as the stored file: stop, nothing is saved or sent. A new hash: save the body, the hash, the change time and a new version number, in one transaction.
  7. It queues the new version for every destination. Redis is one more destination. Each destination has its own queue, and it keeps only the newest waiting version of each file.
  8. Each destination sends on its own. It has its own time limits, retries and alerts. A send writes its proof to delivery_state: version, time, size, hash, the receiver's confirmation and the file's age.
  9. Pull clients read Redis. An API server checks the client's secret key in memory, asks Redis for the version, and serves its copy in memory or the new body. Postgres is read only when Redis has nothing.

A worked example

Example · A console fix flips the football final

The values are made up for teaching, except where a source is named. The steps follow the agreed design.

21:14:03.0. Staff in the console correct the football final of a Games. A goal given to one team was scored by the other, so the result flips and gold and silver swap. The command saves the match (fixture row X) and the event's medal standings, and writes outbox rows.

21:14:03.1. The delivery-worker wakes. It looks up the touched records:

SELECT doc_key FROM feed_dependency
 WHERE (record_kind = 'fixture'        AND record_id = :match_x)
    OR (record_kind = 'medal_standing' AND record_id = ANY(:standing_ids));

Two files come back: the medal table and this match's result file. The calendar, team medals and the other files are not touched. Today, all 7 files would be rebuilt.

21:14:04.1. The 1 second timer ends with no new change. The delivery-worker builds both files from the main database. It hashes each body without its time stamp. Both hashes differ from the stored ones, so it saves both files with new versions, say 9120 and 9121, and their changed_at time 21:14:03.0.

21:14:04.2. The two new versions are put on every destination queue that wants them:

DestinationWhat happens
Redis (pull)Write-if-newer, about 0.3 ms per file (measured for a 3.3 KB file). Google's next poll gets the new file
NDTV over S3Sends both files. S3 answers with an ETag. The file reaches NDTV at 21:14:05.0, 2 seconds after the change
News18 over FTPSends both files, one at a time (FTP default). Records the remote file size as its receipt
DailyHunt over SFTPThe host accepts the connection, then hangs

21:15:04. DailyHunt's send reaches its 60 second limit and fails. Its queue waits 5 seconds and tries again. Nothing else waited for it. NDTV, News18 and Redis finished a minute ago.

21:15:04. The oldest file DailyHunt is owed is now over 60 seconds old. An alert goes to the central alert service: "DailyHunt SFTP is 61 s behind". After 3 failures in a row, a second rule fires too.

21:16:30. Staff make one more fix. A new medal-table version 9135 arrives while 9120 still waits in DailyHunt's queue. 9120 is dropped and only 9135 waits. When the host comes back, DailyHunt gets the newest file, not every version in between.

On the client page, an operator sees: NDTV, medal table v9120, sent 21:14:05, ETag confirmed, age 2.0 s. DailyHunt, red, behind since 21:14:04, last error "send timed out after 60 s".

The same change today: the admin process rebuilds all 7 files of the Games, each with a new version number. The delivery worker finds all 7 owed to all 4 destinations, about 24 MB, and sends them in one round of 8 at a time. The round waits for DailyHunt's hanging host, so the next round, for every client, waits too.

Low-level design

Today: how a file reaches a client Built today

Checked in the code on 7 Oct 2026 (branch feat/result-feeds).

  1. A change is saved. A rebuild is started by the in-process seam in the admin service, or by the delivery worker's 60 second sweep (feed_rebuild_seconds, packages/core/src/omnium_core/settings.py:154).
  2. After a 1 second quiet period (feed_debounce_seconds, settings.py:143), the builder rebuilds all 7 feeds of each Games that changed. Each one runs its saved GraphQL query over HTTP against the public API.
  3. Each result is saved to feed_document with a new sequence number, changed or not.
  4. Every 5 seconds (delivery_poll_seconds, settings.py:238) the delivery worker asks the database who is owed a file.
  5. It sends up to 8 at a time (delivery_concurrency, settings.py:237), waits for all of them, and only then goes round again.
  6. Pull clients call /getfeeds, served through a Redis cache kept for 3 seconds, or straight from Postgres when there is no REDIS_URL.

Every rebuild takes a new version number, before the query even runs:

# packages/core/src/omnium_core/publishing/feed.py:292-312 (trimmed)
for definition in definitions:
    try:
        seq = await feed_store.next_seq(session)          # line 295: a new number, always
        body = await self._executor.execute(
            definition.query,
            variables_for(definition, event_code, team),
            definition.root,
        )
        if body is None:
            continue
        payload = _serialize(body)
        await feed_store.upsert(
            session, event=event_feed_id, feed_type=definition.feed_type,
            filter=definition.filter, body=payload, etag=_etag(payload), seq=seq,
        )
        await session.commit()

next_seq draws from the Postgres sequence feed_document_seq. The ETag is a sha256 of the body (feed.py:205-207), so it would show an unchanged file. But nothing compares it before saving or sending.

The store only guards against an older build winning:

# packages/core/src/omnium_core/publishing/feed_store.py:63-84
stmt = pg_insert(FeedDocument).values(
    doc_key=feed_doc_key(event, feed_type, filter),
    event=event, feed_type=feed_type, filter=filter,
    body=body, etag=etag, content_type=content_type, seq=seq,
)
stmt = stmt.on_conflict_do_update(
    index_elements=[FeedDocument.doc_key],
    set_={
        "body": stmt.excluded.body,
        "etag": stmt.excluded.etag,
        "content_type": stmt.excluded.content_type,
        "seq": stmt.excluded.seq,
        "updated_at": sa.func.now(),
    },
    where=stmt.excluded.seq >= FeedDocument.seq,
)
await session.execute(stmt)

In plain words: a slower, older build (a lower seq) cannot overwrite a newer one. That part is right and stays. But an unchanged body with a higher seq always wins, so the file looks new.

The dispatcher decides who is owed what, from one query:

# packages/core/src/omnium_core/publishing/delivery/dispatcher.py:118-147 (trimmed)
"""
  FROM delivery_assignment a
  JOIN delivery_client c   ON c.id = a.client_id AND c.is_active
  JOIN feed_definition def ON def.code = a.feed_key AND def.enabled
  JOIN feed_document d     ON d.feed_type = def.feed_type AND d.filter = def.filter
  ...
  LEFT JOIN delivery_state s ON s.assignment_id = a.id AND s.doc_key = d.doc_key
 WHERE a.is_active
   AND a.delivery_type = ANY(:push_types)
 ORDER BY d.seq
"""
for row in rows:
    last_seq = int(row["last_seq"] or 0)
    changed = last_seq < int(row["seq"])                # line 142
    ...
    schedule = float(row["schedule_seconds"] or 0)
    on_timer = schedule > 0 and (since_delivery is None or since_delivery >= schedule)
    if not (changed or on_timer):
        continue

In plain words: a destination is owed a file when its last sent version is lower than the file's version, or when its schedule timer is due. Two problems follow. A new seq on every rebuild makes every file "changed". And the join at line 121 has no event filter, so every Games' documents match every assignment of that feed type.

The send round waits for everyone:

# packages/core/src/omnium_core/publishing/delivery/runner.py:214-224
jobs = list(self._pending.values())
self._pending.clear()
if not jobs:
    return []
gate = asyncio.Semaphore(max(1, self._settings.delivery_concurrency))

async def one(job: _PendingJob) -> DeliveryResult:
    async with gate:
        return await self.deliver_now(job.spec, job.payload, trigger=job.trigger)

return list(await asyncio.gather(*(one(job) for job in jobs)))

In plain words: up to 8 sends run together, and gather returns only when the last one ends. Retries and their back-off waits happen inside the slot (runner.py:340-368). The worker's main loop is serial (packages/worker/src/omnium_worker/delivery.py:469-495): poll, then rebuild, then flush. So while one round waits for a hanging host, nothing is built, polled or sent.

SFTP and S3 have no time limit set in code:

# packages/core/src/omnium_core/publishing/delivery/transports.py:425-437
def build_s3_transport(settings: Settings) -> S3Transport:
    return S3Transport(
        region=settings.delivery_s3_region,
        endpoint_url=settings.delivery_s3_endpoint_url,
    )

def build_ftp_transport(settings: Settings) -> FtpTransport:
    return FtpTransport(timeout=settings.delivery_transport_timeout)   # 30 s

def build_sftp_transport(settings: Settings) -> SftpTransport:
    return SftpTransport()          # no timeout, and known_hosts=None

FTP gets the 30 second delivery_transport_timeout (settings.py:224). SFTP gets nothing. Its known_hosts=None means host-key checking is off (transports.py:337-343). S3 builds a new client for every send and every retry (transports.py:169-171), with no botocore time-out set.

Nobody is told. The Slack alerter is a seam that only logs:

# packages/core/src/omnium_core/publishing/delivery/deadletter.py:137-162 (trimmed)
class SlackDeadLetterAlerter:
    """Slack alerter **seam** — logs the message it would post; no real call yet."""
    ...
    async def alert(self, entry: DeadLetterEntry) -> None:
        log.warning("slack alert (seam, not sent): %s", self.format_message(entry))

A pull hit reads Postgres four times before the cache:

# packages/api/src/omnium_api/routers/getfeeds.py:159-175 (trimmed)
denied = await _denied_if_unentitled(session, client, feed_type)  # 2 reads: client, entitlement
if denied is not None:
    return denied
definition = await _definition_by_code(session, feed_type)          # read 3: feed name
if definition is None:
    return Response(status_code=404)
event = await resolve_event(session, seriesid)                      # read 4: event
if event is None:
    return Response(status_code=404)
doc_key = feed_doc_key(event.feed_id, definition.feed_type, definition.filter or None)
rendered = await _render_for_client(...)   # cache.get_many, LIVE TTL = 3 s; transform on a miss

And /getfeeds is open without an API key:

# packages/api/src/omnium_api/auth.py:42
_EXEMPT_PATHS = frozenset({"/healthz", "/readyz", "/metrics", "/getfeeds"})

The client name in the address is the only check. Anyone who knows or guesses it gets the feed. The /v1 documents (packages/api/src/omnium_api/routers/documents.py:49-52) read published_document from Postgres on every hit and use no Redis at all.

The gaps, in one table

GapWhat happens todayWhere
Every rebuild resends all 7 filesA new seq per rebuild; the hash is never comparedfeed.py:295; dispatcher.py:142
One hanging host stalls everyoneShared round; no SFTP or S3 time limit; serial worker looptransports.py:436, 169-177; runner.py:224
Nobody is toldThe Slack alerter only logsdeadletter.py:137-162
Builds go through the public APIOnce about 30,000 API calls for one rebuild (now fixed); still under the API's 15 s time-out and 600-a-minute rate limitfeed.py:114-194; api settings.py:43-45
Debounce with no maximum waitA steady stream keeps pushing the rebuild back; feeds were 3 hours behind on 18 Sepfeed.py:494-501; settings.py:143-154
A read copy could serve out-of-date filesGraphQL uses the replica when one is set, and builds go through GraphQLapi app.py:51-60
No event filter when sendingEvery Games' files match every assignment of that feeddispatcher.py:121
A second send path is wrongThe listener matches by feed family: country feeds to the wrong clientsdelivery/listener.py:92
Stats can skip a late commitThe tail reads "id greater than my place"packages/stats/.../triggers.py:69
Failed stats tell no oneAfter 3 tries a stat.failed event is written; no alertpackages/stats/.../worker.py:87, 142-160
"Last updated" can be nowFalls back to the current timeapi graphql/medal_types.py:336
Credentials and hostsSFTP host keys not checked; new S3 client per sendtransports.py:343, 169
Every pull hit reads Postgres4 reads before the cache; /v1 has no cachegetfeeds.py:159-170; documents.py:52
The Redis cache is a 3 s timerNever refreshed on rebuild; off with no REDIS_URLgetfeeds.py:132-137; settings.py:77, 91
Pull data is not protected/getfeeds needs no keyauth.py:42

Whether prod has Redis today: the design doc records "prod has no Redis", not checked in AWS.

The tables Agreed, to build

Two existing tables gain columns, and one table is new. This is agreed design, not in the code yet.

-- Agreed design, not in the code yet.
-- feed_document (exists today) gains the content hash and drops "a new seq every rebuild".
ALTER TABLE feed_document
  ADD COLUMN content_hash bytea,            -- sha256 of the body with the time stamp removed
  ADD COLUMN changed_at  timestamptz;       -- the newest real change inside it ("last updated")
-- seq moves only when content_hash changes.

-- NEW table: which records each built file came from, so a change finds its files.
CREATE TABLE feed_dependency (
  doc_key     text NOT NULL,
  record_kind text NOT NULL,                -- fixture, medal_standing, entry
  record_id   uuid NOT NULL,
  PRIMARY KEY (record_kind, record_id, doc_key)
);

-- delivery_state (exists today) gains the proof of delivery.
ALTER TABLE delivery_state
  ADD COLUMN sent_hash bytea, ADD COLUMN sent_bytes bigint,
  ADD COLUMN receipt   text,                -- S3 ETag, HTTP status, FTP size after upload
  ADD COLUMN age_ms    int;                 -- change to delivery, for the freshness target

When a file is built, its feed_dependency rows are replaced in the same transaction. The primary key starts with (record_kind, record_id), so "which files contain match X?" is one indexed lookup.

What the delivery tables hold today (packages/contract/src/omnium_contract/models/delivery.py):

TableOne row perHolds today
delivery_clientclientcode, name, parent, active
feed_packageproductwhich feed keys it grants
delivery_assignmentclient destination (client, feed, transport, target)delivery_type, destination, credential_ref, throttle_seconds, schedule_seconds, options
delivery_state(assignment, file)last_seq, last_delivered_at, status, failures, last_bytes, last_duration_ms, last_target
delivery_attemptsendoutcome, detail, bytes, duration, trigger; pruned after 14 days
delivery_deadletterfailed versionlast error, attempts, resolved_at
delivery_credentialcredential refthe secret, encrypted

delivery_state already holds the version, time and size of the last send. The agreed columns add the hash, the receiver's confirmation and the age. That turns it into proof of delivery.

The build Agreed, to build

  1. Woken by the Section 8 wake-up, read new outbox rows with the late-commit-safe reader, collect the records they touched.
  2. Look up the files built from those records in feed_dependency. Start a 1 second timer per file, never more than 3 seconds from the first change.
  3. Build each due file in the delivery-worker, from the main database, calling the feed resolvers directly. A direct executor replaces HttpQueryExecutor (feed.py:158).
  4. Hash the body without its time stamp. Same as the stored content_hash: stop. Different: save the body, hash, changed_at and the next seq, and replace the file's dependency rows, in one transaction.
  5. After the commit, put the new version on every destination queue, Redis included.

changed_at is the newest real change inside the file. It is never "now". Today last_updated in the medal table falls back to the current time (medal_types.py:336), so a rebuild with no real change still changes the bytes. That must be fixed, or the hash would always differ.

Cost. The code records the 4.3 MB calendar being built in 1.3 s (a comment at api app.py:157-165). Hashing it takes 1.74 ms (measured, design doc). The 7 Asian Games files total about 6 MB (final copy, 5 Oct). Today all of it goes to each of the 4 client destinations on every rebuild: about 24 MB a rebuild.

The sender: one queue per destination Agreed, to build

Two panels. Left, Today: one yellow box, One shared round, 8 at a time, with arrows down to three clients: NDTV, News18 and DailyHunt. DailyHunt is red and labelled hangs. Under them a red box says: Round waits for all. Right, Agreed: four separate yellow queues, NDTV queue, News18 queue, DailyHunt queue and Redis queue, each with an arrow to its own client. DailyHunt is red and labelled retry later; NDTV, News18 and Pull clients are green. Under them a green box says: Others send on time.
Today one shared round waits for its slowest send. Agreed: each destination has its own queue.

Each destination (a row in delivery_assignment) gets one queue, drained by its own task:

  • Its own pace. A destination sends its files on its own. Destinations run in parallel, up to 32 per worker (first proposal, E3).
  • Coalescing. If version 813 arrives while 812 waits, 812 is dropped and 813 is sent. ("Coalescing" means folding several waiting versions into the newest one.)
  • Time limits on every transport. 10 s to connect, 60 s per send: asyncssh for SFTP, a botocore config for S3, FTP and webhook as today. One S3 client per destination, reused.
  • Proof. After a send, write version, hash, size, receipt and age to delivery_state, in one statement.
  • Back-off on that destination only, then an alert through the central alert service.
  • SFTP checks the host key it was set up with.

The design doc's sketch of the queue:

# Agreed design, not in the code yet.
# core/src/omnium_core/publishing/delivery/destination.py   (new)
class Destination:
    def __init__(self, assignment, transport):
        self.queue: dict[str, Owed] = {}          # doc_key -> newest owed version (coalesced)
        self.wake = asyncio.Event()

    def owe(self, owed: Owed) -> None:
        current = self.queue.get(owed.doc_key)
        if current is None or owed.seq > current.seq:
            self.queue[owed.doc_key] = owed       # only the newest version waits
        self.wake.set()

    async def run(self) -> None:
        while True:
            await self.wake.wait(); self.wake.clear()
            while self.queue:
                doc_key, owed = self.queue.popitem()
                try:
                    receipt = await asyncio.wait_for(self.transport.send(owed), timeout=60)
                    await state.record_sent(owed, receipt)      # version, hash, size, receipt, age
                    self.failures = 0
                except Exception as exc:
                    self.owe(owed)                              # keep it, unless a newer one arrived
                    self.failures += 1
                    if self.failures >= 3:
                        await alerts.raise_("delivery_failing", self.assignment, exc)
                    await asyncio.sleep(backoff(self.failures))

In plain words: owe keeps only the newest version per file. run sends one file at a time for this destination. A failure puts the file back (unless a newer one already came), counts the failure, alerts at 3, and waits before the next try. Nothing here touches any other destination.

What today's runner already has and the new sender keeps: the safety switches (delivery_enabled and the destination fence, runner.py:226-259), per-client transforms, encrypted credentials in delivery_credential, dead letters with replay, and catch-up from delivery_state.

Per-client send settings Agreed, to build

Every destination has these settings, edited in the admin panel. They are new columns on delivery_assignment (E12).

SettingDefaultWhat it does
connect_timeout_s10Give up on connecting after this
send_timeout_s60Give up on one send after this
fast_retries5Retries at 5 s, 15 s, 45 s, 2 min and 5 min
slow_retry_s300After the fast retries, try every 5 min until it lands or a person pauses it
parallel_files1 for FTP and SFTP, 4 for S3 and HTTPHow many files go to this destination at once
alert_after_failures3Failures in a row before an alert
alert_behind_s60How far behind before an alert

Retries never simply stop. A final result must reach the client. fast_retries sets the quick tries. After those, the destination keeps trying every 5 minutes and shows red until a person acts. Today a version is dead-lettered after delivery_max_attempts (5, settings.py:184), and then left alone.

Process-wide settings (agreed; only feed_debounce_seconds exists today):

SettingDefaultToday
feed_debounce_seconds / feed_max_wait_seconds1 / 3debounce exists (1.0); no maximum wait
delivery_connect_timeout / delivery_send_timeout10 / 60delivery_transport_timeout 30, FTP only
delivery_destinations_parallel32delivery_concurrency 8, one shared round
delivery_alert_behind_seconds / delivery_alert_failures60 / 3none

Per destination, the schedule, file name, transform, country and the marks to include stay on delivery_assignment as today.

The pull path: memory, Redis, then Postgres Agreed, to build

A sketch titled A pull hit: memory, then Redis, then Postgres. A pull client such as Google, holding a secret key, sends GET with an ETag to an API server, which keeps a copy in memory. Step 1: the API server asks Redis for the version. Step 2: Redis returns the body only if it changed. Step 3, dashed: only if Redis is empty, the API server reads Postgres (main), at most 4 at once. The API server answers with gzip, or 304. On the right, the delivery-worker first saves each file to Postgres, then writes it to Redis only if newer. A green note: no database read on a normal hit.
Agreed design: a normal pull hit makes no database read.

Pull clients (Google for cricket, apps) read from Redis. There is no CDN: this is our protected data, and every pull needs the client's secret key.

Why Postgres first, then Redis. Postgres and Redis cannot be written in one transaction. If Redis were written first and the Postgres commit then failed, clients would get a file our records say never existed. So Postgres comes first, and Redis right after, retried until it lands. The gap is about 1 ms. During it a pull gets the previous version, which is correct and only a moment old.

Writing (the delivery-worker):

  1. Build the file, hash it, and commit feed_document if the hash changed. Postgres is the record.
  2. After the commit, put the file on the Redis destination's queue. Redis is one more destination, with the same queue, retries and proof as a client (E8).
  3. The Redis job runs the write-if-newer script. An older version is refused, so a retry or a second copy of the delivery-worker can never put an old file back.
  4. A client with its own transform gets its own key, built here at publish time, not on each request.
-- Agreed design, not in the code yet.
-- KEYS[1] = feed:{doc_key} or feed:{client}:{doc_key}
-- ARGV = version (feed_document.seq), etag (content hash), gzip body
local cur = redis.call('HGET', KEYS[1], 'v')
if cur and tonumber(cur) >= tonumber(ARGV[1]) then return 0 end
redis.call('HSET', KEYS[1], 'v', ARGV[1], 'etag', ARGV[2], 'gz', ARGV[3])
return 1

In plain words: read the stored version. If it is the same or newer, do nothing and return 0. Otherwise save version, ETag and gzip body together and return 1. Redis runs a Lua script as one step, so no other write can come in between.

Reading (each API server):

  1. Check the client's key and entitlement in memory. No database read (D13). Clients, keys, entitlements, feed names and events are held in each API server's memory and refreshed when they change.
  2. Ask Redis for the version only. Same as the copy in memory: serve that copy. Changed: fetch the body once and keep it.
  3. Redis has no key: one Postgres read per file per server (single-flight), at most 4 at once per server, then write-if-newer to Redis. ("Single-flight" means many waiting requests share one read instead of each making its own.)
  4. The client already holds this ETag: answer 304 with no body.
# Agreed design, not in the code yet.
# api/src/omnium_api/pull_cache.py   (new)
class PullCache:
    """The last body per key in memory, checked against Redis on each hit."""

    async def get(self, key: str) -> Doc | None:
        mem = self._mem.get(key)
        try:
            v = await self._redis.hget(f"feed:{key}", "v")  # 0.27 ms
        except RedisError:
            # Redis down: serve what we have, marked with its age; else Postgres (limited)
            return mem.with_age_header() if mem else await self._from_db(key)
        if v is None:
            return await self._from_db(key)  # single-flight per key, max 4 per server
        if mem is not None and mem.version == int(v):
            return mem
        v2, etag, gz = await self._redis.hmget(f"feed:{key}", "v", "etag", "gz")
        mem = Doc(version=int(v2), etag=etag, gz=gz)
        self._mem.put(key, mem)  # LRU, capped by PULL_MEMORY_MB
        return mem

In plain words: most hits cost one tiny Redis call that returns a number. The big body crosses the network only when it changed. If Redis is down, the server serves the copy it has, with its age in a header.

The secret key (D14). Each pull client gets a secret key, sent in a header, or in the address for old clients that cannot send headers. The client name alone stops working. A leaked key is revoked in the admin panel. Servers drop it within 1 second by NOTIFY, 60 seconds at worst. More in Section 13 (security).

Gzip and ETag (E11). The body is stored in Redis already gzipped, once. A client that has the current version gets 304 and no body. Today the API already gzips responses over 1 KB (GZipMiddleware, api app.py:166), but it does it again on every request.

Bandwidth is the real limit, not Redis. 5,000 hits a second on a 5 KB file is 25 MB a second (5 KB per cricket scorecard is an estimate). The same rate on the 233 KB gzipped calendar would be 1.2 GB a second. So big files are split per match or per day for pull clients, and 304 answers cost almost nothing.

Redis itself. All files of one Games are about 6 MB, or about 0.5 MB after gzip, so memory is small. Eviction is safe, because a missing key falls back to Postgres, limited by single-flight. Redis keeps nothing we cannot rebuild from Postgres, so it needs no disk copy. A managed Redis with a replica is set up in Section 14.

Public match documents use the same Redis Agreed, to build

The public API's match documents (/v1/fixtures/{id}, standings, careers, medals) are not client files. The write path builds them, inside the scoring, stats and medal transactions (materialize.py:3-6). Today the API reads them from Postgres on every hit (documents.py).

Decided 8 Oct: the delivery-worker also puts them into Redis. It is woken by the same outbox note, reads the new document version from Postgres, and writes it with the same write-if-newer script, as one more destination with the same queue, retries and proof (E8). The publishing-api reads them exactly like client files: memory, then Redis, then Postgres only when Redis has nothing.

So one service writes everything pull clients and fans read into Redis, and the public API never reads the database in normal running.

Freshness and alerts Agreed, to build

The target: 10 seconds from a saved change to the client's file. Build under 3 s, send under 5 s.

  • Each send records the file's age: the time from changed_at to the confirmed send (age_ms).
  • An alert fires when a destination is behind by more than 60 s, or fails 3 times in a row. Alerts go through the central alert service (Section 3, I13). See Monitoring.
  • The test: replay a day of Games data and measure change-to-file age; 95% must be under 10 s.

Stats Partly built

Stats are worked out by the stats engine (packages/stats). It is a queue on Postgres (Procrastinate, packages/stats/src/omnium_stats/queue.py) with three lanes, so a long backfill never sits in front of a live match.

  • Today. The tail reads domain_event rows with "id greater than my place" (triggers.py:65-71). A row from a transaction that commits late can get a lower id and be skipped. A 15-minute sweep repairs it. A failed stat job is retried 3 times (worker.py:87), then writes one stat.failed event (worker.py:142-160). Only the Runs screen shows it.
  • Agreed (D7). The stats queue reads the outbox with the shared late-commit-safe reader (Section 4, L9). A failed stat job raises an alert that names the stat and the match.

A sport's own stats are written in the flows plugin. See How workflow code is written.

Clients are settings Agreed, to build

Client, package, destination, schedule, file names, transform and which country a country feed is about are all edited in the admin panel (D9). One transform setting applies to both pull and push. Credentials are typed into the admin panel (encrypted, delivery_credential) or held in Secrets Manager, never in an environment file. Each file can carry its version number and, per unit, whether the result is provisional or official; each client chooses (D10).

Screens

The admin panel already has a delivery status page (frontend-admin/src/spa/pages/DeliveryStatusPage.tsx) and dead letters with replay. It shows rows by assignment. The agreed screens answer, per client, the question a client asks on the phone:

ScreenWhoWhat they seeWhat they can do
Clients overviewOperatorOne row per client: green, behind (by how long), failing, or switched off; the oldest file ageOpen a client
One clientOperatorEach destination and file: version, sent at, size, confirmation, age when sent; a chart of age over the last hourSend again now; pause a destination; test the connection
FreshnessAdminFor every client file, change-to-delivery time over the day: median and worstFind the slowest destination
Feed settingsAdminClient, package, destinations, file names, transform, country, schedule, marks to includeAdd a client with no code; preview the file

Every switch-off asks for a reason and shows who did it and when. The status page lists all switched-off clients at the top.

Order of the work

From the design doc, in this order (stop the stalls first):

  1. Time limits on SFTP and S3, and one S3 client per destination.
  2. One task per destination instead of one shared round.
  3. content_hash on feed_document; send only on a change; last_updated from real changes.
  4. Remove the listener send path; add the event filter to the dispatcher.
  5. Proof of delivery and age in delivery_state; the client screens.
  6. Alerts through the central service.
  7. The delivery-worker builds on the Section 8 wake-up instead of its 60 s sweep; builds move out of the admin service; direct executor instead of HTTP.
  8. feed_dependency and building only touched files.
  9. The stats reader fix; version and provisional or official marks per client.

The Redis pull path, per-client settings and secret keys (D11 to D15, E7 to E12) were added on 7 Oct and are not yet placed in this order in the design doc.

Tests that prove it

TestProves
An SFTP test host that accepts and then hangsIt times out at 60 s; other destinations deliver on time (D4)
Rebuild with no real changeNo new version; nothing sent (D3)
50 changes in 2 secondsEach file is built at most once per 3 s; each destination gets only the newest (D2, E3)
Kill the sender mid-uploadThe same version is sent again; the receipt is recorded once (E5)
A host down for 10 minutes, then backIt receives the newest version of each file, not every version; the alert fired at 60 s (E3, E6)
Change one match on a 7,000-unit GamesOnly the files containing it are rebuilt (E2)
Replay a day of Games data95% of change-to-file ages under 10 s (D6)
Load test pull (Section 12)The full HTTP path holds 5,000 hits a second (D15); not run yet

When things go wrong

What goes wrongWhat happens (agreed design)
DailyHunt's SFTP host hangsIts send stops at the 60 s limit; its queue backs off and alerts; NDTV and News18 keep receiving on time
A rebuild produces the same bytesSame hash, no new version, nothing sent to anyone
50 matches finish within a minuteRebuilds gather within 1 s and never wait more than 3 s; each destination sends only the newest version of each file
The delivery-worker crashes mid-buildNothing half built is saved (each file is written whole with its hash); on restart it reads the outbox from its last place and builds again
The delivery worker crashes mid-sendThe send is retried with the same version; the receiver gets the same file again, which is harmless
A client says "we never got the final table"The client page shows: calendar v812, sent 12:04:31, 3.29 MB, hash, S3 ETag confirmed. Or the failure and since when
A client host is down for an hourIts queue keeps only the newest version of each file; when it returns it gets the newest file, not 60 old ones; the alert fired after 60 s
The database is slowBuilds wait; push clients keep the last file they got; pull clients keep reading the last good file from Redis and server memory
A read copy of the database lagsThe delivery-worker never builds from a read copy
Redis is downAPI servers serve the copy in memory with its age in a header; misses read Postgres, at most 4 at once per server; an alert fires
The Redis write fails after the Postgres commitThe Redis destination retries; pulls get the previous version meanwhile; an alert at 3 failures
Two Redis writes arrive out of orderWrite-if-newer refuses the older one
Redis restarts empty, or evicts keysThe first hit per file per server reads Postgres once and refills Redis
An API server restartsIts memory is empty; the first hits read Redis
A client key leaksRevoked in the admin panel; servers drop it within 1 s by NOTIFY, 60 s at worst
A stat job fails 3 timesIt stays failed and an alert names the stat and the match
A client needs a different country's feedA setting on their destination, no code
A client is switched off on purposeIt is listed as switched off on the status page and in every update until it is back on

Decisions

All decisions below were agreed on 7 Oct 2026. D11 to D15 and E7 to E12 came from the review comments.

#DecisionIn plain words
D1One build path: the delivery-worker, woken by the Section 8 wake-up, builds every file from the main database by calling the resolvers directly. Section 9 named this "one publisher service"; on 7 Oct 2026 it was placed inside the delivery-worker, not in a container of its ownOne build path instead of two; never built from an out-of-date copy; no 30,000 API calls
D2Each feed declares the records it is built from; a change rebuilds only the feeds that contain it; wait 1 s for a burst, never more than 3 sA change to one match does not rebuild unrelated files, and a stream of changes cannot delay a file
D3A new version only when the file's content hash changes; "last updated" is the newest real change, never nowUnchanged files are never sent again
D4Each client destination has its own queue; every transport has a 10 s connect and 60 s send limit; only the newest waiting version is sentOne hanging host delays only itself
D5Proof of delivery per destination and file: version, time, size, hash and the receiver's confirmation"Did NDTV get the final table?" is answered on a screen
D6Freshness target 10 s from a saved change to the client's file; each send records the file's age; alert at 60 s behind or 3 failures in a row"Within seconds" becomes a number we watch
D7The stats queue reads the outbox without skipping late commits, and a failed stat job raises an alertNo stat is silently missing
D8Pull clients and the public API read from Redis first and Postgres only on a miss; no CDN, because the data is protected; nothing is built on requestTraffic never slows scoring, and clients keep reading during an outage
D9Clients are settings: package, destinations, names, transform (one for pull and push), country, schedule, credentials in the admin panel or Secrets Manager; switch-offs need a reason and stay visibleA new client is set up in the admin panel the same day
D10Each file can carry its version number and, per unit, provisional or official; each client choosesClients can check they have the latest and the final result
D11The delivery-worker writes each file to Postgres first, then to Redis with write-if-newer; a pull reads Redis, and Postgres only when Redis has nothingPull clients read from memory; an old file can never replace a newer one
D12Each destination has its own settings: connect and send limits, fast retries, files at once, alert levels; all destinations send in parallelA slow or strict client is tuned without touching anyone else
D13A pull hit makes zero Postgres reads: clients, keys, entitlements, feed names and events are held in each API server's memory and refreshed when they change3 million hits in a match put no load on the scoring database
D14Pull needs a secret key per client, in a header (or in the address for old clients); the client name alone stops workingProtected data is no longer open to anyone who guesses a client name
D15Design and load-test pull for 5,000 hits a second at peakA number we test against before a series goes live
E1feed_document gains content_hash and changed_at; seq moves only when the hash changesNo new version means nothing to send
E2A new feed_dependency table, replaced with each build; a change finds its files by one indexed lookupOnly the files a change touches are rebuilt
E3One queue and task per destination, coalescing to the newest version; up to 32 destinations at once per workerA slow destination never blocks another; a returning host gets the newest file only
E410 s connect and 60 s send limits on all four transports; one reused S3 client per destination; SFTP host keys checkedNo send can hang forever, and memory stays flat
E5delivery_state keeps version, hash, size, receipt and age for each destination and fileProof of delivery and freshness come from one row
E6Per-destination back-off; alerts at 3 failures or 60 s behind, through the central alert serviceA failing client is retried calmly and a person is told
E7Redis keys feed:{doc_key}, plus feed:{client}:{doc_key} for a client with its own transform, hold version, ETag and gzip body; a Lua script writes only if newerOne small atomic write per changed file
E8Redis is written by the delivery-worker as one more destination (own queue, retries, proof), after the Postgres commitIf Redis is down, Postgres still has the file, and Redis is filled when it returns
E9Each API server keeps the last body per key in memory and asks Redis only for the versionBig files are not copied over the network on every hit
E10On a Redis miss, one request per file per server reads Postgres, at most 4 at once; when Redis is down, serve the copy we have with its ageRedis trouble cannot flood Postgres and slow scoring
E11Answers are gzip with an ETag; a client with the current version gets 304 and no bodyFiles are 8 to 15 times smaller, and repeat polls cost almost nothing
E12delivery_assignment gains connect_timeout_s, send_timeout_s, fast_retries, slow_retry_s, parallel_files, alert_after_failures, alert_behind_sThe per-client send settings, edited in the admin panel

Built today, or still to build

PieceToday (checked 7 Oct)Agreed designStatus
Feed definitions as saved queriesfeed_definition rows with a GraphQL query eachKeptBuilt today
Who buildsAdmin service (same-process seam) and the delivery worker's 60 s sweep (feed.py:569-589; worker delivery.py:489-492)The delivery-worker alone, on the Section 8 wake-upAgreed, to build
How it buildsOver HTTP through the public API (feed.py:158-194)Resolvers called directly, main database onlyAgreed, to build
Which files are rebuiltAll 7 of each changed Games (feed.py:351, 386-437)Only files in feed_dependency for the touched recordsAgreed, to build
Debounce1 s, no maximum (feed.py:494-501)1 s, maximum 3 sPartly built
Version numberNew seq every rebuild (feed.py:295)New seq only when content_hash changesAgreed, to build
Older build cannot winseq guard in the upsert (feed_store.py:73-83)KeptBuilt today
Who is owed whatDispatcher query every 5 s (dispatcher.py:76-192); no event filter (:121)Per-destination queues fed by the build step; event filterPartly built
Second send pathListener matches by feed family (listener.py:92)RemovedAgreed, to build
Public match documents (/v1)Built in the write path (materialize.py); read from Postgres on every hit (documents.py)The delivery-worker also writes them into Redis; the API reads Redis first (8 Oct)Agreed, to build
SendingOne shared round, 8 at a time, gather (runner.py:214-224)One task per destination, up to 32 per workerAgreed, to build
CoalescingPer throttle key, inside one round (runner.py:185-193)Per destination and file, newest onlyPartly built
Time limitsFTP 30 s, webhook 10 s; SFTP and S3 none (transports.py:425-437)10 s connect, 60 s send, every transport, per clientPartly built
S3 clientNew one per send (transports.py:169-171)One per destination, reusedAgreed, to build
SFTP host keysNot checked (transports.py:343)CheckedAgreed, to build
Retries5 tries, 0.5 s doubling to 30 s, then dead letter (settings.py:184-186)5 fast retries, then every 5 min until it landsAgreed, to build
Safety switchesdelivery_enabled and destination fence (runner.py:226-259)KeptBuilt today
Proof of deliverydelivery_state with seq, time, size, status (delivery.py:197-253 in contract models)Plus hash, receipt, agePartly built
AlertsSlack seam only logs (deadletter.py:137-162)Central alert service, 60 s behind or 3 failuresAgreed, to build
Encrypted credentialsdelivery_credential tableKept, or Secrets ManagerBuilt today
Pull reads4 Postgres reads per hit, then a 3 s Redis cache if REDIS_URL is set (getfeeds.py:159-186)Memory and Redis; Postgres only on a miss, single-flightAgreed, to build
Pull security/getfeeds exempt from the API key (auth.py:42)Secret key per clientAgreed, to build
Gzip and 304Gzip per request over 1 KB (app.py:166); ETag and 304 on /getfeeds (getfeeds.py:188-196)Gzip stored once in Redis; ETag and 304 keptPartly built
Stats tail"id greater than my place", 15 min sweep (triggers.py:65-71)Late-commit-safe readerAgreed, to build
Failed statsstat.failed event, no alert (worker.py:142-160)Alert naming the stat and matchPartly built
Client screensDeliveryStatusPage.tsx, rows by assignmentClients overview, one client, freshnessPartly built

Numbers

NumberWhatSource
0.30 / 0.43 msRedis write-if-newer, 3.3 KB file, median / 99th percentileMeasured: laptop, Redis 7 in Docker, Python client with 50 connections, 7 Oct
0.67 / 2.05 msSame, 44 KB fileSame
1.95 / 4.53 msSame, 233 KB fileSame
0.27 to 0.30 msRedis version check, median, all three sizesSame
0.29 / 1.07 msFull Redis read, median, 3.3 KB / 233 KBSame
11,754 / 4,280Redis reads a second, one Python process, 3.3 KB / 233 KBSame
8 to 15 timesGzip size cut: calendar 3.29 MB to 233 KB, result details 586 KB to 44 KB, top-5 medals 42 KB to 3.3 KB, team medals 16 KB to 2 KBMeasured on the final Asian Games files
1.74 mssha256 of a 4.3 MB fileMeasured, design doc
1.3 sBuilding the 4.3 MB calendarA comment in the code, api app.py:157-165
about 6 MBAll 7 Asian Games filesFinal copy, 5 Oct
about 24 MBSent per rebuild today to 4 destinationsWorked out from the file sizes
5,000 a secondPull design targetEstimate: 3 million hits over a 3.5-hour T20 is about 240 a second; peaks assumed 10 times that; then doubled
25 MB a second5,000 hits a second on a 5 KB fileEstimate (5 KB per scorecard is an estimate)

These are Redis calls only. The full HTTP request through one API process is not measured yet. That is the Section 12 load test against D15.

  • Live updates: the wake-up that starts the delivery-worker's build, and the late-commit-safe outbox reader.
  • The database: the main database, the read-only copy, and why feeds never read the copy.
  • Monitoring: the central alert service that delivery alerts go through.
  • How workflow code is written: where a sport's stats are written.