Model Context Protocol (MCP)

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.

Polling every thirty seconds against a pushPoll the table• Latency is half the interval, at best• Queries run whetheror not anything changed• Missed rows when thecursor moves wrongly• Cost scales withclients, not with eventsPush over SSE• Latency is the publish, plus the hop• Traffic only when something happens• One stream fans outthrough a pub/sub bus• Reconnect and replay become your problem
Push moves the hard part from noticing to not losing events across a reconnect.

Push instead of poll

PollingEvent stream
Requests per hour, idle120 per client at 30 s intervalsNone: one open connection carrying a 15 s heartbeat
Detection lagHalf the interval on average; the full interval at worstTens of milliseconds
CorrectnessWatermarks race with commit orderThe writer decides when to emit
Cost when nothing happensFull priceNearly zero
Cost when everything happensBounded by the intervalUnbounded — needs its own limits
Failure modeLate, or silently missedDropped 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.

Text
   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 id

Push 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.

Text
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.

Python
import asyncio, json, uuidfrom collections import dequefrom fastapi import FastAPI, Requestfrom fastapi.responses import StreamingResponseapp = FastAPI()class EventHub:    def __init__(self, replay_size: int = 500):        self._subs: dict[str, asyncio.Queue] = {}        self._recent: deque = deque(maxlen=replay_size)   # for Last-Event-ID resume        self._seq = 0    def subscribe(self, maxsize: int = 100) -> tuple[str, asyncio.Queue]:        sub_id = uuid.uuid4().hex        self._subs[sub_id] = asyncio.Queue(maxsize=maxsize)        return sub_id, self._subs[sub_id]    def unsubscribe(self, sub_id: str) -> None:        self._subs.pop(sub_id, None)    def publish(self, event: str, data: dict) -> None:        self._seq += 1        item = {"id": self._seq, "event": event, "data": data}        self._recent.append(item)        for sub_id, q in list(self._subs.items()):            try:                q.put_nowait(item)            except asyncio.QueueFull:                # Slow consumer policy: drop the oldest, keep the newest.                # For alerting, recent state beats complete history.                try:                    q.get_nowait()                    q.put_nowait(item)                except asyncio.QueueEmpty:                    pass    def replay_after(self, last_id: int) -> list[dict]:        return [e for e in self._recent if e["id"] > last_id]hub = EventHub()def sse(item: dict) -> str:    payload = json.dumps(item["data"], separators=(",", ":"))    return f"id: {item['id']}\nevent: {item['event']}\ndata: {payload}\n\n"@app.get("/events")async def events(request: Request):    last_id = int(request.headers.get("Last-Event-ID", 0) or 0)    sub_id, queue = hub.subscribe()    async def stream():        try:            for missed in hub.replay_after(last_id):    # resume where we dropped                yield sse(missed)            while True:                if await request.is_disconnected():                    break                try:                    item = await asyncio.wait_for(queue.get(), timeout=15.0)                    yield sse(item)                except asyncio.TimeoutError:                    yield ": heartbeat\n\n"             # keeps proxies from closing        finally:            hub.unsubscribe(sub_id)                     # always, on any exit path    return StreamingResponse(stream(), media_type="text/event-stream",        headers={"Cache-Control": "no-cache", "X-Accel-Buffering": "no",                 "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 that publish keeps filling. Fifty disconnects an hour becomes a memory leak nobody can explain a week later.

Reading the stream

Python
import httpx, json, asyncio, randomasync def consume(url: str, on_event) -> None:    last_id, backoff = 0, 1.0    while True:        try:            headers = {"Accept": "text/event-stream"}            if last_id:                headers["Last-Event-ID"] = str(last_id)            async with httpx.AsyncClient(timeout=None) as c:                async with c.stream("GET", url, headers=headers) as r:                    r.raise_for_status()                    backoff = 1.0                     # reset only after connecting                    fields = {}                    async for line in r.aiter_lines():                        if line == "":                # blank line dispatches                            if "data" in fields:                                last_id = int(fields.get("id", last_id))                                await on_event(fields.get("event", "message"),                                               json.loads(fields["data"]))                            fields = {}                        elif line.startswith(":"):                            continue                  # comment / heartbeat                        else:                            k, _, v = line.partition(":")                            fields[k] = v.lstrip()        except Exception:            await asyncio.sleep(backoff + random.uniform(0, 1))            backoff = min(backoff * 2, 30.0)          # capped exponential backoff

Two 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 storeDeliverySurvives restartUse when
In-process queuesAt-most-onceNoSingle replica, local dev, stdio servers
Redis pub/subAt-most-once, fan-outNoSeveral replicas, loss is tolerable
Redis StreamsAt-least-once, with consumer groupsYesYou need replay and acknowledgements
Postgres LISTEN/NOTIFYAt-most-once, transactional timingNoThe events already originate in the database
Kafka / NATS JetStreamAt-least-once, retained logYesHigh 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.

