axio-sse¶
Read text/event-stream: a decoder you feed, and a reader for what its payloads mean.
The package knows nothing about HTTP and imports no client. It has no dependencies, not
even on axio.
Two shapes of stream, two ways to read one. A stream whose events are all one shape needs
nothing above payloads(). A stream that says what each event is subclasses Reader and
writes one @on(...) method per event.
from collections.abc import Iterator
from dataclasses import dataclass
from axio_sse import Decoder, Event, Handled, Payload, Reader, Wire, on
# Chunks are cut wherever the network cut them; nothing dispatches until the blank line.
decoder = Decoder()
assert decoder.decode('event: delta\ndata: {"type": "message.delta", ') == []
events = decoder.decode('"text": "hi"}\n\n')
assert [event.event for event in events] == ["delta"]
@dataclass(frozen=True, slots=True)
class Delta(Wire, name="message.delta"):
text: str = ""
index: int = 0
class Messages(Reader[str]):
@on(Delta)
def _delta(self, wire: Delta) -> Iterator[str]:
yield wire.text
def unmatched(self, name: str, payload: Payload) -> Handled[str]:
return [f"<{name}>"]
reader = Messages()
assert Messages.names() == {"message.delta"}
assert reader.read(events[0]) == ["hi"]
# Not named, so not interpreted - and forwarded rather than dropped.
assert reader.read(Event(data='{"type": "message.tool.started"}')) == ["<message.tool.started>"]
What a Reader names, and what it does not¶
A Reader names only the events it interprets. Everything else reaches unmatched(),
which returns nothing by default. A reader overrides it to forward instead of drop. Both
readers in this repository do that, as ProviderEvent under the provider’s own name.
The reason is not obvious. The instinct it contradicts - add a handler for every event in the provider’s documentation - is the wrong one. An endpoint that runs tools publishes one event family per tool. That set therefore depends on which tools exist and which the caller declared, not on the protocol. Named one by one, the list is stale the day a tool is added. A new tool then reads as news about the protocol when it is news about the tools.
Reading with strict=True still raises UnknownEvent for any name no method claims. A
test can therefore hold names() against the schema the provider publishes, without the
reader carrying a list it cannot keep true.
Reading a stream¶
- async axio_sse.payloads(chunks: AsyncIterable[bytes | str], *, until: str = '') AsyncIterator[Payload][source]¶
The JSON object of every event in this stream. Comments, keep-alives and junk do not arrive.
All a stream with no discriminator needs: its events are one shape, read field by field.
- async axio_sse.events(chunks: AsyncIterable[bytes | str], *, until: str = '') AsyncIterator[Event][source]¶
Every event in this stream, as the chunks arrive.
Chunks must carry their line terminators:
aiter_lines()strips them, so nothing dispatches. A stream that stops without its final blank line still yields what it collected.untilnames the data payload that closes the stream, and is not yielded.
- class axio_sse.Event(data: str = '', event: str = '', id: str = '', retry: int | None = None)[source]¶
One dispatched event, with the four fields the format defines.
- event: str¶
What the
event:field carried, empty where the stream sent none.
- id: str¶
The stream position for a client that reconnects, not an id of this event.
- property name: str¶
The event’s type. An unnamed event is of type
message, which the format defines.Dispatched on the raw field instead, an
@on("message")handler never runs for the ordinary unnamed event, and a strict read rejects it as unknown.
- payload() Payload | None[source]¶
This event’s JSON object, or None where the event carries no data at all.
Data that will not parse raises. Skipped instead, a text or tool-call event whose JSON arrived corrupt was dropped, the completion event after it still reported success, and the caller got a short answer or half a tool’s arguments with nothing saying anything was lost.
A sentinel such as
[DONE]is data too, and reaches here as junk. Name it inuntil, which ends the stream before it is read.
The format, with no loop¶
- class axio_sse.Decoder(limit: int = 33554432)[source]¶
The format as a state machine: feed it chunks, take the events they completed.
Same shape as
codecs.IncrementalDecoder:decode(chunk, final)andreset(). The problem is the same one. Input is cut at arbitrary points, and output only sometimes completes. It takes chunks and never lines.aiohttp’sreaduntilraisesLineTooLongpast 131072 bytes, andLineTooLongis not aClientError. A large reasoning event killed a turn with no answer.Held text costs time for its size, and never for its square. Chunks with no terminator wait in a list. A read line is left behind rather than sliced out. A scanned tail is never scanned twice.
- decode(chunk: bytes | str = b'', final: bool = False) list[Event][source]¶
Every event this chunk completed.
final=Truecloses the stream: what is still pending is discarded, which the format requires of an event that never reached its blank line. Without that last call a stream cut mid-character keeps the half character instead of replacing it.
What the payloads mean¶
- class axio_sse.Wire[source]¶
One payload shape, named by the wire name it arrives under:
@dataclass(frozen=True, slots=True) class OutputTextDelta(Wire, name="response.output_text.delta"): delta: str = "" output_index: int = 0
Every field is read by its declared name and type, so a misspelled key is a type error at the place that uses it rather than a default quietly standing in for the value. A field the provider did not send, sent as null, or sent as the wrong type takes its default. That is what an optional provider field is, and one bad field must not lose the whole event.
A nested object is another
Wire; a list of them islist[ThatWire]. Give a shape noname=and it is only ever nested, never dispatched to.Declare a field
raw: Payloadand it receives the whole payload, for a shape that varies too much to declare whole. A citation arrives under five shapes and each names its span differently, so the fields worth reading are declared and the rest travels beside them.Declaring a shape registers it nowhere. A
Readerclaims it with@on(ThatShape).- names: ClassVar[tuple[str, ...]] = ()¶
Every name this shape arrives under, from
name=andalso=on the class line.
- class axio_sse.Reader[source]¶
What one endpoint sends, as one method per event.
- class Messages(Reader[StreamEvent], by=EVENT_NAME):
@on(“content_block_delta”) def _delta(self, payload: Payload) -> Iterator[StreamEvent]: …
One instance reads one stream. The turn’s running totals and id maps live on
selfinstead of travelling through a call. A reader used for a second response would carry the first one’s state into it. Construct one per response. Being that state, a reader must not be frozen.readlatches the caller’sstrictonself, and@dataclass(frozen=True)refuses that assignment. Withslots=Trueas well, the refusal comes from inside the rebuilt class and says only thatsuper()got the wrong type.@dataclass(slots=True)alone is fine.A handler returns what the event became — an iterable, or None where the event only moved that state.
bynames the payload key that holds the event’s name, orEVENT_NAMEfor the format’s ownevent:field. A subclass that does not givebyinherits it.- classmethod names() frozenset[str][source]¶
Every name this reader claims, for a test to hold against the provider’s own list.
- async over(chunks: AsyncIterable[bytes | str], *, strict: bool = False, until: str = '') AsyncIterator[T][source]¶
Read a whole stream of chunks, yielding what each event became.
An event that becomes nothing yields nothing, so the outputs do not line up with the events.
- read(event: Event, *, strict: bool = False) list[T][source]¶
Everything this one event became, empty where it became nothing.
- unknown(name: str) None[source]¶
The one policy for a name nothing here reads: DEBUG, or refuse under
strict.A handler calls it for a second discriminator inside one event, such as a block that names the kind of its own chunks. A nested name nobody read then fails the same replay a new event fails, instead of disappearing.
- unmatched(name: str, payload: Payload) Handled[T][source]¶
What a payload no method claims becomes. Nothing, unless a reader says otherwise.
Override it to forward instead of drop. That is what a stream whose vocabulary grows on its own needs. An endpoint that runs tools names an event per tool. That set is a function of which tools exist and which were asked for, not of the protocol. Naming them one by one makes the reader stale the day a tool is added. It also reports a new tool as news about the protocol, when it is news about the tools.
Name here only what this reader interprets.
strictstill refuses anything unnamed. A test can therefore hold the interpreted set against the schema, and the reader carries no list it cannot keep true.
- axio_sse.on(*claimed: str | type[Wire]) Callable[[_Handler], _Handler][source]¶
Give a
Readermethod the payloads it reads.Give it a
Wireshape and the method is handed that shape, its fields read by declared name and type. Give it wire names and the method is handed thePayloaditself, which is what a method that only forwards an event wants. Declaring a shape for a payload nobody reads a field of would be a schema written for nothing.Several names on one method is how a stream that sends one thing under two names is written. Both stay in the class body, so
stricthas nothing to fire on and no second list exists to keep in step with the first.
EVENT_NAMEThe
by=sentinel,"event:".by=EVENT_NAMEdispatches on the format’s ownevent:field. Any otherbynames a key inside the payload. A JSON key holding a colon is not a name a provider gives a field, so the two can never mean each other.Handled[T]What a handler returns:
Iterable[T] | None. One rule for none, one, and many — an event that only moved the reader’s own state returns nothing.
- class axio_sse.Payload[source]¶
The JSON object inside one event, read by path.
payload.number("message", "usage", "input_tokens")walks the path and gives the default wherever a step is missing, null, or the wrong type — which is what an optional provider field is. It is adict, sopayload["x"],in, andjson.dumpsall still work. The four readers exist so a handler carries noAnyand no chain of.get({}).- number(*keys: str, default: int = 0) int[source]¶
The whole number at this path, or the default where the provider sent none.
- obj(*keys: str) Payload[source]¶
The object at this path, empty where there is none, so a path can be walked in steps.