Start with an experiment that has nothing to do with realtime UIs. A pipe has 2 KB in it. You are told it is readable, you read 1 KB, and you ask again. Does the kernel say it is readable?
It depends on how you registered it, and the answer carries over, almost line for line, to how a client should consume an event stream.
TL;DR
- epoll’s edge-triggered mode reports a transition and does not repeat it. Read part of the data and wait again, and you can wait forever. Level-triggered mode reports a condition (“still readable”) and repeats it until the condition is gone.
- An event stream has the same two designs. A diff (“item 7 is now X”) is an edge: miss one and the client is wrong with no way to know. A dirty flag (“item 7 may be stale, go and look”) is a level: duplicates, reordering and coalescing do not matter, because the client always asks for the current value.
- A hint alone is not yet a level. What makes epoll robust is that the kernel keeps the condition. For an application the equivalent is something the client can re-read and compare (a version, a cursor) plus a periodic check. In a simulation on a lossy channel, hints without that check converged in 179 of 300 runs; with it, in 300 of 300.
- There is a race that is easy to miss: an event that arrives while the refetch is in flight. The answer must be to refetch again. Delete the one line that does that and two of three tests fail.
- Level-triggered costs fetches. In the same simulation a well-built diff consumer needed about 1 fetch per run; the invalidating consumer needed about 64. Pick levels when correctness under messy delivery matters more than that cost, or when the repair path and the normal path being the same code is worth more than the saving.
Experiment 1: what epoll does after a partial read
The epoll(7) manual page describes this exact scenario: a pipe’s read side is registered, a writer writes 2 kB, epoll_wait reports it, the reader reads 1 kB, and epoll_wait is called again. For the edge-triggered flag EPOLLET, the manual says the second call “will probably hang despite the available data still present in the file input buffer”. Here it is, in about twenty lines (Linux only):
epoll_demo.py
# Edge-triggered vs level-triggered epoll: what happens after a partial read?
# Linux only. Run: python3 epoll_demo.py
import os, select
def run(flags, label):
r, w = os.pipe()
os.set_blocking(r, False)
ep = select.epoll()
ep.register(r, select.EPOLLIN | flags)
os.write(w, b"x" * 2048) # (2) the writer writes 2 kB
first = ep.poll(timeout=1) # (3) the fd is reported ready
os.read(r, 1024) # (4) the reader consumes only 1 kB
second = ep.poll(timeout=0.5) # (5) wait again
print(f"{label:15} first={bool(first)} second={bool(second)} (1024 bytes still unread)")
if not second:
# The documented fix for edge-triggered use: read until EAGAIN.
try:
while os.read(r, 1024):
pass
except BlockingIOError:
print(f"{'':15} drained until EAGAIN")
ep.close(); os.close(r); os.close(w)
run(0, "level-triggered")
run(select.EPOLLET, "edge-triggered")
Output on Linux 6.12:
level-triggered first=True second=True (1024 bytes still unread)
edge-triggered first=True second=False (1024 bytes still unread)
drained until EAGAIN
The edge-triggered registration reported the pipe once, when data arrived. Nothing changed afterwards, so nothing was reported, although 1024 bytes were sitting there. The manual’s remedy is to use non-blocking descriptors and to read until EAGAIN: an edge-triggered consumer must consume the whole transition. A level-triggered consumer does not have that obligation; it can take a little, come back, and be told again.
The distinction is older than epoll (it comes from signal levels versus signal transitions in hardware), but epoll makes it easy to feel in your fingers: edge triggering moves the burden of never missing anything onto the consumer. Level triggering removes the burden by making the notification a statement about the present.
The same choice in an event stream
Replace “pipe” with “list of items shown on a screen” and “readable” with “stale”:
| Edge: send the diff | Level: send “this is dirty” | |
|---|---|---|
| Event says | item:7 title = "B" | item:7 may have changed |
| Consumer does | applies the payload | refetches item:7 |
| Missed event | The client is wrong and cannot tell | The client is late until the next event or check |
| Duplicate event | Must be idempotent | Free: marking dirty twice is marking dirty |
| Reordered events | An older diff can overwrite a newer one | Free: the refetch returns the current value |
| Burst of 100 events | 100 applications | Folds into one dirty entry per key |
| Payload | Carries data | Carries a key |
This is how Kubernetes describes its controllers. The design archive says controllers are level-based “to maximize fault tolerance”, operating correctly “regardless of how many intermediate state updates may have been missed”, while still using watch notifications to cut latency. The notification is an optimization; the state is the truth. PostgreSQL’s NOTIFY documentation points the same way: if you need to communicate large amounts of information, put it in a table and send the key of the record.
Experiment 2: which consumer survives a bad channel?
A claim like “diffs are fragile” deserves a test. The simulation below has a server with five keys, about one update every other tick for 200 ticks (one tick is 10 ms of virtual time), and five consumers that each receive the same updates through their own bad channel: each message is dropped with probability p, duplicated with probability 0.2, and delayed by 0–7 ticks (so messages can overtake each other). A consumer’s fetch returns the value as of the moment of the request and arrives 1–8 ticks later, so fetch responses can also overtake each other. After the updates stop, the server sends a small “level” message every 200 ms, carrying the current version of each key; it goes through the same bad channel, so it can be dropped too. Each consumer is judged on whether, at the end, every key matches the server.
- A: apply each diff in arrival order.
- B: apply a diff only if its version is newer than what the client holds. (Idempotent and order-insensitive.)
- C: hint only: every event marks the key dirty, and a coalescing invalidator refetches.
- D: C, plus the periodic level message: any key whose server version is newer than the client’s is marked dirty.
- E: B, plus the same periodic message: any key that is behind is fetched.
sim.mjs: deterministic, 300 seeded trials per setting
// Which consumer design converges when the channel drops, duplicates and reorders messages?
// Deterministic: a seeded PRNG and a virtual clock (1 tick = 10 ms). Run: node sim.mjs
import { createInvalidator } from "./invalidator.mjs";
const KEYS = ["a", "b", "c", "d", "e"];
const rng = (a) => () => { // mulberry32: small, seedable, good enough for a simulation
a = (a + 0x6d2b79f5) | 0;
let t = Math.imul(a ^ (a >>> 15), 1 | a);
t = (t + Math.imul(t ^ (t >>> 7), 61 | t)) ^ t;
return ((t ^ (t >>> 14)) >>> 0) / 2 ** 32;
};
async function trial(seed, { pDrop, pDup, maxDelay }) {
const rand = rng(seed), pick = (n) => Math.floor(rand() * n);
let tick = 0, q = [];
const at = (d, fn) => { const h = { t: tick + d, fn, dead: false }; q.push(h); return h; };
const timers = { setTimeout: (fn, ms) => at(Math.max(1, Math.round(ms / 10)), fn), clearTimeout: (h) => (h.dead = true) };
const server = Object.fromEntries(KEYS.map((k) => [k, 0])); // key -> version (the value is the version)
const mk = () => ({ val: Object.fromEntries(KEYS.map((k) => [k, 0])), fetches: 0 });
const A = mk(), B = mk(), C = mk(), D = mk(), E = mk(); // five consumers, same updates, own channel each
const refetch = (c) => (k) => new Promise((done) => {
c.fetches++;
const answer = server[k]; // the server answers with the value NOW
at(1 + pick(maxDelay), () => { if (answer > c.val[k]) c.val[k] = answer; done(); }); // late answers may arrive out of order
});
const invC = createInvalidator({ refetch: refetch(C), timers });
const invD = createInvalidator({ refetch: refetch(D), timers });
const fetchE = refetch(E);
const handlers = {
A: (m) => m.type === "event" && (A.val[m.k] = m.v), // apply payload in arrival order
B: (m) => m.type === "event" && m.v > B.val[m.k] && (B.val[m.k] = m.v), // apply payload if newer
C: (m) => m.type === "event" && invC.invalidate(m.k), // hint only
D: (m) => (m.type === "event" ? invD.invalidate(m.k) : KEYS.forEach((k) => m.vers[k] > D.val[k] && invD.invalidate(k))),
E: (m) => (m.type === "event" ? m.v > E.val[m.k] && (E.val[m.k] = m.v) : KEYS.forEach((k) => m.vers[k] > E.val[k] && fetchE(k))),
};
const send = (m) => { for (const h of Object.values(handlers)) { // an independent lossy channel per consumer
if (rand() < pDrop) continue;
for (let i = 0; i < (rand() < pDup ? 2 : 1); i++) at(pick(maxDelay), () => h(m));
} };
const END = 200, STOP = 600;
for (tick = 0; tick <= STOP; tick++) {
if (tick < END && rand() < 0.5) { const k = KEYS[pick(KEYS.length)]; send({ type: "event", k, v: ++server[k] }); }
if (tick % 20 === 0 && tick > END) send({ type: "heartbeat", vers: { ...server } }); // the "level": cheap, repeated, idempotent
for (let moved = true; moved;) {
moved = false;
for (const h of q.filter((x) => x.t <= tick && !x.dead)) { h.dead = true; h.fn(); moved = true; }
q = q.filter((x) => !x.dead);
await new Promise(setImmediate); // let promise continuations run
}
}
const ok = (c) => KEYS.every((k) => c.val[k] === server[k]);
return { A, B, C, D, E, ok };
}
for (const cfg of [{ pDrop: 0, pDup: 0.2, maxDelay: 8 }, { pDrop: 0.2, pDup: 0.2, maxDelay: 8 }]) {
const N = 300, tally = Object.fromEntries("ABCDE".split("").map((n) => [n, { ok: 0, fetches: 0 }]));
for (let s = 1; s <= N; s++) {
const r = await trial(s, cfg);
for (const n of "ABCDE") { tally[n].ok += r.ok(r[n]) ? 1 : 0; tally[n].fetches += r[n].fetches; }
}
console.log(`\ndrop=${cfg.pDrop} dup=${cfg.pDup} reorder<=${cfg.maxDelay} ticks, ${N} trials`);
const names = { A: "A delta, apply in arrival order", B: "B delta, apply if newer", C: "C hint only (invalidate+refetch)", D: "D hint + periodic level check", E: "E delta + periodic level check" };
for (const n of "ABCDE") console.log(`${names[n].padEnd(36)} converged ${String(tally[n].ok).padStart(3)}/${N} fetches/trial ${(tally[n].fetches / N).toFixed(1)}`);
}
Output (seeds 1–300, Node.js 20.19; the run is deterministic, so you will get the same numbers):
drop=0 dup=0.2 reorder<=8 ticks, 300 trials
A delta, apply in arrival order converged 203/300 fetches/trial 0.0
B delta, apply if newer converged 300/300 fetches/trial 0.0
C hint only (invalidate+refetch) converged 300/300 fetches/trial 75.2
D hint + periodic level check converged 300/300 fetches/trial 75.2
E delta + periodic level check converged 300/300 fetches/trial 0.0
drop=0.2 dup=0.2 reorder<=8 ticks, 300 trials
A delta, apply in arrival order converged 68/300 fetches/trial 0.0
B delta, apply if newer converged 101/300 fetches/trial 0.0
C hint only (invalidate+refetch) converged 179/300 fetches/trial 64.2
D hint + periodic level check converged 300/300 fetches/trial 64.6
E delta + periodic level check converged 300/300 fetches/trial 1.0
Four things stand out, and the last one is the least flattering to the idea of sending hints.
- Duplicates and reordering alone break the naive diff consumer (A: 203 of 300), and a version guard fixes it completely (B).
- Loss breaks everything that depends on the last event for a key. With 20% loss, B converged in 101 runs and C in 179. A lost hint is as final as a lost diff, if that was the last message about the key.
- What fixes loss is the level, not the hint. D and E both converged every time, and the difference between C and D is just the periodic comparison. The kernel’s readiness flag in epoll plays this role; an application has to build it: a version, a cursor, or a digest that the client can compare against the server at any time.
- A well-built diff consumer with a repair path (E) is cheaper. It fetched about once per run, against about 64 for the invalidating consumer.
So the case for level-triggered design is not “diffs cannot work”. It is that in D the repair path and the normal path are the same code. E has two paths: an apply path that runs constantly, and a fetch path that runs only after something went wrong, which is exactly the code that is least exercised and most likely to be broken the first time it matters. If fetches are cheap, take D. If they are expensive, take E, and test the repair path as hard as the apply path.
Building the invalidator
The consumer in the simulation is the code below. It does four things, each one answering a specific mistake.
invalidator.mjs
// A level-triggered consumer: events only say "this key may be stale".
// The consumer owns the question "what is the current value?" and always answers it by refetching.
export function createInvalidator({
refetch, // async (key) => void; must read the CURRENT value and store it
windowMs = 50, // coalescing window: fixed, so latency is bounded
maxBackoffMs = 5000,
timers = { setTimeout, clearTimeout },
}) {
const dirty = new Set(); // keys that may be stale and are not being fetched right now
const fetching = new Set(); // keys with a refetch in flight
const failures = new Map(); // key -> consecutive failures
let timer = null;
function schedule(ms = windowMs) {
if (timer === null) timer = timers.setTimeout(flush, ms);
}
function flush() {
timer = null;
for (const key of [...dirty]) {
if (fetching.has(key)) continue; // stays dirty; re-armed when that fetch settles
dirty.delete(key);
fetching.add(key);
Promise.resolve(refetch(key)).then(
() => { failures.delete(key); settle(key, windowMs); },
() => {
const n = (failures.get(key) ?? 0) + 1;
failures.set(key, n);
dirty.add(key); // still stale
// full jitter backoff so that many clients do not retry in lockstep
settle(key, Math.random() * Math.min(maxBackoffMs, windowMs * 2 ** n));
},
);
}
}
function settle(key, nextMs) {
fetching.delete(key);
if (dirty.has(key)) schedule(nextMs); // an event arrived during the fetch, or the fetch failed
}
return {
invalidate(key) { dirty.add(key); schedule(); }, // any order, any number of times
invalidateAll(keys) { for (const k of keys) dirty.add(k); schedule(); }, // after (re)connect
pending: () => dirty.size + fetching.size,
};
}
Fold events into a set of keys. invalidate(key) adds to dirty and arms one timer. A hundred events for two keys leave two entries.
Use a fixed window, not a debounce. The timer is armed once, by the first event, and is not reset by later ones. A true debounce, which restarts on every event, never fires while events keep arriving; a steady stream starves the UI. The fixed window bounds latency at windowMs.
Handle the event that arrives mid-fetch. This is the race. Suppose the refetch for k is running and has already read the server. An event for k arrives: the thing it announces is newer than what the fetch will return. If the consumer drops it (“already fetching, nothing to do”), the screen shows the older value until some unrelated event. So the consumer keeps k in dirty while a fetch is in flight and schedules another flush when the fetch settles. That is the if (dirty.has(key)) schedule(nextMs) line. Remove it and see the tests fail:
invalidator.test.mjs
import assert from "node:assert/strict";
import { test } from "node:test";
import { createInvalidator } from "./invalidator.mjs";
const sleep = (ms) => new Promise((r) => setTimeout(r, ms));
test("100 events on 2 keys become 2 fetches", async () => {
const fetched = [];
const inv = createInvalidator({ refetch: async (k) => void fetched.push(k), windowMs: 20 });
for (let i = 0; i < 100; i++) inv.invalidate(i % 2 ? "list" : "item:1");
await sleep(80);
assert.deepEqual(fetched.sort(), ["item:1", "list"]);
});
test("an event that arrives during a fetch causes a second fetch", async () => {
let version = 1, shown = 0, calls = 0;
const inv = createInvalidator({
windowMs: 10,
refetch: async () => {
calls++;
const v = version; // the server answers with the value at request time...
await sleep(40); // ...and the response is slow
shown = v;
},
});
inv.invalidate("k"); // version 1 is announced
await sleep(20); // fetch #1 is now in flight, holding version 1
version = 2; inv.invalidate("k"); // version 2 is announced while fetch #1 is running
await sleep(200);
assert.equal(calls, 2);
assert.equal(shown, 2); // without the re-arm, the UI would stay at 1
});
test("a failed fetch is retried with backoff instead of being forgotten", async () => {
let calls = 0;
const inv = createInvalidator({
windowMs: 5,
maxBackoffMs: 20,
refetch: async () => { if (++calls < 3) throw new Error("boom"); },
});
inv.invalidate("k");
await sleep(300);
assert.equal(calls, 3); // two failures, then success, then it stops
assert.equal(inv.pending(), 0);
});
ok 1 - 100 events on 2 keys become 2 fetches
ok 2 - an event that arrives during a fetch causes a second fetch
ok 3 - a failed fetch is retried with backoff instead of being forgotten
With that line commented out, the second and third tests fail (I checked); only the first still passes.
Retry failures with jittered backoff, and one fetch per key at a time. A failed fetch puts the key back in dirty. The delay is randomized (full jitter, as in this analysis of backoff and jitter) so that many clients that failed together do not retry together. And because a key is never fetched twice concurrently, responses for the same key cannot overtake each other.
Two things the 40 lines do not do, deliberately:
- Reconnect. After a reconnect you do not know what you missed, so call
invalidateAll(keys)for the keys you display. - The periodic level check. In the simulation it was a message from the server. In a real system it might be a version on each response header, an ETag, a cursor, or a cheap digest endpoint. The shape matters more than the mechanism: something the client can compare, that persists even if the event was lost.
Trade-offs
- Fetch amplification. A dirty key costs a fetch even when the client could have computed the new value from the event. In the simulation the invalidating consumers fetched about 64 to 75 times in a run in which the diff consumer with a repair path fetched about once.
- Thundering herd. If a server broadcasts one hint to many clients, they all refetch at once. Add jitter to the first window (not only to retries) or let the server spread the hints.
- Hints that carry a version. A cheap middle path: send
{key, version}. The client skips the fetch if it already holds that version or newer, which removes the duplicates you caused yourself. - Hints over best-effort channels. Channels that may drop or collapse messages, such as push notifications, are fine for hints. Do not make them the only carrier of the level.
How it goes wrong
| Failure | Symptom | Defence |
|---|---|---|
| Event arrives during a refetch and is ignored | Screen is stale until an unrelated event | Keep the key dirty and refetch again (the re-arm line). |
| Debounce that restarts on every event | UI never updates during a steady stream | Fixed window from the first event. |
| Last event for a key is lost | Stale until the next change | A periodic comparison of versions, or a refresh on focus/reconnect. |
| Out-of-order responses for one key | An older value replaces a newer one | One fetch per key at a time, or a version guard on write. |
| Failed fetch forgotten | Permanently stale | Return the key to dirty and retry with jittered backoff. |
| Reconnect without resync | Everything may be stale | invalidateAll for visible keys. |
| Many clients refetch on one broadcast | Load spike on the origin | Jitter the first window; hint with a version. |
When to use it, and when not
Use it when the thing the client displays has a cheap “read the current value” operation, when delivery is messy (reconnects, flaky networks, fan-out through intermediaries), and when you want the repair path to be the same as the normal path.
Do not use it for:
- Event streams where every event is the data: audit logs, ledgers, anything that must see each step in order. Those need a log with cursors, not invalidation (see resumable streams).
- Large objects with tiny changes, where refetching the whole value per event is wasteful. Consider diffs with versions, with a hint-style resync as the fallback (consumer E).
- Notifications that are the content, such as a toast that says something happened. A hint cannot be re-read to the same effect.
Run the experiments
# Experiment 1 (Linux only)
python3 epoll_demo.py
# Experiments 2 and the tests (Node.js 18+; I ran Node.js 20)
node --test invalidator.test.mjs
node sim.mjs
Then break things on purpose. Delete the if (dirty.has(key)) schedule(nextMs) line and rerun the tests. Change pDrop in sim.mjs to 0.5 and see which consumers survive. The simulation only sends the level message after the updates stop (tick > END); remove that condition so that it also runs during the updates, and see what changes in the fetch counts.
Facts or hints
An event stream can carry facts about the world, or it can carry hints that the world is different now. The first is cheaper when everything is delivered once, in order. The second keeps working when delivery is messy, because a hint cannot be wrong, only late. Either way, give the client a way to compare its state with the server’s, because the stream is never the whole truth.
What was verified, and what is only a model
Experiment 1 ran on Linux 6.12 with Python 3.13. Experiments 2 and the tests ran on Node.js 20.19. The simulation is a model: its channel (independent drop, duplication and delay), its tick length and its parameters are choices of mine, not measurements of any network, and the fetch counts depend on them. I did not test against real browsers or real networks, and I did not benchmark the invalidator. The epoll behavior is the documented one, reproduced on this one kernel.