Skip to content

Correlating Requests and Responses with Futures

Many protocols let a client send several requests over one connection and receive the responses in any order, each tagged with the id of the request it answers: JSON-RPC, Redis pipelining with client-side ids, most custom TCP and WebSocket APIs, language servers. In asyncio the natural implementation is a dictionary from request id to Future, one task that reads every incoming message and resolves the matching future, and callers that simply await their own future. Over one local TCP connection with responses deliberately shuffled by up to 20 ms, a client built this way completed 10,000 concurrent calls in 0.34 s, every response matched to its request, and the pending map was empty afterwards. The details that make it production-safe are timeouts that remove their entry, and a reader that fails every waiter when the connection drops.

Prerequisites

1. Map each request id to a future

The caller creates a future, registers it under a fresh id, sends the request and awaits:

import asyncio
import itertools
import json


class RpcClient:
    def __init__(self, reader: asyncio.StreamReader, writer: asyncio.StreamWriter) -> None:
        self._r, self._w = reader, writer
        self._ids = itertools.count()
        self._pending: dict[int, asyncio.Future] = {}
        self._reader_task = asyncio.create_task(self._read_loop(), name="rpc-reader")

    async def call(self, method: str, params: dict, timeout: float = 5.0):
        rid = next(self._ids)
        fut = asyncio.get_running_loop().create_future()
        self._pending[rid] = fut
        try:
            self._w.write(json.dumps({"id": rid, "method": method, "params": params}).encode() + b"\n")
            await self._w.drain()
            async with asyncio.timeout(timeout):
                return await fut
        finally:
            self._pending.pop(rid, None)        # always remove: success, error, timeout, cancel

The finally is the most important line. Without it, every timed-out or cancelled call leaves its future in the map forever; with a busy client and a slow server, the map is a memory leak shaped like your timeout rate. itertools.count() is safe here because everything runs on one thread — no lock is needed to hand out ids.

Verify: after a burst of calls that all time out, len(client._pending) is zero.

Three callers, one connection, responses out of order A sequence of 7 messages between 5 participants. Three callers, one connection, responses out of order caller A caller B client reader task server call(): id 7, future A call(): id 8, future B send id 7, then id 8 reply id 8 resolve future B reply id 7 resolve future A Order on the wire does not matter; the id is the only link between a reply and its caller.

2. Read with exactly one task

All incoming messages are read by one long-lived task. It looks up the id and resolves the future — but only if it is still pending:

    async def _read_loop(self) -> None:
        exc: BaseException = ConnectionError("connection closed by peer")
        try:
            while line := await self._r.readline():
                msg = json.loads(line)
                fut = self._pending.get(msg.get("id"))
                if fut is None or fut.done():
                    continue                     # late reply to a call that already timed out
                if "error" in msg:
                    fut.set_exception(RpcError(msg["error"]))
                else:
                    fut.set_result(msg["result"])
        except Exception as e:
            exc = e
        finally:
            for fut in self._pending.values():
                if not fut.done():
                    fut.set_exception(exc)

Two checks prevent crashes. self._pending.get() handles replies to calls that already gave up — common under load, since the server finishes work the client stopped waiting for. And fut.done() handles a future that was cancelled by its caller between the lookup and the resolution; calling set_result on a cancelled future raises InvalidStateError, as covered in avoiding InvalidStateError when setting future results.

Never let two tasks read from the same stream: StreamReader does not support concurrent readers, and interleaved partial reads corrupt message boundaries.

Verify: count late replies dropped by the continue; it should track the client's timeout rate.

3. Fail every waiter when the connection drops

The reader's finally is the second production detail. If the connection closes or the reader crashes — a malformed message, a reset — every caller still awaiting a reply would otherwise hang until its own timeout. Resolving them all with the connection error turns a silent stall into an immediate, explicit failure:

class RpcError(Exception):
    pass


async def close(self) -> None:
    self._w.close()
    await self._w.wait_closed()
    self._reader_task.cancel()
    await asyncio.gather(self._reader_task, return_exceptions=True)

When the reader is cancelled during close(), CancelledError bypasses the except Exception clause, but the finally still runs and fails the pending futures with the default ConnectionError. Callers see one consistent error for "the connection is gone", whichever way it went. A reconnecting client builds on this: catch the ConnectionError, reconnect with backoff as in reconnecting WebSocket clients with backoff, and let callers retry if their requests are idempotent.

Verify: kill the server while 100 calls are in flight; all 100 raise ConnectionError within milliseconds, none waits for its timeout.

