Events and processors

Convert a Harness Run into AG-UI events, read snapshots, and apply Host processors.

Use one HarnessAguiObserver per (thread_id, run_id) and feed it each public source item in order. It converts observations, not private state. Start with the offline example, then add your Host's transport and persistence.

Observe a Harness Run

Create one observer for each root or exposed child Run, then route every public source item by its Thread and Run correlation:

from a13n_stream_protocol import HarnessAguiObserver

observers: dict[tuple[str, str], HarnessAguiObserver] = {}

async with executable.stream(input_value, bindings=bindings) as run_stream:
    async for item in run_stream:
        correlation = (item.thread_id, item.run_id)
        observer = observers.get(correlation)
        if observer is None:
            observer = HarnessAguiObserver()
            observers[correlation] = observer

        new_events = observer.observe(item)
        await host.persist_and_publish(new_events)

snapshots = {
    correlation: observer.snapshot()
    for correlation, observer in observers.items()
}

A parent Harness stream can forward inline-child observations while preserving each child's own thread_id and run_id. Routing first prevents a root observer from rejecting a child Run. If an integration guarantees that its source contains exactly one Run, it can use one HarnessAguiObserver directly.

observe() returns only the AG-UI events produced by that source item. A single source item can produce several events—for example, a text part with initial content produces both TEXT_MESSAGE_START and TEXT_MESSAGE_CONTENT.

The first successfully observed item binds observer.thread_id and observer.run_id. Later items passed to that observer must have the same correlation. Use another observer when the Harness starts another Run, even when both Runs advance the same Thread.

Event mappings

The observer uses standard AG-UI events when the semantics match directly:

Harness observationAG-UI output
Text part start, delta, and endTEXT_MESSAGE_START, TEXT_MESSAGE_CONTENT, TEXT_MESSAGE_END
Reasoning start, content, signature, and endReasoning message events and REASONING_ENCRYPTED_VALUE
Completed tool callTOOL_CALL_START, TOOL_CALL_ARGS, TOOL_CALL_END
Successful tool resultTOOL_CALL_RESULT
Completed Run resultRUN_FINISHED
Failed or cancelled Run resultRUN_ERROR
Native Pydantic AI CapabilityEventCUSTOM event named by the native kind
Suspended result or another unmatched public observationNamespaced CUSTOM event

Unmatched Harness extensions use names such as a13n.harness.lifecycle. Unmatched Pydantic AI events use names such as a13n.pydantic_ai.final_result. A native CapabilityEvent retains its concrete kind, Capability ID, optional Tool-call correlation, and public payload. Every custom value also retains the public Thread, Run, sequence, timestamp, and source-event representation.

Content and large custom events

Capability events preserve their native name and payload, including user-defined kinds. A file edit remains before/after data, a summary remains summary data, and a shell status remains a status observation. The protocol does not turn them into assistant answers or pre-rendered panels. Clients decide their presentation. Actual model input uses user-role text events with shared ContentMetadata; normal displays omit content marked display: false.

Custom events larger than 48 KiB use generic a13n.stream.fragment frames. Reassemble them before inspecting the original event:

from a13n_stream_protocol import CustomEventAssembler

assembler = CustomEventAssembler()

# For each CUSTOM payload in one live subscription:
complete = assembler.accept(payload)
if complete is not None:
    render_custom(complete["name"], complete["value"])

The framing preserves the complete JSON structure rather than truncating the source. The assembler bounds pending content and rejects incomplete or inconsistent sequences; check assembler.gap and replace the assembler when resetting a subscription. These are best-effort observations, not a durable event log. See the framing contract for limits and fields.

Serialize events

The returned values are typed ag_ui.core.Event models. Use the upstream Pydantic adapter when a transport or store needs JSON-compatible values:

from ag_ui.core import Event
from pydantic import TypeAdapter

EVENT_ADAPTER = TypeAdapter(Event)

payloads = [
    EVENT_ADAPTER.dump_python(event, mode="json", by_alias=True)
    for event in new_events
]

The Host should wrap serialized events in its own durable IDs, ordering records, and delivery metadata rather than rewriting protocol correlation.

Read the Accumulated Snapshot

snapshot() returns detached copies of all retained post-processor events observed so far:

events = observer.snapshot()

Mutating an event returned by observe() or snapshot() does not mutate the observer. The snapshot is process-local convenience state, not a durable event log or Harness continuation value.

The observer does not compact streaming chunks or enforce a retention limit. A long-running Host should persist incremental results and apply its own bounded retention or projection policy.

Apply a Host Processor

A processor can omit an event or replace approved content fields before the event is accumulated and returned:

from typing import Any

from ag_ui.core import Event
from ag_ui.core.events import CustomEvent, TextMessageContentEvent
from a13n_harness import HarnessStreamEvent
from a13n_stream_protocol import HarnessAguiObserver


def process_event(
    source: HarnessStreamEvent[Any],
    event: Event,
) -> Event | None:
    del source

    if (
        isinstance(event, CustomEvent)
        and event.name == "a13n.harness.diagnostic"
    ):
        return None

    if isinstance(event, TextMessageContentEvent):
        return event.model_copy(update={"delta": event.delta.strip("\x00")})

    return event


observer = HarnessAguiObserver(processor=process_event)

A replacement must preserve the AG-UI event type and structural correlation, including message, tool, Thread, Run, lifecycle, and source fields. Invalid replacements raise AguiObservationError without committing the current source item.

A processor used with resume() must be replay-stable: the same ordered source history and stable Host configuration must produce the same retained event sequence. It must not retain mutable processing state or perform persistence, publication, acknowledgements, or other externally observable effects.

On this page