skip to content

What changes when an autogen-core app moves to the distributed gRPC runtime?

level: seniorimportance: nice to knowfreq 28%

answer

  1. Same API, different failure model
  2. Objects must become bytes
  3. No more shared module globals
  4. A host process relays between workers
  5. Experimental, not the default path

basics

~20 s

Agents stop sharing a process: messages must be serializable and serializers registered, a host process relays traffic and subscriptions between workers, and shared Python objects, in-process exceptions and instant delivery all stop being available. The gRPC runtime is experimental.

solid answer

~50 s

`SingleThreadedAgentRuntime` runs every agent in one process on one asyncio event loop, so messages are passed as live Python objects and a handler exception is an ordinary traceback. The distributed alternative in `autogen_ext.runtimes.grpc` replaces that with a `GrpcWorkerAgentRuntimeHost` process that relays messages and subscription registrations, and one `GrpcWorkerAgentRuntime` per worker process. Three things change materially. **Serialization**: every message type crossing the wire needs a registered serializer, so ad-hoc types and anything holding non-serializable handles stop working. **Shared state**: agents in different workers no longer share module globals, caches or clients — anything shared must become a message or an external store. **Failure and latency**: delivery is a network hop, a worker can be down or restart mid-flight, and errors arrive as transport failures rather than local exceptions. The API shape of `send_message` and `publish_message` is unchanged, which is the point, but the operational model is not. Treat the gRPC runtime as experimental.

go deeper

for a junior

Know that the default SingleThreadedAgentRuntime runs everything in one process, and that AutoGen also offers an experimental gRPC-based runtime for spreading agents across processes.

for a middle

Explain the host-plus-worker topology and the serialization requirement, and say why the send/publish API staying identical does not make the port trivial.

for a senior

Talk through what actually breaks: shared in-memory state, exception propagation, delivery guarantees under worker restarts, and chatty message designs that were free in-process.

for a principal

Own the decision itself — what isolation or independent-scaling need justifies a distributed topology, what it costs in schema versioning and operations, and why model latency usually dominates anyway.

## What the local runtime quietly gives you `SingleThreadedAgentRuntime` is a single process running a single asyncio loop. Every agent instance lives in the same interpreter, so a message is just a Python object handed from one coroutine to another — no copying, no encoding, no schema. Handler exceptions propagate to the caller of `send_message` as ordinary exceptions with a full traceback. Latency is microseconds. Agents can, if you let them, touch the same module-level cache or the same HTTP client. `runtime.start()` begins draining the queue and `await runtime.stop_when_idle()` blocks until it is empty. Every one of those conveniences is a hidden dependency on co-location, and moving to distribution removes them. ## The distributed topology The distributed runtime lives in `autogen-ext` under `autogen_ext.runtimes.grpc`, behind the `grpc` extra. It has two pieces: - **`GrpcWorkerAgentRuntimeHost`** — a standalone relay process bound to an address. It is the rendezvous point: workers connect to it, register their agent types and subscriptions with it, and it routes messages between them. - **`GrpcWorkerAgentRuntime`** — the runtime object each worker process constructs with the host's address. Agents register against it exactly as they would locally, and it exposes the same `send_message` / `publish_message` surface. Application code that only calls those two methods is largely portable between the runtimes. That deliberate API symmetry is the framework's selling point here, and also the trap: identical code, very different failure model. ## Serialization is the first wall you hit In-process, a message can be any Python object. Across processes it must be encoded. The distributed runtime requires message types to have serializers registered with the runtime, using the helper that derives serializers for known shapes — dataclasses and Pydantic models are the well-trodden paths. Consequences: - A message holding a database connection, an open socket, a callback or a model client cannot cross the wire. Those must be reconstructed on the receiving side, usually from configuration, not shipped. - Both sides need the *same* message definitions. A worker deployed with a stale copy of the message module receives something it cannot deserialize or, worse, deserializes into a type no handler is annotated for — and unhandled types are only logged. Versioning message schemas becomes a real deployment concern, not a code-organisation preference. ## Shared state disappears In one process, two agents can share a cache by accident and it works. Split across workers, each has its own interpreter: module globals diverge, in-memory caches are per-worker, and an agent instance's state lives only in the worker that created it. Anything two agents must agree on has to become either a message or an external store — Redis, a database, a blob store. This is usually the largest rewrite when moving an existing local design, and it is where an interviewer is looking for scars rather than API recall. ## Failure, ordering and latency Locally, a handler either returns or raises. Distributed, the same call can fail because the recipient's worker is starting up, has crashed, has been redeployed, or because the host is unreachable. A direct send now has a timeout dimension and deserves a cancellation token and a bounded wait. A publish, already fire-and-forget locally, becomes fire-and-forget across a network, so a subscriber that is down simply misses the event — there is no durable queue behind it. Do not assume cross-worker ordering, and do not build a protocol that depends on two publishes arriving in sequence. Latency also changes cost calculus: a chatty design that made hundreds of tiny sends per task is fine in-process and expensive across a network. Coarser messages carrying more work per hop is the usual correction. ## When it is worth it Distribution earns its complexity when agents need independent scaling or isolation — a code-executing agent on a hardened host, a GPU-bound agent on different hardware, teams owning agents deployed on their own release cadence, or a workload where one agent type must scale horizontally while others stay singletons. It is not the answer to "my agent loop is slow", which is almost always model latency and is unaffected by moving processes around. ## Say this in an interview Lead with the three concrete changes — serialization, no shared memory, network failure semantics — then note that the messaging API is intentionally the same so the port is mostly about state and message design, and finish by flagging that the gRPC runtime is experimental and that most production AutoGen work still runs single-process with the model API as the real bottleneck. That last piece of honesty reads as experience, not evasion.

  • Why can't you just ship any Python object as a message once you distribute?
    Because it has to be encoded, sent and reconstructed in another interpreter. The distributed runtime requires registered serializers, so dataclasses and Pydantic models work while anything holding live handles — sockets, DB connections, callbacks, model clients — does not. Those get rebuilt on the receiving side from config. Both processes also need matching message definitions, which turns message schemas into a versioned deployment artefact.
  • Would moving to the gRPC runtime speed up a slow agent loop?
    Almost never. In a typical AutoGen application the dominant cost is model latency and token generation, not message dispatch, and distribution adds a network hop rather than removing work. Distribution is for isolation and independent scaling — sandboxing a risky executor, placing a GPU-bound agent on different hardware, letting one agent type scale out — not for shaving milliseconds off in-process message passing.
  • What breaks first when a worker process restarts mid-conversation?
    The agent instances it held, along with their in-memory state: instances are created lazily per AgentId in a specific worker, so a restart loses whatever they accumulated. Direct sends in flight surface as transport failures, and published events aimed at that worker are simply missed, since publish is fire-and-forget with no durable queue. Anything that must survive has to live in an external store or be reconstructible from a persisted record.

saying these in an interview costs you the question

  • Assuming any message type just works across processes
  • Expecting shared in-memory caches to still be shared
  • Treating distribution as a latency optimization
  • Believing handler exceptions still surface as local tracebacks
  • Presenting the gRPC runtime as the production default

context