From e0ae22e83baf823d158c9913ff8b09b164f2e8af Mon Sep 17 00:00:00 2001 From: euvre <93761161+euvre@users.noreply.github.com> Date: Wed, 22 Jul 2026 14:59:29 +0800 Subject: [PATCH] fix(agent): prevent double reply when Agent node is followed by Message node (#17233) --- internal/agent/component/message.go | 8 +++++- internal/agent/component/message_test.go | 34 ++++++++++++++++++++++++ 2 files changed, 41 insertions(+), 1 deletion(-) diff --git a/internal/agent/component/message.go b/internal/agent/component/message.go index 1b02c3688d..06a057f04a 100644 --- a/internal/agent/component/message.go +++ b/internal/agent/component/message.go @@ -216,7 +216,13 @@ func (m *MessageComponent) Invoke(ctx context.Context, inputs map[string]any) (m Text: resolved, }) } - if rendered != "" { + // Skip emission when an upstream Agent component already streamed its + // answer via the agent message emitter. Without this guard the user sees + // the same content twice: once from the Agent's StreamCallback during the + // ReAct loop, and again from the Message node's own emit call here. + // The rendered content is still stored in outputs["content"] for state + // persistence and the service-layer answer collector. + if rendered != "" && !runtime.AgentMessageEventsEmitted(ctx) { if !runtime.EmitCanvasMessage(ctx, rendered) { runtime.EmitAgentMessage(ctx, rendered, "") } diff --git a/internal/agent/component/message_test.go b/internal/agent/component/message_test.go index 6ff4ac7901..5d686cd57c 100644 --- a/internal/agent/component/message_test.go +++ b/internal/agent/component/message_test.go @@ -193,6 +193,40 @@ func TestMessage_EmitsDirectCanvasMessage(t *testing.T) { } } +// TestMessage_SkipsEmissionWhenAgentAlreadyStreamed verifies that the +// Message component does not re-emit content when an upstream Agent +// component already streamed its answer. This prevents the double-reply +// bug (Agent → Message): the Agent streams via StreamCallback during the +// ReAct loop, and the Message node must not send the same content again. +func TestMessage_SkipsEmissionWhenAgentAlreadyStreamed(t *testing.T) { + c, _ := NewMessageComponent(nil) + state := canvas.NewCanvasState("run-skip", "task-skip") + ctx := withStateForTest(context.Background(), state) + var emitted []string + ctx = runtime.WithAgentMessageEmitter(ctx, func(contentDelta, thinkingDelta string) { + if contentDelta != "" { + emitted = append(emitted, contentDelta) + } + }) + + // Simulate the Agent having already streamed its answer. + runtime.EmitAgentMessage(ctx, "agent streamed this", "") + emitted = nil // clear setup emission so we only capture Message's output + + out, err := c.Invoke(ctx, map[string]any{"content": "agent streamed this"}) + if err != nil { + t.Fatalf("Invoke: %v", err) + } + // Content is still set in outputs for state persistence. + if got, _ := out["content"].(string); got != "agent streamed this" { + t.Fatalf("content: got %q, want %q", got, "agent streamed this") + } + // But no additional SSE emission should occur. + if len(emitted) != 0 { + t.Fatalf("emitted = %#v, want empty (Agent already streamed)", emitted) + } +} + func TestMessage_FormalizedContentFallback(t *testing.T) { c, _ := NewMessageComponent(nil) state := canvas.NewCanvasState("run-5", "task-5")