Adapt the Event Stream to a Product Surface¶
The terminal renderer accepts Python dataclasses. A browser does not. Passing a
SessionEndEvent directly to send_json() fails on enums, exceptions, and
binary output.
Do not change the agent to solve a delivery problem. Add one adapter at the edge of the harness.
Outcome¶
The same CloudHarness.stream_turn() output can drive a terminal, WebSocket,
or HTTP stream through a versioned JSON envelope.
Fast Track¶
Walk each event dataclass without copying unknown field values.
Normalize enums, exceptions, bytes, and nested containers.
Add the concrete event class name and a wire-format version.
Bind authenticated identity to the internal session ID.
Propagate disconnect cancellation to the active
AgentStream.
Hands-on delta¶
1. Build the adapter¶
Axio intentionally does not define a public web protocol. Applications differ in authentication, reconnect behavior, binary transport, and compatibility requirements. This small adapter is one explicit application contract:
def json_value(value: Any) -> Any:
if isinstance(value, Enum):
return value.value
if isinstance(value, BaseException):
return {
"code": "internal_error",
"message": "The agent turn failed.",
}
if isinstance(value, bytes):
return {
"encoding": "base64",
"data": base64.b64encode(value).decode("ascii"),
}
if is_dataclass(value) and not isinstance(value, type):
return {
item.name: json_value(getattr(value, item.name)) for item in fields(value)
}
if isinstance(value, dict):
return {str(key): json_value(item) for key, item in value.items()}
if isinstance(value, (list, tuple)):
return [json_value(item) for item in value]
return value
def event_to_dict(event: StreamEvent) -> dict[str, Any]:
return {
"version": 1,
"type": type(event).__name__,
"data": json_value(event),
}
The field walk is intentionally shallow before json_value() recurses. In
contrast, dataclasses.asdict() deep-copies unknown values and can fail on a
provider exception that contains a lock, stream, or client object.
Message.to_dict() serializes persisted conversation messages. It is not an
event serializer. Keep message storage and the public stream protocol as two
separate contracts.
Do not send str(exception) to an untrusted client. Provider responses,
credentials, paths, or query data can appear in exception text. Record the
full exception only in access-controlled server logs, with the application’s
normal secret-redaction policy.
2. Keep the endpoint small¶
The endpoint receives an authenticated identity, selects its server-issued session ID, and forwards events:
@dataclass(frozen=True, slots=True)
class AuthenticatedIdentity:
subject: str
session_id: str
def __post_init__(self) -> None:
if re.fullmatch(r"[0-9a-f]{32}", self.session_id) is None:
raise ValueError("session_id must be a server-issued UUID hex value")
async def handle_prompt(websocket, harness, identity, prompt):
stream = harness.stream_turn(identity.session_id, prompt)
async with aclosing(stream):
async for event in stream:
await websocket.send_json(event_to_dict(event))
aclosing() closes the async generator when send_json() fails or the request
is cancelled. This lets CloudHarness.stream_turn() close AgentStream and
release the per-session turn lock.
The authentication layer must issue and retain session_id. The example uses
a UUID hex value, which is also safe in Docker resource names. Never accept an
arbitrary client string as proof that the caller owns a context or sandbox.
3. Preserve meaning across surfaces¶
Different surfaces render the same event differently:
Event |
Terminal |
Browser |
Service log |
|---|---|---|---|
|
append text |
append to message |
usually omit |
|
show tool name |
open tool card |
record start time |
|
show status |
complete tool card |
record duration and error flag |
|
show failure |
show recoverable error |
record exception metadata |
|
show usage |
close turn |
record stop reason and totals |
The adapter can omit internal fields or split binary data into another channel.
Once clients depend on that choice, change the version when compatibility
breaks.
Try It¶
Run uv run python examples/tutorial/adapt_the_event_stream.py from the
repository root. It checks enum, exception, and binary serialization without a
server or model API.
Then send those envelopes through your selected HTTP or WebSocket framework. Disconnect during a tool call and verify that the session lock is released for the next request.
Done when¶
[ ] Every Axio event becomes JSON without provider-specific logic.
[ ] Enum, exception, and binary values have explicit wire representations.
[ ] The envelope has a compatibility version.
[ ] Authentication, not client input, selects the internal session.
Next failure¶
The product surface now receives stable events, but a later refactor can still break dispatch, ordering, or cleanup silently. The final lesson turns the complete harness boundary into a deterministic contract check.
Continue with Test the Harness, Not the Model.