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

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. 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.)
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 and leverages asyncio for non-blocking operation. The Runners.streamed_run_task() method yields events as they occur rather than buffering until completion.

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 demonstrates this pattern:

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:

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 demonstrates swarm coordination patterns.

Streaming Events to Web Clients via SSE

The repository includes a FastAPI-based server in gaia_agent_server.py that exposes events through Server-Sent Events:

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. When a message contains traceparent or tracestate headers, the instrumentation layer extracts and re-attaches them to downstream calls:


# 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 Message base class and tracing helpers
experience_agent.py Async runner with streamed_run_task()
chapter9/phone-agent/agent.py PhoneAgent with on_event callback
aworld/trace/instrumentation/eventbus.py OpenTelemetry context propagation
AWorld/examples/gaia/gaia_agent_server.py FastAPI SSE streaming endpoint
AWorld/tests/test_state_manager.py Swarm coordination examples

Summary

  • Uniform events: The Message class in 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

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

Have a question about this repo?

These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:

Share the following with your agent to get started:
curl -s "https://instagit.com/install.md"

Works with
Claude Codex Cursor VS Code OpenClaw Any MCP Client

Maintain an open-source project? Get it listed too →