Skip to content

[Feature]: Shared streaming agent-loop for builder surface, forge run, and A2A #476

Description

@initializ-mk

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:

  1. Per iteration: build the same ChatRequest, call ChatStream.
  2. Forward StreamDelta.Content to the consumer as it arrives; accumulate assistant text.
  3. Accumulate streamed ToolCalls (provider-aware assembly).
  4. 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.
  5. 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

  • BuilderSession (native surface) — stream deltas to the REPL for anthropic + openai (the concrete ask that prompted this).
  • forge run A2A server — real incremental SSE via the shared events.
  • forge try LocalSession — optional, once the construct lands.

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.

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    enhancementNew feature or request

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions