# How to Implement Event-Driven Asynchronous Agent Architecture: A Production-Ready Guide

> Implement event driven asynchronous agent architecture using asyncio and uniform event messages. Learn production ready techniques from bojieli ai agent book repository for scalable AI agents.

- Repository: [Bojie Li/ai-agent-book](https://github.com/bojieli/ai-agent-book)
- Tags: how-to-guide
- Published: 2026-08-06

---

**Build scalable AI agents using `asyncio`, uniform event messages, and callback-based orchestration as demonstrated in the bojieli/ai-agent-book repository.**

The *ai-agent-book* repository provides a complete reference implementation for building event-driven asynchronous AI agents that handle concurrent tasks without blocking. This architecture wraps all interactions—LLM calls, tool invocations, and state changes—in a standardized `Message` event format, then processes them through an `asyncio`-based execution kernel with full observability.

## What Is Event-Driven Asynchronous Agent Architecture?

At its core, this architecture decouples *what* an agent does from *when* it reacts. Instead of synchronous request-response chains, every operation emits an event that subscribers consume at their own pace. The repository implements this through three interconnected layers: a **uniform event abstraction**, an **async execution engine**, and **event-driven orchestration** components.

The benefits are substantial: multiple agents run concurrently in a single event loop, UI updates stream in real-time, and OpenTelemetry tracing propagates through every event for end-to-end observability.

## The Event Abstraction Layer

All interactions flow through `aworld.core.event.base.Message`, defined in [`chapter8/gaia-experience/AWorld/core/event/base.py`](https://github.com/bojieli/ai-agent-book/blob/main/chapter8/gaia-experience/AWorld/core/event/base.py). This class standardizes payloads, headers, and tracing information across the system.

Key attributes include:

- `event.id`: Unique identifier for distributed tracing
- `topic`: Event category (e.g., `Constants.PLAN` for planning messages)
- `payload`: The actual data (LLM tokens, tool results, user input)
- `kind`: Event type discriminator (`"tool_call"`, `"llm_response"`, etc.)

```python
from aworld.core.event.base import Message

# Creating a message with tracing context

msg = Message(
    payload={"query": "search Wikipedia for Python"},
    topic="tool.execute",
    kind="tool_call"
)

# event.id is auto-generated; tracing headers attach automatically

```

## The Async Execution Kernel

The execution engine lives in [`experience_agent.py`](https://github.com/bojieli/ai-agent-book/blob/main/experience_agent.py) and leverages `asyncio` for non-blocking operation. The `Runners.streamed_run_task()` method yields events as they occur rather than buffering until completion.

```python
from aworld.core.runner import Runners

async def execute_with_streaming(task: dict):
    """Run a task and process events as they arrive."""
    runner = Runners.streamed_run_task(task)
    
    async for event in runner.stream_events():
        # event is a Message instance

        yield event.kind, event.payload

```

This pattern enables **real-time responsiveness**—UI elements update as the LLM generates tokens, not after the full response completes.

## Event-Driven Orchestration with Callbacks

Individual agents expose an `on_event` callback for external consumption. The `PhoneAgent` implementation in [`chapter9/phone-agent/agent.py`](https://github.com/bojieli/ai-agent-book/blob/main/chapter9/phone-agent/agent.py) demonstrates this pattern:

```python
from chapter9.phone_agent.agent import PhoneAgent

def ui_handler(kind: str, payload):
    """Push events to websocket, log, or update terminal UI."""
    if kind == "llm_response":
        print(f"🤖 {payload['token']}", end="", flush=True)
    elif kind == "tool_call":
        print(f"\n🔧 Calling tool: {payload['tool_name']}")

agent = PhoneAgent(on_event=ui_handler)

# Execute asynchronously

import asyncio
asyncio.run(agent.run("What's the weather in Tokyo?"))

```

The `on_event` signature is `Callable[[str, Any], None]` where the first argument is the event `kind` and the second is the deserialized payload.

## Coordinating Multiple Agents with Swarm Mode

For multi-agent scenarios, the `Swarm` class accepts `event_driven=True` to forward events from all members through a shared bus:

```python
from aworld.core.agent import Swarm

swarm = Swarm(
    research_agent,
    writer_agent,
    critique_agent,
    event_driven=True  # Enable unified event streaming

)

async def run_collaborative_task(query: str):
    async for event in swarm.run(query):
        # Route by topic to appropriate handler

        if event.topic == "research.finding":
            await update_knowledge_base(event.payload)
        elif event.topic == "writing.draft":
            await render_preview(event.payload)

```

Reference: [`chapter8/gaia-experience/AWorld/tests/test_state_manager.py`](https://github.com/bojieli/ai-agent-book/blob/main/chapter8/gaia-experience/AWorld/tests/test_state_manager.py) demonstrates swarm coordination patterns.

## Streaming Events to Web Clients via SSE

The repository includes a FastAPI-based server in [`gaia_agent_server.py`](https://github.com/bojieli/ai-agent-book/blob/main/gaia_agent_server.py) that exposes events through Server-Sent Events:

```python
from fastapi import FastAPI
from fastapi.responses import StreamingResponse
import json

app = FastAPI()

@app.get("/agent/stream")
async def stream_agent(query: str):
    async def event_generator():
        task = {"messages": [{"role": "user", "content": query}]}
        
        async for event in Runners.streamed_run_task(task).stream_events():
            # SSE format: data: <json>\n\n

            payload = json.dumps({
                "kind": event.kind,
                "payload": event.payload,
                "trace_id": event.id
            })
            yield f"data: {payload}\n\n"
    
    return StreamingResponse(
        event_generator(),
        media_type="text/event-stream"
    )

```

Clients connect to this endpoint and receive events in real-time without polling overhead.

## Observability with OpenTelemetry Tracing

Every `Message` carries tracing context through [`instrumentation/eventbus.py`](https://github.com/bojieli/ai-agent-book/blob/main/instrumentation/eventbus.py). When a message contains `traceparent` or `tracestate` headers, the instrumentation layer extracts and re-attaches them to downstream calls:

```python

# From aworld/trace/instrumentation/eventbus.py

def propagate_trace_context(message: Message) -> None:
    """Extract trace context from message headers and activate."""
    carrier = message.headers.get("trace_context", {})
    context = extract(carrier)
    
    # Attach to current span for downstream propagation

    with use_span(get_current_span(context)):
        yield

```

This ensures **end-to-end visibility** across LLM inference, tool execution, and agent handoffs.

## Minimal Implementation Checklist

1. **Define your Message schema** — Subclass or extend `aworld.core.event.base.Message` with domain-specific payload types.

2. **Wrap blocking I/O in async** — Convert all LLM and tool calls to `async def` coroutines that return `Message` instances.

3. **Consume with `async for`** — Use streaming iteration to process events as they arrive, not after completion.

4. **Register `on_event` callbacks** — Wire UI updates, logging, or side effects through the callback interface.

5. **Enable distributed tracing** — Include trace headers in all outbound messages; extract and propagate in inbound handlers.

6. **Scale with Swarm mode** — Set `event_driven=True` when coordinating multiple agents.

## Key Source Files

| File | Purpose |
|------|---------|
| [`aworld/core/event/base.py`](https://github.com/bojieli/ai-agent-book/blob/main/aworld/core/event/base.py) | `Message` base class and tracing helpers |
| [`experience_agent.py`](https://github.com/bojieli/ai-agent-book/blob/main/experience_agent.py) | Async runner with `streamed_run_task()` |
| [`chapter9/phone-agent/agent.py`](https://github.com/bojieli/ai-agent-book/blob/main/chapter9/phone-agent/agent.py) | `PhoneAgent` with `on_event` callback |
| [`aworld/trace/instrumentation/eventbus.py`](https://github.com/bojieli/ai-agent-book/blob/main/aworld/trace/instrumentation/eventbus.py) | OpenTelemetry context propagation |
| [`AWorld/examples/gaia/gaia_agent_server.py`](https://github.com/bojieli/ai-agent-book/blob/main/AWorld/examples/gaia/gaia_agent_server.py) | FastAPI SSE streaming endpoint |
| [`AWorld/tests/test_state_manager.py`](https://github.com/bojieli/ai-agent-book/blob/main/AWorld/tests/test_state_manager.py) | Swarm coordination examples |

## Summary

- **Uniform events**: The `Message` class in [`base.py`](https://github.com/bojieli/ai-agent-book/blob/main/base.py) standardizes all agent interactions
- **Async execution**: `Runners.streamed_run_task()` yields events through `async for` iteration
- **Callback architecture**: `on_event` hooks decouple agent logic from consumption
- **Swarm coordination**: `event_driven=True` enables multi-agent event streaming
- **Production observability**: OpenTelemetry tracing integrates automatically via [`eventbus.py`](https://github.com/bojieli/ai-agent-book/blob/main/eventbus.py)

## Frequently Asked Questions

### What makes event-driven architecture better than synchronous loops for AI agents?

**Event-driven agents scale horizontally without thread overhead.** The `asyncio` event loop in [`experience_agent.py`](https://github.com/bojieli/ai-agent-book/blob/main/experience_agent.py) handles thousands of concurrent connections, while synchronous code would block on each LLM call. Events also enable streaming UI updates and persistent audit logs that synchronous returns cannot provide.

### How does tracing work across agent boundaries?

**Trace context propagates through Message headers.** The [`instrumentation/eventbus.py`](https://github.com/bojieli/ai-agent-book/blob/main/instrumentation/eventbus.py) module extracts `traceparent` from incoming events and re-attaches it to outgoing calls. This creates a single trace spanning multiple agents, tools, and LLM invocations without manual context management.

### Can I mix event-driven and non-event-driven agents in the same system?

**Yes, but Swarm mode requires consistency.** Individual agents can run synchronously, but a `Swarm` with `event_driven=True` expects all members to yield `Message` events. The repository's `PhoneAgent` demonstrates a hybrid: it runs asynchronously internally but exposes a synchronous `run()` method for compatibility.

### What is the performance overhead of the event abstraction layer?

**Negligible for I/O-bound workloads.** The `Message` class adds ~50-100 bytes per event and zero serialization overhead for intra-process communication. For high-frequency scenarios, the repository recommends batching small events or using shared memory for large payloads while keeping headers in `Message`.