Skip to content

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.

  1. 1pipeline/upstream_response.rsDetects successful SSE responses and dispatches to provider streaming handling.
  2. 2provider/<protocol>/response/streaming.rsConstructs the protocol-specific BodyObserver and returns an outbound streaming body.
  3. 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.

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);
}
  1. 1Pull bytesRead the next chunk from reqwest response bytes_stream().
  2. 2Record carrier metricsUpdate common upstream stream byte/chunk stats.
  3. 3Scan SSE eventsUse SseEventScanner once per chunk inside the provider observer.
  4. 4Update protocol stateFeed parsed events into the provider state machine.
  5. 5Maybe interveneMost chunks continue unchanged; rare semantic failures inject an SSE error and close.
  6. 6Finish snapshotOn EOF/error/timeout/drop, emit a compact stream outcome snapshot.

The three provider protocols now follow the same structure:

FileResponsibility
streaming.rsCarrier hook implementation and lifecycle checks
state.rsProtocol event state and summary/outcome projection
SseEventScannerChunk-to-event scanning, held by the observer
  1. 1
    response.created

    Start of a Responses stream and initial response metadata.

  2. 2
    response.output_item.added

    A new output item such as message, reasoning, function_call, or MCP item.

  3. 3
    response.output_text.delta / response.function_call_arguments.delta

    Incremental text or tool-call argument bytes.

  4. 4
    response.output_text.done / response.function_call_arguments.done

    Semantic completion of a content part or tool-call arguments.

  5. 5
    response.output_item.done

    Output item is complete.

  6. 6
    response.completed

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

  1. 1
    chat.completion.chunk

    Delta chunk containing role/content/tool_calls/finish_reason updates.

  2. 2
    [DONE]

    Terminal sentinel. EOF without this sentinel is treated as a closed/incomplete stream.

  1. 1
    message_start

    Starts an Anthropic message stream.

  2. 2
    content_block_start

    Starts a text, thinking, tool_use, or other content block.

  3. 3
    content_block_delta

    Incremental content for the current block.

  4. 4
    content_block_stop

    Completes a content block.

  5. 5
    message_delta

    Carries stop_reason, stop_sequence, usage, and other message-level deltas.

  6. 6
    message_stop

    Terminal event for a complete Anthropic stream.

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 errors

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

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:

PhaseMeaning
`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 sourceTerminal payloadReason
Anthropic MessagesStreamingPhase<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 Completionspair-specific Tfinish_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 ResponsesSResponses 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 sourceTyped event modelExtra source lifecycle checks
Anthropic MessagesMessageStreamEventRequires message_start before semantic events, message_delta.stop_reason before message_stop, and rejects events after message_stop.
OpenAI Chat CompletionsCreateChatCompletionStreamResponseFirst assistant chunk initializes the stream; later chunks must keep the same identity; finish_reason closes semantic deltas before [DONE]/EOF.
OpenAI ResponsesResponseStreamEventresponse.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.

Streaming identity (StreamIdentity) belongs to the protocol-neutral inbound lifecycle state machine, while source-protocol wrappers decide when to initialize or validate it.

Inbound sourceWhere identity comes fromWhy it is handled this way
Anthropic Messagesmessage_start.message.id and message_start.message.modelAnthropic 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 CompletionsEvery chat.completion.chunk repeats id and modelChat 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 Responsesresponse.created.response.id / model or response.in_progress.response.id / modelResponses 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.

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.

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 ordering
translation/<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 types
  • 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.

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.

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.

PairS (streaming)T (terminal)Why
chat → antChatStreamingStatePendingAnthropicTerminalKeeps finish_reason, refusal, and the latest usage snapshot until [DONE], EOF, or a trailing usage-only chunk finalizes message_delta + message_stop.
chat → responsesStreamingStatePendingResponsesTerminalKeeps 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 by finish_reason only 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.completed are emitted all at once. Data must stay in state until the end, so terminal must be the full streaming state.