SQL
create or replace function notify_incident_change() returns trigger as $fn$begin  perform pg_notify('incident_events', json_build_object(    'op',       tg_op,    'id',       new.id,    'severity', new.severity,    'status',   new.status  )::text);  return new;end;$fn$ language plpgsql;create trigger incident_change  after insert or update on incidents  for each row execute function notify_incident_change();
Python
async def listen_forever(dsn: str) -> None:    while True:        conn = None        try:            conn = await asyncpg.connect(dsn)     # a dedicated connection, never pooled            await conn.add_listener("incident_events",                lambda c, pid, channel, payload:                    hub.publish("incident.changed", json.loads(payload)))            while True:                            # keep the session alive and healthy                await asyncio.sleep(20)                await conn.execute("select 1")        except Exception as e:            print(f"listener died: {e}", file=sys.stderr)            await asyncio.sleep(2)        finally:            if conn:                await conn.close()

Three constraints on NOTIFY that bite in production:

  1. 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.
  2. 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.
  3. The listening connection must not come from the pool. A pooled connection gets handed to someone else, who does not have your LISTEN registered. Dedicate one connection and supervise it.

Where losing events is unacceptable, pair the notification with a durable outbox written in the same transaction:

SQL
create table event_outbox (  id            bigserial primary key,  topic         text        not null,  payload       jsonb       not null,  created_at    timestamptz not null default now(),  processed_at  timestamptz);create index on event_outbox (processed_at) where processed_at is null;
Python
async def drain_outbox(conn) -> int:    """Claim unprocessed rows atomically; safe to run in several replicas."""    rows = await conn.fetch("""        with claimed as (          select id from event_outbox           where processed_at is null           order by id           limit 100             for update skip locked        )        update event_outbox e set processed_at = now()          from claimed c where e.id = c.id        returning e.id, e.topic, e.payload""")    for r in rows:        hub.publish(r["topic"], json.loads(r["payload"]))    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.

JSON
{  "jsonrpc": "2.0",  "method": "notifications/resources/updated",  "params": {    "_meta": { "io.modelcontextprotocol/subscriptionId": 3 },    "uri": "incidents://active"  }}

In the Python SDK, the server fans these out through a subscription bus, and your event listener publishes to it:

Python
from mcp.server.mcpserver import MCPServerfrom mcp.server.subscriptions import InMemorySubscriptionBus, ResourceUpdatedbus = InMemorySubscriptionBus()mcp = MCPServer("incidents", subscriptions=bus)async def on_incident_changed(event: str, data: dict) -> None:    # Called by the hub or outbox drainer; no MCP request is in progress here.    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:

PolicyBehaviourGood forCost
AccumulateBuffer events; inject a digest with the user's next messageMost assistantsUp to one turn of staleness
InterruptWake the agent immediately on qualifying eventsSEV-1 alerting, tradingExtra model calls; interrupts the user
Refresh on readIgnore pushes; re-read subscribed resources before each turnSlowly changing contextOne 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:

Python
class EventBuffer:    def __init__(self, capacity: int = 200):        self.events = deque(maxlen=capacity)   # oldest silently dropped        self.dropped = 0    def add(self, event: dict) -> None:        if len(self.events) == self.events.maxlen:            self.dropped += 1        self.events.append(event)    def digest(self) -> str:        """One compact block, safe to prepend to the next user turn."""        if not self.events:            return ""        counts: dict[str, int] = {}        for e in self.events:            counts[e["event"]] = counts.get(e["event"], 0) + 1        head = ", ".join(f"{n} x {k}" for k, n in sorted(counts.items()))        recent = "\n".join(f"- {e['event']}: {json.dumps(e['data'])[:160]}"                           for e in list(self.events)[-5:])        extra = f" ({self.dropped} older events dropped)" if self.dropped else ""        self.events.clear(); self.dropped = 0        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.

SymptomCauseFix
Events arrive in bursts every few KBThe reverse proxy is buffering the responseproxy_buffering off; and send X-Accel-Buffering: no
Connection drops at exactly 60 sProxy idle-read timeoutHeartbeat every 15 s; raise proxy_read_timeout
Works on one replica, flaky on threePublisher and subscriber in different processesRedis pub/sub or Streams behind the hub
Memory climbs steadily all daySubscriber queues leaked on disconnectUnsubscribe in finally; bound every queue
Deploys drop every streamNo graceful shutdownOn SIGTERM, send a final event and close; clients resume by ID
Compressed responses stallgzip buffers until its window fillsDisable 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.