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.
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.
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.
| # | Question | Why it matters |
|---|---|---|
| 1 | How fresh must each client file and API answer be, and how do we prove it? | Clients were promised updates "within seconds" |
| 2 | When is a file rebuilt, and what is sent when it changes? | Every rebuild resent all 7 files to every client, changed or not |
| 3 | How does one slow or broken client server not delay everyone else? | One hanging SFTP host stalled every client's send round |
| 4 | How do we know each client really received each file? | "Did NDTV get the final medal table?" must be answered from a screen |
| 5 | Where are stats worked out, and what happens when one fails? | A failed stat job told no one |
| 6 | How are new clients and file shapes added? | A new client must be settings, not code |
| 7 | What 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 (
/getfeedsand the/v1documents).
How it works

This is the agreed design. The steps, in order:
- 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.)
- 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.
- 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. - 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.
- 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.
- 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.
- 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.
- 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. - 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:
| Destination | What 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 S3 | Sends both files. S3 answers with an ETag. The file reaches NDTV at 21:14:05.0, 2 seconds after the change |
| News18 over FTP | Sends both files, one at a time (FTP default). Records the remote file size as its receipt |
| DailyHunt over SFTP | The 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).
- 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). - 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. - Each result is saved to
feed_documentwith a new sequence number, changed or not. - Every 5 seconds (
delivery_poll_seconds,settings.py:238) the delivery worker asks the database who is owed a file. - It sends up to 8 at a time (
delivery_concurrency,settings.py:237), waits for all of them, and only then goes round again. - Pull clients call
/getfeeds, served through a Redis cache kept for 3 seconds, or straight from Postgres when there is noREDIS_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
| Gap | What happens today | Where |
|---|---|---|
| Every rebuild resends all 7 files | A new seq per rebuild; the hash is never compared | feed.py:295; dispatcher.py:142 |
| One hanging host stalls everyone | Shared round; no SFTP or S3 time limit; serial worker loop | transports.py:436, 169-177; runner.py:224 |
| Nobody is told | The Slack alerter only logs | deadletter.py:137-162 |
| Builds go through the public API | Once about 30,000 API calls for one rebuild (now fixed); still under the API's 15 s time-out and 600-a-minute rate limit | feed.py:114-194; api settings.py:43-45 |
| Debounce with no maximum wait | A steady stream keeps pushing the rebuild back; feeds were 3 hours behind on 18 Sep | feed.py:494-501; settings.py:143-154 |
| A read copy could serve out-of-date files | GraphQL uses the replica when one is set, and builds go through GraphQL | api app.py:51-60 |
| No event filter when sending | Every Games' files match every assignment of that feed | dispatcher.py:121 |
| A second send path is wrong | The listener matches by feed family: country feeds to the wrong clients | delivery/listener.py:92 |
| Stats can skip a late commit | The tail reads "id greater than my place" | packages/stats/.../triggers.py:69 |
| Failed stats tell no one | After 3 tries a stat.failed event is written; no alert | packages/stats/.../worker.py:87, 142-160 |
| "Last updated" can be now | Falls back to the current time | api graphql/medal_types.py:336 |
| Credentials and hosts | SFTP host keys not checked; new S3 client per send | transports.py:343, 169 |
| Every pull hit reads Postgres | 4 reads before the cache; /v1 has no cache | getfeeds.py:159-170; documents.py:52 |
| The Redis cache is a 3 s timer | Never refreshed on rebuild; off with no REDIS_URL | getfeeds.py:132-137; settings.py:77, 91 |
| Pull data is not protected | /getfeeds needs no key | auth.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):
| Table | One row per | Holds today |
|---|---|---|
delivery_client | client | code, name, parent, active |
feed_package | product | which feed keys it grants |
delivery_assignment | client 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_attempt | send | outcome, detail, bytes, duration, trigger; pruned after 14 days |
delivery_deadletter | failed version | last error, attempts, resolved_at |
delivery_credential | credential ref | the 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
- Woken by the Section 8 wake-up, read new outbox rows with the late-commit-safe reader, collect the records they touched.
- 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. - 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). - Hash the body without its time stamp. Same as the stored
content_hash: stop. Different: save the body, hash,changed_atand the nextseq, and replace the file's dependency rows, in one transaction. - 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

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).
| Setting | Default | What it does |
|---|---|---|
connect_timeout_s | 10 | Give up on connecting after this |
send_timeout_s | 60 | Give up on one send after this |
fast_retries | 5 | Retries at 5 s, 15 s, 45 s, 2 min and 5 min |
slow_retry_s | 300 | After the fast retries, try every 5 min until it lands or a person pauses it |
parallel_files | 1 for FTP and SFTP, 4 for S3 and HTTP | How many files go to this destination at once |
alert_after_failures | 3 | Failures in a row before an alert |
alert_behind_s | 60 | How 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):
| Setting | Default | Today |
|---|---|---|
feed_debounce_seconds / feed_max_wait_seconds | 1 / 3 | debounce exists (1.0); no maximum wait |
delivery_connect_timeout / delivery_send_timeout | 10 / 60 | delivery_transport_timeout 30, FTP only |
delivery_destinations_parallel | 32 | delivery_concurrency 8, one shared round |
delivery_alert_behind_seconds / delivery_alert_failures | 60 / 3 | none |
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

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):
- Build the file, hash it, and commit
feed_documentif the hash changed. Postgres is the record. - 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).
- 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.
- 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):
- 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.
- Ask Redis for the version only. Same as the copy in memory: serve that copy. Changed: fetch the body once and keep it.
- 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.)
- 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_atto 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_eventrows 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 onestat.failedevent (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:
| Screen | Who | What they see | What they can do |
|---|---|---|---|
| Clients overview | Operator | One row per client: green, behind (by how long), failing, or switched off; the oldest file age | Open a client |
| One client | Operator | Each destination and file: version, sent at, size, confirmation, age when sent; a chart of age over the last hour | Send again now; pause a destination; test the connection |
| Freshness | Admin | For every client file, change-to-delivery time over the day: median and worst | Find the slowest destination |
| Feed settings | Admin | Client, package, destinations, file names, transform, country, schedule, marks to include | Add 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):
- Time limits on SFTP and S3, and one S3 client per destination.
- One task per destination instead of one shared round.
content_hashonfeed_document; send only on a change;last_updatedfrom real changes.- Remove the listener send path; add the event filter to the dispatcher.
- Proof of delivery and age in
delivery_state; the client screens. - Alerts through the central service.
- 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.
feed_dependencyand building only touched files.- 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
| Test | Proves |
|---|---|
| An SFTP test host that accepts and then hangs | It times out at 60 s; other destinations deliver on time (D4) |
| Rebuild with no real change | No new version; nothing sent (D3) |
| 50 changes in 2 seconds | Each file is built at most once per 3 s; each destination gets only the newest (D2, E3) |
| Kill the sender mid-upload | The same version is sent again; the receipt is recorded once (E5) |
| A host down for 10 minutes, then back | It 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 Games | Only the files containing it are rebuilt (E2) |
| Replay a day of Games data | 95% 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 wrong | What happens (agreed design) |
|---|---|
| DailyHunt's SFTP host hangs | Its 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 bytes | Same hash, no new version, nothing sent to anyone |
| 50 matches finish within a minute | Rebuilds 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-build | Nothing 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-send | The 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 hour | Its 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 slow | Builds 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 lags | The delivery-worker never builds from a read copy |
| Redis is down | API 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 commit | The Redis destination retries; pulls get the previous version meanwhile; an alert at 3 failures |
| Two Redis writes arrive out of order | Write-if-newer refuses the older one |
| Redis restarts empty, or evicts keys | The first hit per file per server reads Postgres once and refills Redis |
| An API server restarts | Its memory is empty; the first hits read Redis |
| A client key leaks | Revoked in the admin panel; servers drop it within 1 s by NOTIFY, 60 s at worst |
| A stat job fails 3 times | It stays failed and an alert names the stat and the match |
| A client needs a different country's feed | A setting on their destination, no code |
| A client is switched off on purpose | It 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.
| # | Decision | In plain words |
|---|---|---|
| D1 | One 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 own | One build path instead of two; never built from an out-of-date copy; no 30,000 API calls |
| D2 | Each 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 s | A change to one match does not rebuild unrelated files, and a stream of changes cannot delay a file |
| D3 | A new version only when the file's content hash changes; "last updated" is the newest real change, never now | Unchanged files are never sent again |
| D4 | Each client destination has its own queue; every transport has a 10 s connect and 60 s send limit; only the newest waiting version is sent | One hanging host delays only itself |
| D5 | Proof 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 |
| D6 | Freshness 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 |
| D7 | The stats queue reads the outbox without skipping late commits, and a failed stat job raises an alert | No stat is silently missing |
| D8 | Pull 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 request | Traffic never slows scoring, and clients keep reading during an outage |
| D9 | Clients 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 visible | A new client is set up in the admin panel the same day |
| D10 | Each file can carry its version number and, per unit, provisional or official; each client chooses | Clients can check they have the latest and the final result |
| D11 | The 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 nothing | Pull clients read from memory; an old file can never replace a newer one |
| D12 | Each destination has its own settings: connect and send limits, fast retries, files at once, alert levels; all destinations send in parallel | A slow or strict client is tuned without touching anyone else |
| D13 | A 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 change | 3 million hits in a match put no load on the scoring database |
| D14 | Pull needs a secret key per client, in a header (or in the address for old clients); the client name alone stops working | Protected data is no longer open to anyone who guesses a client name |
| D15 | Design and load-test pull for 5,000 hits a second at peak | A number we test against before a series goes live |
| E1 | feed_document gains content_hash and changed_at; seq moves only when the hash changes | No new version means nothing to send |
| E2 | A new feed_dependency table, replaced with each build; a change finds its files by one indexed lookup | Only the files a change touches are rebuilt |
| E3 | One queue and task per destination, coalescing to the newest version; up to 32 destinations at once per worker | A slow destination never blocks another; a returning host gets the newest file only |
| E4 | 10 s connect and 60 s send limits on all four transports; one reused S3 client per destination; SFTP host keys checked | No send can hang forever, and memory stays flat |
| E5 | delivery_state keeps version, hash, size, receipt and age for each destination and file | Proof of delivery and freshness come from one row |
| E6 | Per-destination back-off; alerts at 3 failures or 60 s behind, through the central alert service | A failing client is retried calmly and a person is told |
| E7 | Redis 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 newer | One small atomic write per changed file |
| E8 | Redis is written by the delivery-worker as one more destination (own queue, retries, proof), after the Postgres commit | If Redis is down, Postgres still has the file, and Redis is filled when it returns |
| E9 | Each API server keeps the last body per key in memory and asks Redis only for the version | Big files are not copied over the network on every hit |
| E10 | On 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 age | Redis trouble cannot flood Postgres and slow scoring |
| E11 | Answers are gzip with an ETag; a client with the current version gets 304 and no body | Files are 8 to 15 times smaller, and repeat polls cost almost nothing |
| E12 | delivery_assignment gains connect_timeout_s, send_timeout_s, fast_retries, slow_retry_s, parallel_files, alert_after_failures, alert_behind_s | The per-client send settings, edited in the admin panel |
Built today, or still to build
| Piece | Today (checked 7 Oct) | Agreed design | Status |
|---|---|---|---|
| Feed definitions as saved queries | feed_definition rows with a GraphQL query each | Kept | Built today |
| Who builds | Admin 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-up | Agreed, to build |
| How it builds | Over HTTP through the public API (feed.py:158-194) | Resolvers called directly, main database only | Agreed, to build |
| Which files are rebuilt | All 7 of each changed Games (feed.py:351, 386-437) | Only files in feed_dependency for the touched records | Agreed, to build |
| Debounce | 1 s, no maximum (feed.py:494-501) | 1 s, maximum 3 s | Partly built |
| Version number | New seq every rebuild (feed.py:295) | New seq only when content_hash changes | Agreed, to build |
| Older build cannot win | seq guard in the upsert (feed_store.py:73-83) | Kept | Built today |
| Who is owed what | Dispatcher query every 5 s (dispatcher.py:76-192); no event filter (:121) | Per-destination queues fed by the build step; event filter | Partly built |
| Second send path | Listener matches by feed family (listener.py:92) | Removed | Agreed, 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 |
| Sending | One shared round, 8 at a time, gather (runner.py:214-224) | One task per destination, up to 32 per worker | Agreed, to build |
| Coalescing | Per throttle key, inside one round (runner.py:185-193) | Per destination and file, newest only | Partly built |
| Time limits | FTP 30 s, webhook 10 s; SFTP and S3 none (transports.py:425-437) | 10 s connect, 60 s send, every transport, per client | Partly built |
| S3 client | New one per send (transports.py:169-171) | One per destination, reused | Agreed, to build |
| SFTP host keys | Not checked (transports.py:343) | Checked | Agreed, to build |
| Retries | 5 tries, 0.5 s doubling to 30 s, then dead letter (settings.py:184-186) | 5 fast retries, then every 5 min until it lands | Agreed, to build |
| Safety switches | delivery_enabled and destination fence (runner.py:226-259) | Kept | Built today |
| Proof of delivery | delivery_state with seq, time, size, status (delivery.py:197-253 in contract models) | Plus hash, receipt, age | Partly built |
| Alerts | Slack seam only logs (deadletter.py:137-162) | Central alert service, 60 s behind or 3 failures | Agreed, to build |
| Encrypted credentials | delivery_credential table | Kept, or Secrets Manager | Built today |
| Pull reads | 4 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-flight | Agreed, to build |
| Pull security | /getfeeds exempt from the API key (auth.py:42) | Secret key per client | Agreed, to build |
| Gzip and 304 | Gzip 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 kept | Partly built |
| Stats tail | "id greater than my place", 15 min sweep (triggers.py:65-71) | Late-commit-safe reader | Agreed, to build |
| Failed stats | stat.failed event, no alert (worker.py:142-160) | Alert naming the stat and match | Partly built |
| Client screens | DeliveryStatusPage.tsx, rows by assignment | Clients overview, one client, freshness | Partly built |
Numbers
| Number | What | Source |
|---|---|---|
| 0.30 / 0.43 ms | Redis write-if-newer, 3.3 KB file, median / 99th percentile | Measured: laptop, Redis 7 in Docker, Python client with 50 connections, 7 Oct |
| 0.67 / 2.05 ms | Same, 44 KB file | Same |
| 1.95 / 4.53 ms | Same, 233 KB file | Same |
| 0.27 to 0.30 ms | Redis version check, median, all three sizes | Same |
| 0.29 / 1.07 ms | Full Redis read, median, 3.3 KB / 233 KB | Same |
| 11,754 / 4,280 | Redis reads a second, one Python process, 3.3 KB / 233 KB | Same |
| 8 to 15 times | Gzip 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 KB | Measured on the final Asian Games files |
| 1.74 ms | sha256 of a 4.3 MB file | Measured, design doc |
| 1.3 s | Building the 4.3 MB calendar | A comment in the code, api app.py:157-165 |
| about 6 MB | All 7 Asian Games files | Final copy, 5 Oct |
| about 24 MB | Sent per rebuild today to 4 destinations | Worked out from the file sizes |
| 5,000 a second | Pull design target | Estimate: 3 million hits over a 3.5-hour T20 is about 240 a second; peaks assumed 10 times that; then doubled |
| 25 MB a second | 5,000 hits a second on a 5 KB file | Estimate (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.
Read next
- 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.