Course Content
Model Context Protocol (MCP)
3 sections · 8 lessons
Event Streams and Real-Time Context for Agents
An on-call assistant watches the incident tracker. Its MCP server polls every thirty seconds: select * from incidents where updated_at > $last_seen. It works, in the sense that it eventually notices things.
Then two problems arrive on the same shift. The first is cost. A twelve-hour shift is 1,440 polls per agent; across forty engineers that is 57,600 queries per shift, and on a normal day roughly 99% of them return zero rows. You are paying for a database round trip, a connection checkout and an index scan, 57,000 times, to be told nothing happened.
The second is worse, because it is silent. A SEV-1 is opened at 09:14:58 and the agent sees it at 09:15:20 — a 22-second lag on the incident everyone is looking at. And occasionally the agent never sees an incident at all. That bug is worth understanding, because it catches everyone who writes a polling watermark: updated_at is stamped when the row is written, but the row becomes visible when the transaction commits. A transaction that stamps 09:14:58 and commits at 09:15:03 is invisible to a poll that ran at 09:15:00 and advanced its watermark to 09:15:00. The row is now permanently in the past. Nothing will ever pick it up.
Polling is not just inefficient. Done naively over a wall clock, it is incorrect. The fix is to stop asking and start being told.
Push instead of poll
| Polling | Event stream | |
|---|---|---|
| Requests per hour, idle | 120 per client at 30 s intervals | None: one open connection carrying a 15 s heartbeat |
| Detection lag | Half the interval on average; the full interval at worst | Tens of milliseconds |
| Correctness | Watermarks race with commit order | The writer decides when to emit |
| Cost when nothing happens | Full price | Nearly zero |
| Cost when everything happens | Bounded by the interval | Unbounded — needs its own limits |
| Failure mode | Late, or silently missed | Dropped connections, backpressure |
Note the last two rows, because they are the honest trade. Polling has a natural rate limit: however chaotic the system gets, you never receive more than one batch per interval. Push does not. A deployment that flips 4,000 rows in one transaction will try to deliver 4,000 events, and if you have not thought about that, the agent's context window is the thing that breaks.
writer event bus MCP server agent host (any code) (Redis/NOTIFY) (SSE endpoint) (MCP client) | | | | | publish | | | |-------------------> | | | | fan out | | | |---------------------> | | | | notifications/ | | | | resources/updated | | | |--------------------> | | | | bounded queue heartbeat every 15 s per subscriber resumable by event idPush moves the decision about when something is interesting from the reader to the writer. That is the whole benefit, and also the whole risk.
Server-Sent Events
SSE is a one-directional stream of text over an ordinary HTTP response that never ends. It is the format underneath MCP's Streamable HTTP responses, and it is worth knowing at the byte level because almost every SSE bug is a formatting or buffering bug.
The examples in this section build your own event feed — the stream between the event bus and the MCP server in the diagram above. Keep one difference from MCP in mind: since the 2026-07-28 revision, MCP's own streams are not resumable (no event IDs, no Last-Event-ID), and a client that loses a stream simply re-sends the request. In your own feed, resumption is yours to provide, and the code below does.
HTTP/1.1 200 OKContent-Type: text/event-streamCache-Control: no-cacheConnection: keep-aliveX-Accel-Buffering: noid: 1042event: incident.updateddata: {"id":"INC-2231","severity":"SEV1","status":"mitigating"}: heartbeatid: 1043event: incident.resolveddata: {"id":"INC-2230","duration_s":1840}Four rules govern that format. A field is name: value. A blank line dispatches the event — forget it and the browser or client buffers your event forever. A line starting with : is a comment, which is how heartbeats are sent. And multi-line payloads need one data: line each, which is why JSON payloads must not contain raw newlines.
A correct multi-client server
The naive implementation keeps one global list of messages and hands every connection the same iterator. It breaks the moment two clients connect: they consume each other's events, and a slow client blocks a fast one. Each subscriber needs its own bounded queue.
1import asyncio, json, uuid2from collections import deque3from fastapi import FastAPI, Request4from fastapi.responses import StreamingResponse56app = FastAPI()78class EventHub:9 def __init__(self, replay_size: int = 500):10 self._subs: dict[str, asyncio.Queue] = {}11 self._recent: deque = deque(maxlen=replay_size) # for Last-Event-ID resume12 self._seq = 01314 def subscribe(self, maxsize: int = 100) -> tuple[str, asyncio.Queue]:15 sub_id = uuid.uuid4().hex16 self._subs[sub_id] = asyncio.Queue(maxsize=maxsize)17 return sub_id, self._subs[sub_id]1819 def unsubscribe(self, sub_id: str) -> None:20 self._subs.pop(sub_id, None)2122 def publish(self, event: str, data: dict) -> None:23 self._seq += 124 item = {"id": self._seq, "event": event, "data": data}25 self._recent.append(item)26 for sub_id, q in list(self._subs.items()):27 try:28 q.put_nowait(item)29 except asyncio.QueueFull:30 # Slow consumer policy: drop the oldest, keep the newest.31 # For alerting, recent state beats complete history.32 try:33 q.get_nowait()34 q.put_nowait(item)35 except asyncio.QueueEmpty:36 pass3738 def replay_after(self, last_id: int) -> list[dict]:39 return [e for e in self._recent if e["id"] > last_id]4041hub = EventHub()4243def sse(item: dict) -> str:44 payload = json.dumps(item["data"], separators=(",", ":"))45 return f"id: {item['id']}\nevent: {item['event']}\ndata: {payload}\n\n"4647@app.get("/events")48async def events(request: Request):49 last_id = int(request.headers.get("Last-Event-ID", 0) or 0)50 sub_id, queue = hub.subscribe()5152 async def stream():53 try:54 for missed in hub.replay_after(last_id): # resume where we dropped55 yield sse(missed)56 while True:57 if await request.is_disconnected():58 break59 try:60 item = await asyncio.wait_for(queue.get(), timeout=15.0)61 yield sse(item)62 except asyncio.TimeoutError:63 yield ": heartbeat\n\n" # keeps proxies from closing64 finally:65 hub.unsubscribe(sub_id) # always, on any exit path6667 return StreamingResponse(stream(), media_type="text/event-stream",68 headers={"Cache-Control": "no-cache", "X-Accel-Buffering": "no",69 "Connection": "keep-alive"})Five design decisions in that code are the difference between a demo and a service:
- One queue per subscriber, so clients cannot steal each other's events.
- Bounded at 100, so a client that stops reading cannot grow the server's memory without limit. An unbounded queue turns one stuck laptop into an out-of-memory kill.
- An explicit slow-consumer policy. Dropping the oldest suits alerting, where the current state matters more than history. For anything auditable, drop the connection instead and force a resync — silently discarding financial events is far worse than a reconnect.
- A 15-second heartbeat, because idle-connection timeouts in load balancers are commonly 60 seconds and a stream with nothing to say looks identical to a dead one.
- Cleanup in
finally. Without it, every disconnect leaks a queue thatpublishkeeps filling. Fifty disconnects an hour becomes a memory leak nobody can explain a week later.
Reading the stream
1import httpx, json, asyncio, random23async def consume(url: str, on_event) -> None:4 last_id, backoff = 0, 1.05 while True:6 try:7 headers = {"Accept": "text/event-stream"}8 if last_id:9 headers["Last-Event-ID"] = str(last_id)10 async with httpx.AsyncClient(timeout=None) as c:11 async with c.stream("GET", url, headers=headers) as r:12 r.raise_for_status()13 backoff = 1.0 # reset only after connecting14 fields = {}15 async for line in r.aiter_lines():16 if line == "": # blank line dispatches17 if "data" in fields:18 last_id = int(fields.get("id", last_id))19 await on_event(fields.get("event", "message"),20 json.loads(fields["data"]))21 fields = {}22 elif line.startswith(":"):23 continue # comment / heartbeat24 else:25 k, _, v = line.partition(":")26 fields[k] = v.lstrip()27 except Exception:28 await asyncio.sleep(backoff + random.uniform(0, 1))29 backoff = min(backoff * 2, 30.0) # capped exponential backoffTwo details matter. timeout=None on the read, because the default read timeout will kill a legitimately idle stream. And resetting backoff only after a successful connect — reset it on the first byte received and a server that accepts then immediately closes will be hammered at one connection per second forever.
Every buffer needs a stated policy for what happens when it fills. Drop oldest, drop newest, or disconnect and resync — refusing to choose means choosing "run out of memory".
The pub/sub bus
The hub above works inside one process. Two facts break that as soon as you run more than one replica: publishers in process A cannot reach subscribers in process B, and a client that reconnects may land on a different replica with a different sequence counter.
| Backing store | Delivery | Survives restart | Use when |
|---|---|---|---|
| In-process queues | At-most-once | No | Single replica, local dev, stdio servers |
| Redis pub/sub | At-most-once, fan-out | No | Several replicas, loss is tolerable |
| Redis Streams | At-least-once, with consumer groups | Yes | You need replay and acknowledgements |
| Postgres LISTEN/NOTIFY | At-most-once, transactional timing | No | The events already originate in the database |
| Kafka / NATS JetStream | At-least-once, retained log | Yes | High volume, multiple independent consumers |
"At-most-once" means an event published while a subscriber is disconnected is gone. That is acceptable for a live dashboard and unacceptable for anything that triggers an action, which is the distinction that decides your choice.
Turning database writes into events
PostgreSQL can push. A trigger calls pg_notify, and any session that has issued LISTEN receives the payload — crucially, only when the transaction commits, which is exactly the correctness property the polling watermark lacked.
1create or replace function notify_incident_change() returns trigger as $fn$2begin3 perform pg_notify('incident_events', json_build_object(4 'op', tg_op,5 'id', new.id,6 'severity', new.severity,7 'status', new.status8 )::text);9 return new;10end;11$fn$ language plpgsql;1213create trigger incident_change14 after insert or update on incidents15 for each row execute function notify_incident_change();1async def listen_forever(dsn: str) -> None:2 while True:3 conn = None4 try:5 conn = await asyncpg.connect(dsn) # a dedicated connection, never pooled6 await conn.add_listener("incident_events",7 lambda c, pid, channel, payload:8 hub.publish("incident.changed", json.loads(payload)))9 while True: # keep the session alive and healthy10 await asyncio.sleep(20)11 await conn.execute("select 1")12 except Exception as e:13 print(f"listener died: {e}", file=sys.stderr)14 await asyncio.sleep(2)15 finally:16 if conn:17 await conn.close()Three constraints on NOTIFY that bite in production:
- The payload is capped at 8,000 bytes. Exceeding it raises an error inside the trigger, which aborts the user's transaction. Never put row contents in the payload — send the primary key and let the consumer fetch what it needs.
- It is not durable. If no session is listening at commit time, the notification evaporates. A listener that restarts loses everything sent during the gap.
- The listening connection must not come from the pool. A pooled connection gets handed to someone else, who does not have your
LISTENregistered. Dedicate one connection and supervise it.
Where losing events is unacceptable, pair the notification with a durable outbox written in the same transaction:
1create table event_outbox (2 id bigserial primary key,3 topic text not null,4 payload jsonb not null,5 created_at timestamptz not null default now(),6 processed_at timestamptz7);8create index on event_outbox (processed_at) where processed_at is null;1async def drain_outbox(conn) -> int:2 """Claim unprocessed rows atomically; safe to run in several replicas."""3 rows = await conn.fetch("""4 with claimed as (5 select id from event_outbox6 where processed_at is null7 order by id8 limit 1009 for update skip locked10 )11 update event_outbox e set processed_at = now()12 from claimed c where e.id = c.id13 returning e.id, e.topic, e.payload""")14 for r in rows:15 hub.publish(r["topic"], json.loads(r["payload"]))16 return len(rows)for update skip locked is what makes this safe across replicas: each worker claims a disjoint set of rows and never blocks on rows another worker holds. The NOTIFY then becomes a latency optimisation — it wakes the drainer immediately — while the outbox table guarantees nothing is lost if the notification is missed.
Getting events into the agent's head
Here is where streaming systems meet an awkward truth: a language model has no ears. It is not sitting there listening. It processes a context window when it is asked to produce a token, and events that arrive between turns must be somewhere it can see them.
MCP's mechanism is resource subscriptions. The client opens a subscriptions/listen request naming the URIs it cares about (for example "resourceSubscriptions": ["incidents://active"]); the server pushes a notification on that stream when one changes; the notification carries only the URI, and the client re-reads if it cares.
1{2 "jsonrpc": "2.0",3 "method": "notifications/resources/updated",4 "params": {5 "_meta": { "io.modelcontextprotocol/subscriptionId": 3 },6 "uri": "incidents://active"7 }8}In the Python SDK, the server fans these out through a subscription bus, and your event listener publishes to it:
1from mcp.server.mcpserver import MCPServer2from mcp.server.subscriptions import InMemorySubscriptionBus, ResourceUpdated34bus = InMemorySubscriptionBus()5mcp = MCPServer("incidents", subscriptions=bus)67async def on_incident_changed(event: str, data: dict) -> None:8 # Called by the hub or outbox drainer; no MCP request is in progress here.9 await bus.publish(ResourceUpdated(uri="incidents://active"))The SDK delivers the notification only to listen streams that asked for that URI. With several replicas, implement the same bus interface over Redis pub/sub — the fan-out problem from the section above. Servers on protocol versions up to 2025-11-25 received a resources/subscribe call instead, but sent the same notification.
The host then decides what to do, and this is a product decision more than a technical one:
| Policy | Behaviour | Good for | Cost |
|---|---|---|---|
| Accumulate | Buffer events; inject a digest with the user's next message | Most assistants | Up to one turn of staleness |
| Interrupt | Wake the agent immediately on qualifying events | SEV-1 alerting, trading | Extra model calls; interrupts the user |
| Refresh on read | Ignore pushes; re-read subscribed resources before each turn | Slowly changing context | One extra read per turn |
Whichever you choose, the buffer must be bounded and summarised. A deployment that touches 4,000 rows produces 4,000 notifications; injecting them raw at roughly 40 tokens each is 160,000 tokens and a destroyed conversation. Collapse them instead:
1class EventBuffer:2 def __init__(self, capacity: int = 200):3 self.events = deque(maxlen=capacity) # oldest silently dropped4 self.dropped = 056 def add(self, event: dict) -> None:7 if len(self.events) == self.events.maxlen:8 self.dropped += 19 self.events.append(event)1011 def digest(self) -> str:12 """One compact block, safe to prepend to the next user turn."""13 if not self.events:14 return ""15 counts: dict[str, int] = {}16 for e in self.events:17 counts[e["event"]] = counts.get(e["event"], 0) + 118 head = ", ".join(f"{n} x {k}" for k, n in sorted(counts.items()))19 recent = "\n".join(f"- {e['event']}: {json.dumps(e['data'])[:160]}"20 for e in list(self.events)[-5:])21 extra = f" ({self.dropped} older events dropped)" if self.dropped else ""22 self.events.clear(); self.dropped = 023 return f"Since your last message: {head}{extra}.\nMost recent:\n{recent}"That digest turns 4,000 events into roughly 120 tokens: the counts by type, the five most recent in full, and an honest note about what was dropped. The model gets the shape of what happened plus the detail that is most likely to matter, and the conversation survives.
A language model has no ears. Events that arrive between turns exist only if something buffered them, bounded them and summarised them into the next prompt.
Deploying a streaming server
Long-lived connections break assumptions that every default in the stack was written around.
| Symptom | Cause | Fix |
|---|---|---|
| Events arrive in bursts every few KB | The reverse proxy is buffering the response | proxy_buffering off; and send X-Accel-Buffering: no |
| Connection drops at exactly 60 s | Proxy idle-read timeout | Heartbeat every 15 s; raise proxy_read_timeout |
| Works on one replica, flaky on three | Publisher and subscriber in different processes | Redis pub/sub or Streams behind the hub |
| Memory climbs steadily all day | Subscriber queues leaked on disconnect | Unsubscribe in finally; bound every queue |
| Deploys drop every stream | No graceful shutdown | On SIGTERM, send a final event and close; clients resume by ID |
| Compressed responses stall | gzip buffers until its window fills | Disable compression on the SSE route |
Capacity planning is different too. A polling service is sized by requests per second; a streaming service is sized by concurrent open connections. Each SSE connection holds a socket, a file descriptor and a task. On an async runtime, 10,000 concurrent streams is routine at a few kilobytes of memory each; on a thread-per-request server the same load needs 10,000 threads and falls over. This is why streaming MCP servers are written on async frameworks.
What this means when you build one
The instinct after learning about push is to stream everything. Resist it, because streams cost something polling does not: a stream is a stateful relationship you must keep alive, resume, bound and clean up, and each of those is a place to get it wrong.
Use the rate of change to decide. Data that changes less often than the user asks about it should be read on demand — a resource the host re-reads before each turn, with no subscription at all. Data that changes faster than the user asks, and where being late is a real cost, earns a stream. Everything in between is usually well served by a subscription with a coarse digest.
When you do stream, three commitments make it survivable. Bound every buffer and decide out loud what happens when it fills — drop oldest, drop newest, or disconnect and resync; there is no fourth option and refusing to choose means choosing "run out of memory". Give every event a monotonic ID and support resumption, because clients on laptops disconnect constantly and a stream you cannot resume is a stream that loses data on every train tunnel. And summarise before the model sees anything, because the context window is the narrowest pipe in the system and a stream is the easiest way to flood it.
The on-call team replaced their poller with a trigger, an outbox and one SSE endpoint. Idle database load from the assistant fell from 57,600 queries a shift to a few hundred outbox drains, detection lag went from a 15-second average to under 200 milliseconds, and the class of bug where an incident was committed just after a watermark advanced stopped existing — not because they fixed it, but because they deleted the code that could contain it.