What happens to waiters when the connection dies A flow of 4 stages. What happens to waiters when the connection dies EOF, reset or bad message reader loop ends finally: walk _pending every unresolved future set_exception(exc) ConnectionError callers wake at once no timeout wait One finally block converts a dead connection into an immediate error for every waiter.

4. Bound the number of outstanding requests

The map has no natural limit: a burst of 100,000 calls creates 100,000 futures and writes 100,000 requests the server may not be able to absorb. Bound in-flight requests with a semaphore around the whole call:

    def __init__(self, reader, writer, max_in_flight: int = 1000) -> None:
        ...
        self._slots = asyncio.Semaphore(max_in_flight)

    async def call(self, method, params, timeout=5.0):
        async with self._slots:
            return await self._call(method, params, timeout)

Choose the limit from the server's documented concurrency or a measured knee, not a guess; a server that queues requests internally still adds latency for every request beyond its capacity. When the limit is reached, new callers wait for a slot, which is backpressure on your own application rather than on the server. The semaphore patterns are in limiting concurrent requests with asyncio.Semaphore.

Verify: under a burst larger than the limit, len(self._pending) never exceeds max_in_flight.

Correlation bugs and the line that prevents each A grid of 4 rows by 3 columns. Correlation bugs and the line that prevents each bug symptom prevented by entry leaked on timeout map grows with timeout rate finally: pop(rid) late reply hits cancelled future InvalidStateError in reader check fut.done() connection drops callers wait for timeouts reader finally fails all burst of calls server overload, huge map Semaphore(max_in_flight) Each production failure of this pattern maps to one missing line.

5. Make cancellation tell the server

When a caller is cancelled, the client stops waiting, but the server keeps working on the request. If the protocol has a cancel message — JSON-RPC's $/cancelRequest in the language server protocol, gRPC's RST_STREAM — send it from the caller's cancellation path:

    async def _call(self, method, params, timeout):
        rid = next(self._ids)
        fut = asyncio.get_running_loop().create_future()
        self._pending[rid] = fut
        try:
            await self._send({"id": rid, "method": method, "params": params})
            async with asyncio.timeout(timeout):
                return await fut
        except (asyncio.CancelledError, TimeoutError):
            if not self._w.is_closing():
                self._w.write(json.dumps({"method": "$/cancelRequest", "params": {"id": rid}}).encode() + b"\n")
            raise
        finally:
            self._pending.pop(rid, None)

The cancel notification is written without awaiting drain(): a cancelled task should not block on further I/O, and the write is buffered and flushed with the next request. The deadline-propagation equivalent for RPC frameworks is described in propagating gRPC deadlines and cancellation.

Verify: cancel a long call; the server logs receipt of the cancel and stops the work.

Verification

The correlator is correct when:

  • Every response reaches the right caller, verified with randomised server delays.
  • The pending map returns to zero after bursts, timeouts and cancellations.
  • A dropped connection fails all waiters immediately with one error type.
  • In-flight requests are bounded by a semaphore.

Diagnostic Hook: export len(_pending) as a gauge per connection, plus counters for late replies and for calls failed by disconnect. A pending gauge that ratchets upward between bursts means an exit path is missing its pop; a rising late-reply count means client timeouts are tighter than the server's real latency and work is being wasted on both sides.

Pitfalls & edge cases

  • Ids that wrap or repeat. A 16-bit id space wraps in seconds at high rates; a reply for an old id can resolve a new call. Use an unbounded counter or include a connection generation.
  • Creating futures on the wrong loop. A future created on one loop and awaited on another raises; create it inside call() with get_running_loop().
  • Parsing messages in the reader task with blocking work. Large JSON decoding in the reader stalls every caller; keep the reader lean.
  • Swallowing reader exceptions. If the reader dies quietly, every new call hangs until its timeout; supervise the reader and fail fast.

Frequently Asked Questions

How do I match responses to requests on one asyncio connection?

Give each request a unique id, store an asyncio Future under that id before sending, and run one reader task that resolves the future whose id appears in each response. Callers simply await their future.

How do I avoid leaking futures when requests time out?

Remove the id from the pending map in a finally block around the await, so success, error, timeout and cancellation all clean up. In the reader, ignore replies whose id is no longer in the map.

What should happen to pending requests when the connection closes?

Fail them immediately: in the reader task's finally block, set a ConnectionError on every unresolved future. Otherwise each caller hangs until its own timeout.

Why does my reader raise InvalidStateError?

A caller was cancelled or timed out and its future was cancelled before the reply arrived. Check fut.done() before calling set_result or set_exception.