Streaming Internals
Streaming Internals
Section titled “Streaming Internals”ProxAI treats streaming as a byte-carrier problem plus protocol-specific event observation. The upstream SSE bytes are preserved, while provider observers scan events, maintain state, and emit compact diagnostics.
Layering
Section titled “Layering”- 1pipeline/upstream_response.rsDetects successful SSE responses and dispatches to provider streaming handling.
- 2provider/<protocol>/response/streaming.rsConstructs the protocol-specific BodyObserver and returns an outbound streaming body.
- 3upstream/streaming.rsWraps reqwest bytes_stream, records metrics, invokes observer hooks, and enforces idle timeout.
translation/ does not perform provider byte-stream observation or compatibility repair. Streaming observation stays inside application provider/; sse_translation.rs adapts SSE bytes to StreamTranslationInput values and invokes core ResponsePipeline, which composes provider normalization with structured translation.
BodyObserver lifecycle
Section titled “BodyObserver lifecycle”pub(crate) trait BodyObserver: Send + Unpin + 'static { fn on_chunk(&mut self, _chunk: &[u8]) -> BodyAction; fn on_stream_error(&mut self, error: &reqwest::Error); fn poll_pending_action(&mut self, _cx: &mut Context<'_>) -> BodyAction; fn on_stream_finished(&self, head: &UpstreamResponseHead, stats: UpstreamBodyStreamStats);}- 1Pull bytesRead the next chunk from reqwest response bytes_stream().
- 2Record carrier metricsUpdate common upstream stream byte/chunk stats.
- 3Scan SSE eventsUse SseEventScanner once per chunk inside the provider observer.
- 4Update protocol stateFeed parsed events into the provider state machine.
- 5Maybe interveneMost chunks continue unchanged; rare semantic failures inject an SSE error and close.
- 6Finish snapshotOn EOF/error/timeout/drop, emit a compact stream outcome snapshot.
Common provider observer shape
Section titled “Common provider observer shape”The three provider protocols now follow the same structure:
| File | Responsibility |
|---|---|
streaming.rs | Carrier hook implementation and lifecycle checks |
state.rs | Protocol event state and summary/outcome projection |
SseEventScanner | Chunk-to-event scanning, held by the observer |
Protocol event timelines
Section titled “Protocol event timelines”OpenAI Responses
Section titled “OpenAI Responses”- 1
response.createdStart of a Responses stream and initial response metadata.
- 2
response.output_item.addedA new output item such as message, reasoning, function_call, or MCP item.
- 3
response.output_text.delta / response.function_call_arguments.deltaIncremental text or tool-call argument bytes.
- 4
response.output_text.done / response.function_call_arguments.doneSemantic completion of a content part or tool-call arguments.
- 5
response.output_item.doneOutput item is complete.
- 6
response.completedTerminal event for a complete stream.
OpenAI Responses has extra tool-call semantics. If tool-call arguments start but never finish, ProxAI can inject a Responses-style SSE error event and close the stream.
OpenAI Chat Completions
Section titled “OpenAI Chat Completions”- 1
chat.completion.chunkDelta chunk containing role/content/tool_calls/finish_reason updates.
- 2
[DONE]Terminal sentinel. EOF without this sentinel is treated as a closed/incomplete stream.
Anthropic Messages
Section titled “Anthropic Messages”- 1
message_startStarts an Anthropic message stream.
- 2
content_block_startStarts a text, thinking, tool_use, or other content block.
- 3
content_block_deltaIncremental content for the current block.
- 4
content_block_stopCompletes a content block.
- 5
message_deltaCarries stop_reason, stop_sequence, usage, and other message-level deltas.
- 6
message_stopTerminal event for a complete Anthropic stream.
Cross-protocol translation boundary
Section titled “Cross-protocol translation boundary”The application and reusable translation core meet at a structured-event boundary:
provider/upstream byte observation → sse_translation.rs parses SSE frames into StreamTranslationInput → core ResponsePipeline normalizes provider events and translates StreamEvent values → sse_translation.rs encodes target SSE frames → pipeline observes failures and renders client-facing errorsRaw chunks, triggering SSE frames, HTTP metadata, diagnostics, and client error rendering remain in the application layer. Carrier-independent provider compatibility policy and structured payload/event normalization live in crates/proxai-core/src/provider/; cross-protocol conversion remains in crates/proxai-core/src/translation/. The core returns semantic failures through Stream<Item = Result<StreamEvent, StreamTranslationError>>; the shared core Observer reports typed compatibility adaptations and legal-but-lossy translations, and is not an error channel.
Translation stream lifecycle
Section titled “Translation stream lifecycle”Cross-protocol streaming translation uses a protocol-neutral four-phase event state machine in translation/stream.rs. Translator owns the route and observer; Translator::translate_stream consumes it, privately creates the selected pair state inside the returned stream, and is the only public structured-stream driver:
| Phase | Meaning |
|---|---|
`Waiting` | No semantic source message/chunk has initialized the translated stream yet. |
`Streaming(StreamingPhase<S>)` | Source deltas are active; pair-private state and target-representable output tracking live in StreamingPhase<S>. |
`Terminal(T)` | The source protocol has emitted its semantic terminal signal, but the final source stop/end marker has not been fully consumed. |
`Stopped` | The translator has emitted its final target stream output. |
Terminal(T) is intentionally generic because source protocols terminate differently:
| Inbound source | Terminal payload | Reason |
|---|---|---|
Anthropic Messages | StreamingPhase<S> | message_delta carries terminal stop/usage semantics, but message_stop is still required. The translator keeps the full phase until message_stop consumes it. |
OpenAI Chat Completions | pair-specific T | finish_reason ends semantic deltas. Later [DONE], EOF, or usage-only chunks use a target-specific pending terminal: PendingAnthropicTerminal for Anthropic output, or PendingResponsesTerminal for Responses output. |
OpenAI Responses | S | Responses terminal events (response.completed, response.incomplete, response.failed) carry response-level terminal semantics directly. There is no extra message_stop; the wrapper moves the streaming state into terminal so pair translators can finalize target output. |
Protocol wrappers such as anthropic_messages::streaming, openai_chat_completions::streaming, and openai_responses::streaming keep source-specific event ordering and error wording. Pair translators remain responsible for target representability checks, output item/block emission, and terminal flush behavior.
All three wrappers use typed protocol event parsing as the source-event allowlist:
| Inbound source | Typed event model | Extra source lifecycle checks |
|---|---|---|
Anthropic Messages | MessageStreamEvent | Requires message_start before semantic events, message_delta.stop_reason before message_stop, and rejects events after message_stop. |
OpenAI Chat Completions | CreateChatCompletionStreamResponse | First assistant chunk initializes the stream; later chunks must keep the same identity; finish_reason closes semantic deltas before [DONE]/EOF. |
OpenAI Responses | ResponseStreamEvent | response.created / response.in_progress initializes the stream; the source wrapper validates output_index → kind/item_id registration and completion; response terminal events move the lifecycle to terminal. |
ResponsesInboundLifecycle also owns the source-side output registry. response.output_item.added binds each output_index to its item kind and optional item id; subsequent content, reasoning, and tool events must carry the same item_id and expected kind, and response.output_item.done must close the same item. Pair state remains responsible only for target block/index projection after these source invariants pass.
Stream identity ownership
Section titled “Stream identity ownership”Streaming identity (StreamIdentity) belongs to the protocol-neutral inbound lifecycle state machine, while source-protocol wrappers decide when to initialize or validate it.
| Inbound source | Where identity comes from | Why it is handled this way |
|---|---|---|
Anthropic Messages | message_start.message.id and message_start.message.model | Anthropic emits identity once at the explicit message_start semantic boundary. Later content_block_*, message_delta, and message_stop events do not repeat id/model, so the Anthropic wrapper initializes lifecycle identity at message_start and then relies on lifecycle ordering to prevent duplicate starts. |
OpenAI Chat Completions | Every chat.completion.chunk repeats id and model | Chat has no separate message_start event. The first semantic chunk initializes lifecycle identity, and later chunks are checked against it so pieces from different upstream responses are never merged into one translated message/response. |
OpenAI Responses | response.created.response.id / model or response.in_progress.response.id / model | Responses has explicit response-level lifecycle events. The wrapper initializes identity at response.created or response.in_progress; later item/content events rely on that response identity while pair-local state tracks target output assembly. |
InboundStreamLifecycle stores identity outside the phase enum because identity is stream-envelope metadata that remains stable across Streaming, Terminal, and Stopped. Pair-local stream state should keep only target-conversion state or derived identifiers, not another copy of the full source identity.
BodyAction
Section titled “BodyAction”pub(crate) enum BodyAction { Continue, InjectAndClose(Bytes),}Most observers return Continue. InjectAndClose is reserved for cases where ProxAI can produce a better client-facing stream failure than silently hanging or closing without context.
Translation pair internals
Section titled “Translation pair internals”Cross-protocol streaming code separates source lifecycle, pair state, and target construction:
translation/<source>/streaming.rs or streaming/mod.rs source protocol parsing, identity, and lifecycle orderingtranslation/<source>/streaming/<dialect>.rs optional source-specific compatibility dialect when a directory module is needed
translation/<source>/to_<target>/streaming/├── mod.rs # routes source events through pair state and owns event/terminal timing├── state.rs # validates/accumulates pair state and returns typed target events└── output.rs # optional pair-local pure helpers when terminal assembly is substantial
translation/<target>/outbound/ shared stateless constructors for target request/response/stream wire typesResponsibility boundaries
Section titled “Responsibility boundaries”- source lifecycle wrappers own source event parsing, envelope identity, and phase ordering.
- pair state owns atomic source-to-target transitions: validation, accumulation, target index allocation, and the typed target events produced by that transition.
- target outbound modules contain reusable stateless constructors. They do not inspect source events or hold streaming state.
- pair mod.rs dispatches source events, invokes pair state, handles terminal outcomes, and maps typed target events to protocol-neutral
StreamEvents. - pair output.rs is optional. Keep it only when pair-local terminal construction is substantial enough to improve readability; do not duplicate constructors already owned by the target
outbound/module.
This makes state methods small transducers rather than passive bags of fields. If a method mutates pair state and returns target events, both operations must describe one atomic protocol transition. Carrier encoding and client-facing error rendering remain outside pair state.
Pair-private types
Section titled “Pair-private types”Use a pair-local types.rs when multiple pair modules share accumulation records, such as StreamTextItem and StreamToolItem. Keep protocol-native wire constructors in the target outbound/ module instead of moving them into pair-private types.
Terminal state type choice
Section titled “Terminal state type choice”InboundStreamLifecycle<S, T> has two type parameters: S is the streaming-phase state type, T is the terminal-phase state type. They may be the same or different, depending on how much data the terminal events need.
| Pair | S (streaming) | T (terminal) | Why |
|---|---|---|---|
chat → ant | ChatStreamingState | PendingAnthropicTerminal | Keeps finish_reason, refusal, and the latest usage snapshot until [DONE], EOF, or a trailing usage-only chunk finalizes message_delta + message_stop. |
chat → responses | StreamingState | PendingResponsesTerminal | Keeps the full accumulated state plus finish_reason until *.done events and the final completed/incomplete response snapshot are emitted. |
Root cause is the timing and data requirements of terminal event construction:
- Anthropic target has an explicit block lifecycle (
content_block_stop); each block is closed and emitted during streaming, so byfinish_reasononly a few snapshot values remain — a lightweight terminal type suffices. - Responses target has no block-close equivalent; text deltas and tool arguments stream until
finish_reason, when*.done+response.completedare emitted all at once. Data must stay in state until the end, so terminal must be the full streaming state.