Summary
Introduce a shared streaming agent-loop construct in forge-core that emits token-level deltas and tool-call events, and wire it through every surface that currently batches model output:
- the native
forge surface builder REPL (forge-cli/runtime BuilderSession),
forge run (the A2A dev server, ExecuteStream → SSE),
- optionally
forge try (LocalSession).
Today there is no real streaming. LLMExecutor.ExecuteStream runs the full non-streaming loop and then emits the final message as one chunk:
// forge-core/runtime/loop.go:1027
// ExecuteStream runs the tool-calling loop non-streaming, then emits the final
// response as a single message on the channel. True word-by-word streaming is v2.
func (e *LLMExecutor) ExecuteStream(...) (<-chan *a2a.Message, error) {
... resp, err := e.Execute(ctx, task, msg) ... ch <- resp ...
}
So the native builder REPL, forge run's SSE, and A2A consumers all receive a reply only after the whole turn (including every tool round-trip) completes. This is the "v2" that comment anticipates.
We deliberately did not build a one-off streaming path just for the builder surface — a per-surface hack would diverge from the executor's behavior (hooks, budgeting, truncation, egress, MCP, OTel). Instead this issue tracks the common construct all three surfaces share.
Motivation
- Native surface (
forge bare → forge native) shows the model reply all at once; users want to watch it stream.
forge run advertises an "A2A-compliant dev server" but streaming responses aren't truly incremental.
- One implementation avoids three drifting copies of the agent loop.
What already exists (building blocks)
llm.Client.ChatStream(ctx, *ChatRequest) (<-chan StreamDelta, error) — implemented for Anthropic, OpenAI, Bedrock, Ollama.
llm.StreamDelta{ Content, ToolCalls, FinishReason, Done, Usage } (forge-core/llm/types.go) already carries text deltas and tool calls + finish reason + usage.
- Provider streaming tool-call emission differs and must be assembled:
- Anthropic (
providers/anthropic.go): input_json_delta fragments per tool_use content block (accumulate arguments per block).
- OpenAI (
providers/openai.go): delta.tool_calls streamed by index (accumulate function.arguments per index).
- The non-streaming reference loop to reach parity with:
LLMExecutor.Execute (forge-core/runtime/loop.go:219) — request build (:461), client.Chat (:499), tool exec (e.tools.Execute, :858), hooks, context budgeting, tool-result truncation, MCP handling, OTel spans.
Proposed design
Add a streaming path to LLMExecutor (not a parallel loop) that reuses the same request construction, tool registry, hooks, budgeting, and truncation, but drives each assistant turn via ChatStream:
- Per iteration: build the same
ChatRequest, call ChatStream.
- Forward
StreamDelta.Content to the consumer as it arrives; accumulate assistant text.
- Accumulate streamed
ToolCalls (provider-aware assembly).
- On finish: if tool calls → run
AfterToolExec/BeforeToolExec hooks + e.tools.Execute, append assistant + tool-result messages, loop; else the accumulated text is the final message.
- Preserve:
OnError/audit hooks, DeferToolResultTruncation, compaction, cancellation, prompt-cache stability, and OTel spans — parity with Execute.
Emit a typed event stream (text delta / tool-call started / tool-result / iteration / done) rather than only *a2a.Message, so:
BuilderSession.RunTurn can print deltas live (and the existing tryview renderer can show tool lines inline);
forge run maps events onto A2A SSE frames;
- A2A streaming stays spec-compliant.
Keep the current Execute as the non-streaming entry point; ExecuteStream becomes a thin adapter over the new construct (delta-batching for callers that want whole messages).
Consumers to wire
Acceptance criteria
- Native
forge (forge native) streams the reply token-by-token for anthropic and openai.
forge run streams incrementally over A2A/SSE.
- Streaming path passes the same tool-loop / hook / truncation tests as
Execute (parity), plus provider tool-call-assembly tests (Anthropic fragment vs. OpenAI index).
- Cancellation (Ctrl-C / context cancel) stops an in-flight stream cleanly.
Notes / risks
- Provider tool-call assembly is the trickiest part (fragmented args); needs dedicated tests.
- Must not regress prompt-cache prefix stability or tool-result truncation timing (
DeferToolResultTruncation).
- OTel span shape should match the non-streaming path so audit/trace consumers are unaffected.
Summary
Introduce a shared streaming agent-loop construct in
forge-corethat emits token-level deltas and tool-call events, and wire it through every surface that currently batches model output:forgesurface builder REPL (forge-cli/runtimeBuilderSession),forge run(the A2A dev server,ExecuteStream→ SSE),forge try(LocalSession).Today there is no real streaming.
LLMExecutor.ExecuteStreamruns the full non-streaming loop and then emits the final message as one chunk:So the native builder REPL,
forge run's SSE, and A2A consumers all receive a reply only after the whole turn (including every tool round-trip) completes. This is the "v2" that comment anticipates.We deliberately did not build a one-off streaming path just for the builder surface — a per-surface hack would diverge from the executor's behavior (hooks, budgeting, truncation, egress, MCP, OTel). Instead this issue tracks the common construct all three surfaces share.
Motivation
forgebare → forge native) shows the model reply all at once; users want to watch it stream.forge runadvertises an "A2A-compliant dev server" but streaming responses aren't truly incremental.What already exists (building blocks)
llm.Client.ChatStream(ctx, *ChatRequest) (<-chan StreamDelta, error)— implemented for Anthropic, OpenAI, Bedrock, Ollama.llm.StreamDelta{ Content, ToolCalls, FinishReason, Done, Usage }(forge-core/llm/types.go) already carries text deltas and tool calls + finish reason + usage.providers/anthropic.go):input_json_deltafragments pertool_usecontent block (accumulate arguments per block).providers/openai.go):delta.tool_callsstreamed by index (accumulatefunction.argumentsper index).LLMExecutor.Execute(forge-core/runtime/loop.go:219) — request build (:461),client.Chat(:499), tool exec (e.tools.Execute,:858), hooks, context budgeting, tool-result truncation, MCP handling, OTel spans.Proposed design
Add a streaming path to
LLMExecutor(not a parallel loop) that reuses the same request construction, tool registry, hooks, budgeting, and truncation, but drives each assistant turn viaChatStream:ChatRequest, callChatStream.StreamDelta.Contentto the consumer as it arrives; accumulate assistant text.ToolCalls(provider-aware assembly).AfterToolExec/BeforeToolExechooks +e.tools.Execute, append assistant + tool-result messages, loop; else the accumulated text is the final message.OnError/audit hooks,DeferToolResultTruncation, compaction, cancellation, prompt-cache stability, and OTel spans — parity withExecute.Emit a typed event stream (text delta / tool-call started / tool-result / iteration / done) rather than only
*a2a.Message, so:BuilderSession.RunTurncan print deltas live (and the existingtryviewrenderer can show tool lines inline);forge runmaps events onto A2A SSE frames;Keep the current
Executeas the non-streaming entry point;ExecuteStreambecomes a thin adapter over the new construct (delta-batching for callers that want whole messages).Consumers to wire
BuilderSession(native surface) — stream deltas to the REPL for anthropic + openai (the concrete ask that prompted this).forge runA2A server — real incremental SSE via the shared events.forge tryLocalSession— optional, once the construct lands.Acceptance criteria
forge(forge native) streams the reply token-by-token for anthropic and openai.forge runstreams incrementally over A2A/SSE.Execute(parity), plus provider tool-call-assembly tests (Anthropic fragment vs. OpenAI index).Notes / risks
DeferToolResultTruncation).