This document explains how token/tool streaming is normalized in @oh-my-pi/pi-ai, then propagated through @oh-my-pi/pi-agent-core and coding-agent session events.
End-to-end flow
streamSimple()(packages/ai/src/stream.ts) maps generic options and dispatches to a provider stream function.- Provider stream functions translate provider-native stream events into the unified
AssistantMessageEventsequence. Current built-ins include Anthropic, OpenAI Responses/Completions/Codex/Azure Responses, Google Gemini/Gemini CLI/Vertex, Bedrock Converse, Ollama, Cursor, pi-native gateway transport, plus GitLab Duo/Kimi/Synthetic/xAI-Grok-Responses wrappers and extension-registered custom APIs. - Each provider pushes events into
AssistantMessageEventStream(packages/ai/src/utils/event-stream.ts), which exposes:- async iteration for incremental updates
result()for finalAssistantMessage
agentLoop(packages/agent/src/agent-loop.ts) consumes those events, mutates in-flight assistant state, and emitsmessage_updateevents carrying the rawassistantMessageEvent.AgentSession(packages/coding-agent/src/session/agent-session.ts) subscribes to agent events, persists messages, drives extension hooks, and applies session behaviors (retry, compaction, TTSR, streaming-edit abort checks).
Unified stream contract in @oh-my-pi/pi-ai
All providers emit the same shape (AssistantMessageEvent in packages/ai/src/types.ts):
start- content block lifecycle triplets:
- text:
text_startβtext_delta* βtext_end - thinking:
thinking_startβthinking_delta* βthinking_end - tool call:
toolcall_startβtoolcall_delta* βtoolcall_end
- text:
- terminal event:
donewithreason: "stop" | "length" | "toolUse"- or
errorwithreason: "aborted" | "error"
AssistantMessageEventStream guarantees:
- final result is resolved by terminal event (
doneorerror) - events are delivered to consumers immediately, in push order (no batching or merging)
Delta throttling behavior
AssistantMessageEventStream itself no longer throttles or merges delta events β every provider event is delivered as pushed. The per-delta cost control moved into tool-call argument parsing: providers accumulate partial JSON and re-parse it via parseStreamingJsonThrottled() (packages/ai/src/utils/json-parse.ts), which skips the re-parse until at least STREAMING_JSON_PARSE_MIN_GROWTH (256) new bytes have arrived, bounding mid-stream parse cost from quadratic to linear. The final toolcall_end parse is always unconditional and authoritative.
There is no provider backpressure: providers still produce at full speed, while the local stream queues.
Provider normalization details
Anthropic (anthropic-messages)
Source: packages/ai/src/providers/anthropic.ts
Normalization points:
message_startinitializes usage (input/output/cache tokens)content_block_startmaps to text/thinking/toolcall startscontent_block_deltamaps:text_deltaβtext_deltathinking_deltaβthinking_deltainput_json_deltaβtoolcall_deltasignature_deltaupdatesthinkingSignatureonly (no event)
content_block_stopemits corresponding*_endmessage_delta.stop_reasonmaps viamapStopReason()
Tool-call argument streaming:
- each tool block carries internal
partialJson - every JSON delta appends to
partialJson argumentsare reparsed on appended deltas viaparseStreamingJsonThrottled()(re-parse only after β₯256 new bytes)toolcall_endreparses once more, then stripspartialJson
OpenAI Responses family (openai-responses, openai-codex-responses, azure-openai-responses)
Sources: packages/ai/src/providers/openai-responses.ts, openai-codex-responses.ts, and azure-openai-responses.ts
Normalization points:
response.output_item.addedstarts reasoning/text/function-call/custom-tool blocks- reasoning summary events (
response.reasoning_summary_text.delta) and raw reasoning events (response.reasoning_text.delta) becomethinking_delta - output/refusal deltas become
text_delta response.function_call_arguments.deltaandresponse.custom_tool_call_input.deltabecometoolcall_deltaresponse.output_item.doneemitsthinking_end/text_end/toolcall_endresponse.completedmaps status to stop reason and usage;response.failed/ SDKerrorevents throw into the wrapperβs terminalerrorpath
Tool-call argument streaming:
- same
partialJsonaccumulation pattern as Anthropic for function-call JSON arguments - custom tools stream raw string input and expose final arguments as
{ input: <raw> } - providers that send only
response.function_call_arguments.donestill populate final args - tool call IDs are normalized as
"<call_id>|<item_id>"
Google Generative AI (google-generative-ai)
Source: packages/ai/src/providers/google.ts (thin request wrapper) and google-shared.ts (streamGoogleGenAI, shared chunk-to-block translation)
Normalization points:
- iterates
candidate.content.parts - text parts are split into thinking vs text by
isThinkingPart(part) - block transitions close previous block before starting a new one
part.functionCallis treated as a complete tool call (start/delta/end emitted immediately)- finish reason mapped by
mapStopReason()fromgoogle-shared.ts
Tool-call argument streaming:
- function call args arrive as structured object, not incremental JSON text
- implementation emits one synthetic
toolcall_deltacontainingJSON.stringify(arguments) - no partial JSON parser needed for Google in this path
Partial tool-call JSON accumulation and recovery
Shared behavior for Anthropic/OpenAI Responses uses parseStreamingJson() / parseStreamingJsonThrottled() (packages/ai/src/utils/json-parse.ts):
- try
JSON.parse - fallback to the in-house
RelaxedJsonparser (relaxed/repairing) for incomplete fragments - if both fail, return
{}
Implications:
- malformed or truncated argument deltas do not crash stream processing immediately
- in-progress
argumentsmay temporarily be{} - later valid deltas can recover structured arguments because parsing is retried as the buffer grows (throttled to β₯256-byte growth steps mid-stream)
- final
toolcall_endperforms one more parse attempt before emission
Stop reasons vs transport/runtime errors
Provider stop reasons are mapped to normalized stopReason:
- Anthropic:
end_turnβstop,max_tokensβlength,tool_useβtoolUse, safety/refusal casesβerror - OpenAI Responses:
completedβstop,incompleteβlength,failed/cancelledβerror - Google:
STOPβstop,MAX_TOKENSβlength, safety/prohibited/malformed-function-call classesβerror
Error semantics are split in two stages:
- Model completion semantics (provider reported finish reason/status)
- Transport/runtime failure (network/client/parser/abort exceptions)
If provider stream throws or signals failure, each provider wrapper catches and emits terminal error event with:
stopReason = "aborted"when abort signal is set- otherwise
stopReason = "error" errorMessage = finalizeErrorMessage(error, rawRequestDump)(packages/ai/src/utils/http-inspector.ts), which wrapsformatErrorMessageWithRetryAfter()and appends any captured HTTP-error body / raw-request dump (thecursorwrapper callsformatErrorMessageWithRetryAfter()directly)
Malformed chunk / SSE parse failure behavior
The OpenAI Completions/Responses paths use the in-repo HTTP+SSE transport postOpenAIStream() (packages/ai/src/utils/openai-http.ts), which decodes frames with readSseJson() and replaced the openai SDK client. Anthropic uses the in-repo AnthropicMessagesClient (packages/ai/src/providers/anthropic-client.ts); the Google paths and the Codex SSE fallback read SSE via readSseJson() directly, and websocket Codex frames are normalized through the same event handler.
Observed behavior in current implementation:
- malformed SSE framing or chunk JSON surfaces as an exception or stream
errorevent - malformed Codex SSE JSON/framing throws from the local SSE reader
- provider wrapper converts failures into unified terminal
errorevents - no provider-specific resume/retry inside the stream function itself, except Codex websocket-to-SSE transport fallback before replay-unsafe output is emitted
- higher-level retries are handled in
AgentSessionauto-retry logic (message-level retry, not stream-chunk replay)
Cancellation boundaries
Cancellation is layered:
- AI provider request:
options.signalis passed into provider client stream call. - Provider wrapper: after stream loop, aborted signal forces error path (
"Request was aborted"). - Agent loop: checks
signal.abortedbefore handling each provider event and can synthesize an aborted assistant message from the latest partial. - Session/agent controls:
AgentSession.abort()->agent.abort()-> shared abort controller cancellation.
Tool execution cancellation is separate from model stream cancellation:
- tool runners use
AbortSignal.any([agentSignal, steeringAbortSignal]) - steering interrupts can abort remaining tool execution while preserving already-produced tool results
Backpressure boundaries
There is no hard backpressure mechanism between provider SDK stream and downstream consumers:
EventStreamuses in-memory queues with no max size- the throttled partial-JSON re-parse reduces per-delta CPU cost but does not slow provider intake
- if consumers lag significantly, queued events can grow until completion
Current design favors responsiveness and simple ordering over bounded-buffer flow control.
How stream events surface as agent/session events
agentLoop.streamAssistantResponse() bridges AssistantMessageEvent to AgentEvent:
- on
start: pushes placeholder assistant message and emitsmessage_start - on block events (
text_*,thinking_*,toolcall_*): updates last assistant message, emitsmessage_updatewith rawassistantMessageEvent - on terminal (
done/error): resolves final message fromresponse.result(), emitsmessage_end
AgentSession then consumes those events for session-level behaviors:
- TTSR watches
message_update.assistantMessageEventfortext_delta,thinking_delta, andtoolcall_delta - streaming edit guard inspects
toolcall_delta/toolcall_endoneditcalls and can abort early - persistence writes finalized messages at
message_end - auto-retry examines assistant
stopReason === "error"pluserrorMessageheuristics
Unified vs provider-specific responsibilities
Unified (common contract):
- event shape (
AssistantMessageEvent) - final result extraction (
done/error) - immediate in-order event delivery
- agent/session event propagation model
Provider-specific (not fully abstracted):
- upstream event taxonomies and mapping logic
- stop-reason translation tables
- tool-call ID conventions
- reasoning/thinking block semantics and signatures
- usage token semantics and availability timing
- message conversion constraints per API
Implementation files
../../ai/src/stream.tsβ provider dispatch, option mapping, API key/session plumbing, custom API dispatch, and provider-specific credential handling.../../ai/src/utils/event-stream.tsβ generic stream queue + final-result resolution.../../ai/src/utils/json-parse.tsβ partial JSON parsing for streamed tool arguments.../../ai/src/providers/anthropic.tsβ Anthropic event translation and tool JSON delta accumulation.../../ai/src/providers/openai-responses.ts,openai-shared.ts,openai-codex-responses.ts,azure-openai-responses.tsβ Responses-family event translation and status mapping.../../ai/src/providers/google.ts,google-gemini-cli.ts,google-vertex.tsβ Gemini stream chunk-to-block translation variants.../../ai/src/providers/google-shared.tsβ Gemini finish-reason mapping and shared conversion rules.../../ai/src/providers/amazon-bedrock.ts,openai-completions.ts,ollama.ts,cursor.ts,pi-native-client.tsβ additional built-in stream adapters using the same event contract.../../agent/src/agent-loop.tsβ provider stream consumption andmessage_updatebridging.../src/session/agent-session.tsβ session-level handling of streaming updates, abort, retry, and persistence.