feat: stream advisor tool output (#25032)

Stream advisor output into the advisor tool card while the nested
advisor call is still running.

This keeps the advisor implementation intentionally advisor-specific:
the parent model still receives the same final structured tool result,
while the frontend receives transient `tool-result.result_delta` parts
to render partial advisor text in the expanded card. The final persisted
chat history remains unchanged.

Refs CODAGT-322.

Generated by Coder Agents.

<details>
<summary>Implementation plan</summary>

- Publish advisor text deltas from the nested `chatloop.Run` via
`RunAdvisorOptions.OnAdviceDelta`.
- Forward those deltas through `chatadvisor.Tool` with the parent
advisor tool call ID.
- Emit transient `ChatMessagePartTypeToolResult` websocket parts with
`ResultDelta` from `chatd`.
- Add `result_delta` to the generated tool-result TypeScript variant.
- Accumulate tool result deltas in frontend stream state and keep the
tool running until the final result arrives.
- Render streamed advisor advice in the existing advisor card using
streaming markdown mode, while retaining the updated advisor UI.

</details>
This commit is contained in:
Thomas Kosiewski
2026-05-11 20:18:49 +02:00
committed by GitHub
parent 6bb88775ab
commit e56381eb61
18 changed files with 619 additions and 23 deletions
+3
View File
@@ -15970,6 +15970,9 @@ const docTemplate = `{
"result_delta": {
"type": "string"
},
"result_reset": {
"type": "boolean"
},
"signature": {
"type": "string"
},
+3
View File
@@ -14391,6 +14391,9 @@
"result_delta": {
"type": "string"
},
"result_reset": {
"type": "boolean"
},
"signature": {
"type": "string"
},
+28 -2
View File
@@ -3,18 +3,29 @@ package chatadvisor
import (
"context"
"strings"
"time"
"charm.land/fantasy"
"golang.org/x/xerrors"
"github.com/coder/coder/v2/coderd/x/chatd/chatloop"
"github.com/coder/coder/v2/coderd/x/chatd/chatretry"
"github.com/coder/coder/v2/codersdk"
)
// RunAdvisorOptions carries optional streaming callbacks for a
// single RunAdvisor invocation.
type RunAdvisorOptions struct {
OnAdviceDelta func(delta string)
OnAdviceReset func()
}
// RunAdvisor executes a single, tool-less nested advisor call.
func (rt *Runtime) RunAdvisor(
ctx context.Context,
question string,
conversationSnapshot []fantasy.Message,
opts *RunAdvisorOptions,
) (AdvisorResult, error) {
// Model, MaxUsesPerRun, and MaxOutputTokens are validated by NewRuntime.
// Runtime fields are unexported so callers cannot bypass that.
@@ -37,7 +48,7 @@ func (rt *Runtime) RunAdvisor(
resetProviderOptionsForNestedCall(nestedProviderOptions)
var persistedStep chatloop.PersistedStep
runOpts := chatloop.RunOptions{
chatLoopOpts := chatloop.RunOptions{
Model: rt.cfg.Model,
Messages: BuildAdvisorMessages(question, conversationSnapshot),
MaxSteps: 1,
@@ -48,8 +59,23 @@ func (rt *Runtime) RunAdvisor(
return nil
},
}
if opts != nil && opts.OnAdviceDelta != nil {
chatLoopOpts.PublishMessagePart = func(role codersdk.ChatMessageRole, part codersdk.ChatMessagePart) {
if role != codersdk.ChatMessageRoleAssistant ||
part.Type != codersdk.ChatMessagePartTypeText ||
part.Text == "" {
return
}
opts.OnAdviceDelta(part.Text)
}
}
if opts != nil && opts.OnAdviceReset != nil {
chatLoopOpts.OnRetry = func(int, error, chatretry.ClassifiedError, time.Duration) {
opts.OnAdviceReset()
}
}
if err := chatloop.Run(ctx, runOpts); err != nil {
if err := chatloop.Run(ctx, chatLoopOpts); err != nil {
// Refund the use so a transient provider failure does not
// permanently exhaust the per-run advisor budget.
rt.release()
+124 -8
View File
@@ -48,7 +48,7 @@ func TestAdvisorRunAdvice(t *testing.T) {
result, err := runtime.RunAdvisor(t.Context(), question, []fantasy.Message{
textMessage(fantasy.MessageRoleSystem, "existing system"),
textMessage(fantasy.MessageRoleUser, "hello"),
})
}, nil)
require.NoError(t, err)
require.Equal(t, chatadvisor.ResultTypeAdvice, result.Type)
require.Equal(t, "Take the smallest safe change.", result.Advice)
@@ -63,6 +63,122 @@ func TestAdvisorRunAdvice(t *testing.T) {
require.Equal(t, question, singleText(t, capturedCall.Prompt[len(capturedCall.Prompt)-1]))
}
func TestAdvisorRunStreamsAdviceDeltas(t *testing.T) {
t.Parallel()
var deltas []string
runtime, err := chatadvisor.NewRuntime(chatadvisor.RuntimeConfig{
Model: &chattest.FakeModel{
ProviderName: "test-provider",
ModelName: "test-model",
StreamFn: func(_ context.Context, _ fantasy.Call) (fantasy.StreamResponse, error) {
return streamFromParts([]fantasy.StreamPart{
{Type: fantasy.StreamPartTypeTextStart, ID: "text-1"},
{Type: fantasy.StreamPartTypeTextDelta, ID: "text-1", Delta: "Use "},
{Type: fantasy.StreamPartTypeTextDelta, ID: "text-1", Delta: "the smaller "},
{Type: fantasy.StreamPartTypeTextDelta, ID: "text-1", Delta: "diff."},
{Type: fantasy.StreamPartTypeTextEnd, ID: "text-1"},
{Type: fantasy.StreamPartTypeFinish, FinishReason: fantasy.FinishReasonStop},
}), nil
},
},
MaxUsesPerRun: 2,
MaxOutputTokens: 128,
})
require.NoError(t, err)
result, err := runtime.RunAdvisor(t.Context(), "what should I do?", nil, &chatadvisor.RunAdvisorOptions{
OnAdviceDelta: func(delta string) {
deltas = append(deltas, delta)
},
})
require.NoError(t, err)
require.Equal(t, []string{"Use ", "the smaller ", "diff."}, deltas)
require.Equal(t, chatadvisor.ResultTypeAdvice, result.Type)
require.Equal(t, "Use the smaller diff.", result.Advice)
require.Equal(t, 1, result.RemainingUses)
}
func TestAdvisorRunResetsAdviceDeltasOnRetry(t *testing.T) {
t.Parallel()
var (
calls int
events []string
)
runtime, err := chatadvisor.NewRuntime(chatadvisor.RuntimeConfig{
Model: &chattest.FakeModel{
ProviderName: "test-provider",
ModelName: "test-model",
StreamFn: func(_ context.Context, _ fantasy.Call) (fantasy.StreamResponse, error) {
calls++
if calls == 1 {
return streamFromParts([]fantasy.StreamPart{
{Type: fantasy.StreamPartTypeTextStart, ID: "text-1"},
{Type: fantasy.StreamPartTypeTextDelta, ID: "text-1", Delta: "stale "},
{Type: fantasy.StreamPartTypeError, Error: xerrors.New("received status 429 from upstream")},
}), nil
}
return streamFromParts([]fantasy.StreamPart{
{Type: fantasy.StreamPartTypeTextStart, ID: "text-1"},
{Type: fantasy.StreamPartTypeTextDelta, ID: "text-1", Delta: "fresh advice"},
{Type: fantasy.StreamPartTypeTextEnd, ID: "text-1"},
{Type: fantasy.StreamPartTypeFinish, FinishReason: fantasy.FinishReasonStop},
}), nil
},
},
MaxUsesPerRun: 2,
MaxOutputTokens: 128,
})
require.NoError(t, err)
result, err := runtime.RunAdvisor(t.Context(), "what should I do?", nil, &chatadvisor.RunAdvisorOptions{
OnAdviceDelta: func(delta string) {
events = append(events, "delta:"+delta)
},
OnAdviceReset: func() {
events = append(events, "reset")
},
})
require.NoError(t, err)
require.Equal(t, []string{"delta:stale ", "reset", "delta:fresh advice"}, events)
require.Equal(t, chatadvisor.ResultTypeAdvice, result.Type)
require.Equal(t, "fresh advice", result.Advice)
}
func TestAdvisorRunErrorAfterPartialDelta(t *testing.T) {
t.Parallel()
var deltas []string
runtime, err := chatadvisor.NewRuntime(chatadvisor.RuntimeConfig{
Model: &chattest.FakeModel{
ProviderName: "test-provider",
ModelName: "test-model",
StreamFn: func(_ context.Context, _ fantasy.Call) (fantasy.StreamResponse, error) {
return streamFromParts([]fantasy.StreamPart{
{Type: fantasy.StreamPartTypeTextStart, ID: "text-1"},
{Type: fantasy.StreamPartTypeTextDelta, ID: "text-1", Delta: "partial advice"},
{Type: fantasy.StreamPartTypeError, Error: xerrors.New("boom after partial")},
}), nil
},
},
MaxUsesPerRun: 1,
MaxOutputTokens: 128,
})
require.NoError(t, err)
result, err := runtime.RunAdvisor(t.Context(), "what should I do?", nil, &chatadvisor.RunAdvisorOptions{
OnAdviceDelta: func(delta string) {
deltas = append(deltas, delta)
},
})
require.NoError(t, err)
require.Equal(t, []string{"partial advice"}, deltas)
require.Equal(t, chatadvisor.ResultTypeError, result.Type)
require.Contains(t, result.Error, "boom after partial")
require.Equal(t, 1, result.RemainingUses)
}
func TestAdvisorRunLimitReached(t *testing.T) {
t.Parallel()
@@ -86,12 +202,12 @@ func TestAdvisorRunLimitReached(t *testing.T) {
})
require.NoError(t, err)
first, err := runtime.RunAdvisor(t.Context(), "first?", nil)
first, err := runtime.RunAdvisor(t.Context(), "first?", nil, nil)
require.NoError(t, err)
require.Equal(t, chatadvisor.ResultTypeAdvice, first.Type)
require.Equal(t, 0, first.RemainingUses)
second, err := runtime.RunAdvisor(t.Context(), "second?", nil)
second, err := runtime.RunAdvisor(t.Context(), "second?", nil, nil)
require.NoError(t, err)
require.Equal(t, chatadvisor.ResultTypeLimitReached, second.Type)
require.Equal(t, 0, second.RemainingUses)
@@ -114,7 +230,7 @@ func TestAdvisorRunError(t *testing.T) {
})
require.NoError(t, err)
result, err := runtime.RunAdvisor(t.Context(), "what failed?", nil)
result, err := runtime.RunAdvisor(t.Context(), "what failed?", nil, nil)
require.NoError(t, err)
require.Equal(t, chatadvisor.ResultTypeError, result.Type)
require.Contains(t, result.Error, "boom")
@@ -149,12 +265,12 @@ func TestAdvisorRunError(t *testing.T) {
})
require.NoError(t, err)
failed, err := runtime2.RunAdvisor(t.Context(), "first?", nil)
failed, err := runtime2.RunAdvisor(t.Context(), "first?", nil, nil)
require.NoError(t, err)
require.Equal(t, chatadvisor.ResultTypeError, failed.Type)
require.Equal(t, 1, failed.RemainingUses)
retried, err := runtime2.RunAdvisor(t.Context(), "retry?", nil)
retried, err := runtime2.RunAdvisor(t.Context(), "retry?", nil, nil)
require.NoError(t, err)
require.Equal(t, chatadvisor.ResultTypeAdvice, retried.Type)
require.Equal(t, "recovered", retried.Advice)
@@ -251,7 +367,7 @@ func TestNewRuntimeDeepClonesOpenAIResponsesProviderOptions(t *testing.T) {
})
require.NoError(t, err)
result, err := runtime.RunAdvisor(t.Context(), "anything?", nil)
result, err := runtime.RunAdvisor(t.Context(), "anything?", nil, nil)
require.NoError(t, err)
require.Equal(t, chatadvisor.ResultTypeAdvice, result.Type)
@@ -316,7 +432,7 @@ func TestAdvisorRunStripsChainStateAndIsConsistentAcrossCalls(t *testing.T) {
require.NoError(t, err)
for i := range 2 {
result, err := runtime.RunAdvisor(t.Context(), fmt.Sprintf("q%d", i), nil)
result, err := runtime.RunAdvisor(t.Context(), fmt.Sprintf("q%d", i), nil, nil)
require.NoError(t, err)
require.Equal(t, chatadvisor.ResultTypeAdvice, result.Type)
}
+19 -2
View File
@@ -24,6 +24,8 @@ const advisorQuestionMaxRunes = 2000
type ToolOptions struct {
Runtime *Runtime
GetConversationSnapshot func() []fantasy.Message
PublishAdviceDelta func(toolCallID string, delta string)
PublishAdviceReset func(toolCallID string)
}
// Tool returns a fantasy.AgentTool that asks a nested model for concise
@@ -33,7 +35,7 @@ func Tool(opts ToolOptions) fantasy.AgentTool {
return fantasy.NewAgentTool(
ToolName,
"Ask a separate advisor pass for strategic guidance about planning, architecture, tradeoffs, or debugging strategy. Provide a brief question. The advisor sees recent conversation context, runs without tools for a single step, and responds to the parent agent rather than the end user.",
func(ctx context.Context, args AdvisorArgs, _ fantasy.ToolCall) (fantasy.ToolResponse, error) {
func(ctx context.Context, args AdvisorArgs, call fantasy.ToolCall) (fantasy.ToolResponse, error) {
if opts.Runtime == nil {
return fantasy.NewTextErrorResponse("advisor runtime is not configured"), nil
}
@@ -51,7 +53,22 @@ func Tool(opts ToolOptions) fantasy.AgentTool {
), nil
}
result, err := opts.Runtime.RunAdvisor(ctx, question, opts.GetConversationSnapshot())
var runOpts *RunAdvisorOptions
if call.ID != "" && (opts.PublishAdviceDelta != nil || opts.PublishAdviceReset != nil) {
runOpts = &RunAdvisorOptions{}
if opts.PublishAdviceDelta != nil {
runOpts.OnAdviceDelta = func(delta string) {
opts.PublishAdviceDelta(call.ID, delta)
}
}
if opts.PublishAdviceReset != nil {
runOpts.OnAdviceReset = func() {
opts.PublishAdviceReset(call.ID)
}
}
}
result, err := opts.Runtime.RunAdvisor(ctx, question, opts.GetConversationSnapshot(), runOpts)
if err != nil {
return fantasy.NewTextErrorResponse(err.Error()), nil
}
+120
View File
@@ -58,6 +58,126 @@ func TestAdvisorToolSuccess(t *testing.T) {
require.Equal(t, 1, result.RemainingUses)
}
func TestAdvisorToolPublishesAdviceDeltasWithToolCallID(t *testing.T) {
t.Parallel()
type publishedDelta struct {
toolCallID string
delta string
}
var published []publishedDelta
runtime, err := chatadvisor.NewRuntime(chatadvisor.RuntimeConfig{
Model: &chattest.FakeModel{
ProviderName: "test-provider",
ModelName: "test-model",
StreamFn: func(_ context.Context, _ fantasy.Call) (fantasy.StreamResponse, error) {
return streamFromParts([]fantasy.StreamPart{
{Type: fantasy.StreamPartTypeTextStart, ID: "text-1"},
{Type: fantasy.StreamPartTypeTextDelta, ID: "text-1", Delta: "Prefer "},
{Type: fantasy.StreamPartTypeTextDelta, ID: "text-1", Delta: "the small diff."},
{Type: fantasy.StreamPartTypeTextEnd, ID: "text-1"},
{Type: fantasy.StreamPartTypeFinish, FinishReason: fantasy.FinishReasonStop},
}), nil
},
},
MaxUsesPerRun: 2,
MaxOutputTokens: 128,
})
require.NoError(t, err)
tool := chatadvisor.Tool(chatadvisor.ToolOptions{
Runtime: runtime,
GetConversationSnapshot: func() []fantasy.Message { return nil },
PublishAdviceDelta: func(toolCallID string, delta string) {
published = append(published, publishedDelta{toolCallID: toolCallID, delta: delta})
},
})
resp := runAdvisorTool(t, tool, chatadvisor.AdvisorArgs{Question: "What's safest?"})
require.False(t, resp.IsError)
require.Equal(t, []publishedDelta{
{toolCallID: "call-1", delta: "Prefer "},
{toolCallID: "call-1", delta: "the small diff."},
}, published)
var result chatadvisor.AdvisorResult
require.NoError(t, json.Unmarshal([]byte(resp.Content), &result))
require.Equal(t, chatadvisor.ResultTypeAdvice, result.Type)
require.Equal(t, "Prefer the small diff.", result.Advice)
}
func TestAdvisorToolPublishesAdviceResetWithToolCallID(t *testing.T) {
t.Parallel()
type publishedEvent struct {
kind string
toolCallID string
delta string
}
var (
calls int
published []publishedEvent
)
runtime, err := chatadvisor.NewRuntime(chatadvisor.RuntimeConfig{
Model: &chattest.FakeModel{
ProviderName: "test-provider",
ModelName: "test-model",
StreamFn: func(_ context.Context, _ fantasy.Call) (fantasy.StreamResponse, error) {
calls++
if calls == 1 {
return streamFromParts([]fantasy.StreamPart{
{Type: fantasy.StreamPartTypeTextStart, ID: "text-1"},
{Type: fantasy.StreamPartTypeTextDelta, ID: "text-1", Delta: "stale "},
{Type: fantasy.StreamPartTypeError, Error: xerrors.New("received status 429 from upstream")},
}), nil
}
return streamFromParts([]fantasy.StreamPart{
{Type: fantasy.StreamPartTypeTextStart, ID: "text-1"},
{Type: fantasy.StreamPartTypeTextDelta, ID: "text-1", Delta: "fresh advice"},
{Type: fantasy.StreamPartTypeTextEnd, ID: "text-1"},
{Type: fantasy.StreamPartTypeFinish, FinishReason: fantasy.FinishReasonStop},
}), nil
},
},
MaxUsesPerRun: 2,
MaxOutputTokens: 128,
})
require.NoError(t, err)
tool := chatadvisor.Tool(chatadvisor.ToolOptions{
Runtime: runtime,
GetConversationSnapshot: func() []fantasy.Message { return nil },
PublishAdviceDelta: func(toolCallID string, delta string) {
published = append(published, publishedEvent{
kind: "delta",
toolCallID: toolCallID,
delta: delta,
})
},
PublishAdviceReset: func(toolCallID string) {
published = append(published, publishedEvent{
kind: "reset",
toolCallID: toolCallID,
})
},
})
resp := runAdvisorTool(t, tool, chatadvisor.AdvisorArgs{Question: "What's safest?"})
require.False(t, resp.IsError)
require.Equal(t, []publishedEvent{
{kind: "delta", toolCallID: "call-1", delta: "stale "},
{kind: "reset", toolCallID: "call-1"},
{kind: "delta", toolCallID: "call-1", delta: "fresh advice"},
}, published)
var result chatadvisor.AdvisorResult
require.NoError(t, json.Unmarshal([]byte(resp.Content), &result))
require.Equal(t, chatadvisor.ResultTypeAdvice, result.Type)
require.Equal(t, "fresh advice", result.Advice)
}
func TestAdvisorToolRejectsEmptyQuestion(t *testing.T) {
t.Parallel()
+22
View File
@@ -7250,6 +7250,28 @@ func (p *Server) runChat(
// no tools. Strip it before handing the snapshot over.
return stripAdvisorGuidanceBlock(slices.Clone(advisorPromptSnapshot))
},
PublishAdviceDelta: func(toolCallID string, delta string) {
if toolCallID == "" || delta == "" {
return
}
p.publishMessagePart(chat.ID, codersdk.ChatMessageRoleTool, codersdk.ChatMessagePart{
Type: codersdk.ChatMessagePartTypeToolResult,
ToolCallID: toolCallID,
ToolName: chatadvisor.ToolName,
ResultDelta: delta,
})
},
PublishAdviceReset: func(toolCallID string) {
if toolCallID == "" {
return
}
p.publishMessagePart(chat.ID, codersdk.ChatMessageRoleTool, codersdk.ChatMessagePart{
Type: codersdk.ChatMessagePartTypeToolResult,
ToolCallID: toolCallID,
ToolName: chatadvisor.ToolName,
ResultReset: true,
})
},
}))
}
+32 -1
View File
@@ -9494,6 +9494,7 @@ func TestAdvisorHappyPath_RootChat(t *testing.T) {
ctx := testutil.Context(t, testutil.WaitLong)
const advisorReply = "break the problem into smaller pieces first"
advisorDeltas := []string{"break the problem ", "into smaller pieces first"}
var (
streamedCallCount atomic.Int32
@@ -9526,7 +9527,7 @@ func TestAdvisorHappyPath_RootChat(t *testing.T) {
streamedCallsMu.Unlock()
advisorCallSeen.Store(true)
return chattest.OpenAIStreamingResponse(
chattest.OpenAITextChunks(advisorReply)...,
chattest.OpenAITextChunks(advisorDeltas...)...,
)
default:
// Parent turn 2: observe the advisor tool result and close
@@ -9612,6 +9613,36 @@ func TestAdvisorHappyPath_RootChat(t *testing.T) {
}
require.True(t, parentSawAdvisorResult,
"parent must see the advisor reply in its continuation call")
snapshot, _, cancelStream, ok := server.Subscribe(ctx, chat.ID, nil, 0)
require.True(t, ok)
cancelStream()
var streamedAdvisorDeltas []string
for _, event := range snapshot {
if event.Type != codersdk.ChatStreamEventTypeMessagePart || event.MessagePart == nil {
continue
}
part := event.MessagePart.Part
if event.MessagePart.Role == codersdk.ChatMessageRoleTool &&
part.Type == codersdk.ChatMessagePartTypeToolResult &&
part.ToolName == chatadvisor.ToolName &&
part.ResultDelta != "" {
streamedAdvisorDeltas = append(streamedAdvisorDeltas, part.ResultDelta)
}
}
require.Equal(t, advisorDeltas, streamedAdvisorDeltas,
"advisor nested text deltas must stream into the parent tool card")
persisted, err := db.GetChatMessagesByChatID(ctx, database.GetChatMessagesByChatIDParams{
ChatID: chat.ID,
AfterID: 0,
})
require.NoError(t, err)
for _, msg := range persisted {
require.NotContains(t, string(msg.Content.RawMessage), "result_delta",
"advisor deltas are stream-only and must not be persisted")
}
}
// TestAdvisorGating_ChildChat guards the second dimension of the advisor