From 7acb8e59fa47fbc06ce0aa2a42bde38fe8c2aecf Mon Sep 17 00:00:00 2001 From: toki Date: Thu, 2 Jul 2026 14:05:19 +0900 Subject: [PATCH] =?UTF-8?q?feat:=20reasoning=5Fdelta=20=EC=9D=B4=EB=B2=A4?= =?UTF-8?q?=ED=8A=B8=20=EC=B2=98=EB=A6=AC=20=EB=B0=8F=20=EC=B6=9C=EB=A0=A5?= =?UTF-8?q?=20=EA=B0=9C=EC=84=A0?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - reasoning_delta 이벤트 전파: edge console와 node sink에서 explicit 케이스로 처리 - 각 줄에 [node-reasoning] prefix를 붙여 reasoning 출력을 구분 - 빈 delta는 무시하고, 멀티라인 content도 각 줄에 prefix 적용 - 관련单元测试 추가 --- apps/edge/internal/opsconsole/console.go | 2 + apps/edge/internal/opsconsole/console_test.go | 22 ++++++ apps/edge/internal/opsconsole/events.go | 15 ++++ apps/edge/internal/opsconsole/events_test.go | 31 ++++++++ apps/node/internal/node/node.go | 20 +++++ apps/node/internal/node/sink_test.go | 76 +++++++++++++++++++ 6 files changed, 166 insertions(+) diff --git a/apps/edge/internal/opsconsole/console.go b/apps/edge/internal/opsconsole/console.go index 7997eed..3f435c5 100644 --- a/apps/edge/internal/opsconsole/console.go +++ b/apps/edge/internal/opsconsole/console.go @@ -282,6 +282,8 @@ func SendRun(ctx context.Context, edgeSvc *edgeservice.Service, events *EventRou switch event.GetType() { case "start": fmt.Fprintf(out, "%s start run_id=%s\n", nodeEventPrefix(label), dispatch.RunID) + case "reasoning_delta": + writeReasoningDelta(out, label, event.GetDelta()) case "delta": response.Write(event.GetDelta()) case "complete": diff --git a/apps/edge/internal/opsconsole/console_test.go b/apps/edge/internal/opsconsole/console_test.go index 1d19ec1..669b35c 100644 --- a/apps/edge/internal/opsconsole/console_test.go +++ b/apps/edge/internal/opsconsole/console_test.go @@ -233,6 +233,28 @@ func TestSendRun_SubmitRunRequest_MetadataSource(t *testing.T) { } } +func TestWriteReasoningDeltaPrefixesEachLine(t *testing.T) { + var out bytes.Buffer + writeReasoningDelta(&out, "node0", "사용자가 인사했습니다.\n\x1f\x1f\n") + got := out.String() + want1 := "[node0-reasoning] 사용자가 인사했습니다." + want2 := "[node0-reasoning] \x1f\x1f" + if !strings.Contains(got, want1) { + t.Errorf("expected %q, got:\n%s", want1, got) + } + if !strings.Contains(got, want2) { + t.Errorf("expected %q, got:\n%s", want2, got) + } +} + +func TestWriteReasoningDeltaIgnoresEmptyDelta(t *testing.T) { + var out bytes.Buffer + writeReasoningDelta(&out, "node0", "") + if out.Len() != 0 { + t.Errorf("expected no output for empty delta, got:\n%s", out.String()) + } +} + func TestFormatNodeCommandView_RendersSortedResult(t *testing.T) { view := edgeservice.NodeCommandView{ NodeLabel: "node0", diff --git a/apps/edge/internal/opsconsole/events.go b/apps/edge/internal/opsconsole/events.go index 18dfe9a..51c8f9f 100644 --- a/apps/edge/internal/opsconsole/events.go +++ b/apps/edge/internal/opsconsole/events.go @@ -202,6 +202,8 @@ func (r *EventRouter) printAsyncLocked(event *iop.RunEvent) { switch event.GetType() { case "start": fmt.Fprintf(r.out, "%s start run_id=%s session=%s background=%v\n", nodeEventPrefix(label), runID, event.GetSessionId(), event.GetBackground()) + case "reasoning_delta": + writeReasoningDelta(r.out, label, event.GetDelta()) case "delta": r.responseStream(runID, label).Write(event.GetDelta()) case "complete": @@ -301,3 +303,16 @@ func (s *ResponseStream) FinishIfStarted() { s.Finish() } } + +func nodeReasoningPrefix(label string) string { + return fmt.Sprintf("[%s-reasoning] ", label) +} + +func writeReasoningDelta(out io.Writer, label, delta string) { + if delta == "" { + return + } + stream := NewResponseStream(out, nodeReasoningPrefix(label)) + stream.Write(delta) + stream.FinishIfStarted() +} diff --git a/apps/edge/internal/opsconsole/events_test.go b/apps/edge/internal/opsconsole/events_test.go index df48f68..8e6f5be 100644 --- a/apps/edge/internal/opsconsole/events_test.go +++ b/apps/edge/internal/opsconsole/events_test.go @@ -148,6 +148,37 @@ func TestEventRouterFallsBackToNodeID(t *testing.T) { } } +func TestEventRouterPrintsReasoningDeltaForAsyncRun(t *testing.T) { + var out bytes.Buffer + reg := edgenode.NewRegistry() + reg.Register(&edgenode.NodeEntry{NodeID: "node-1", Alias: "node1"}) + router := NewEventRouter(&out, reg, nil) + + router.Handle(&iop.RunEvent{ + RunId: "run-rd", + Type: "reasoning_delta", + Delta: "사용자가 인사했습니다.\n\x1f\x1f\n", + NodeId: "node-1", + Background: true, + }) + + got := out.String() + want1 := "[node0-reasoning] 사용자가 인사했습니다." + want2 := "[node0-reasoning] \x1f\x1f" + if !strings.Contains(got, want1) { + t.Errorf("expected %q in output, got:\n%s", want1, got) + } + if !strings.Contains(got, want2) { + t.Errorf("expected %q in output, got:\n%s", want2, got) + } + if strings.Contains(got, "reasoning_delta run_id=") || strings.Contains(got, `detail=""`) { + t.Errorf("reasoning delta should not fall through to detail, got:\n%s", got) + } + if strings.Contains(got, "[node0-msg]") { + t.Errorf("reasoning delta should not appear as msg content, got:\n%s", got) + } +} + func TestEventRouterPrintsNodeLifecycleEvents(t *testing.T) { var out bytes.Buffer reg := edgenode.NewRegistry() diff --git a/apps/node/internal/node/node.go b/apps/node/internal/node/node.go index 1c6a744..c49c548 100644 --- a/apps/node/internal/node/node.go +++ b/apps/node/internal/node/node.go @@ -823,6 +823,16 @@ func (s *sessionSink) printEvent(event runtime.RuntimeEvent) { } fmt.Fprint(s.out, event.Delta) s.lineEnded = strings.HasSuffix(event.Delta, "\n") + case runtime.EventTypeReasoningDelta: + if event.Delta == "" { + return + } + if s.streaming && !s.lineEnded { + fmt.Fprintln(s.out) + } + s.streaming = false + s.lineEnded = true + printPrefixedLines(s.out, "[node-reasoning] ", event.Delta) case runtime.EventTypeComplete: if s.streaming && !s.lineEnded { fmt.Fprintln(s.out) @@ -878,6 +888,16 @@ func printTaggedMessage(out io.Writer, tag, message string) { fmt.Fprintf(out, "[%s] %s\n", tag, message) } +func printPrefixedLines(out io.Writer, prefix, text string) { + text = strings.TrimRight(text, "\n") + if text == "" { + return + } + for _, line := range strings.Split(text, "\n") { + fmt.Fprintf(out, "%s%s\n", prefix, line) + } +} + func structAsMap(s *structpb.Struct) map[string]any { if s == nil { return nil diff --git a/apps/node/internal/node/sink_test.go b/apps/node/internal/node/sink_test.go index aac62d3..890649d 100644 --- a/apps/node/internal/node/sink_test.go +++ b/apps/node/internal/node/sink_test.go @@ -358,3 +358,79 @@ func TestTerminalDeferringSinkSynthesizedTerminalNotDuplicated(t *testing.T) { t.Fatal("expected terminalObserved to be true after flush") } } + +func TestSessionSinkPrintEventPreservesReasoningDelta(t *testing.T) { + ms := &mockSender{} + var out bytes.Buffer + sink := &sessionSink{ + sess: ms, + out: &out, + nodeID: "test-node", + sessionID: "session-x", + } + + runID := "run-reasoning" + + // Emit start. + if err := sink.Emit(context.Background(), runtime.RuntimeEvent{ + RunID: runID, + Type: runtime.EventTypeStart, + Timestamp: time.Now(), + }); err != nil { + t.Fatalf("Emit start: %v", err) + } + + // Emit reasoning delta with multiline content. + reasoning := runtime.RuntimeEvent{ + RunID: runID, + Type: runtime.EventTypeReasoningDelta, + Delta: "사용자가 인사했습니다.\n😀\n", + Timestamp: time.Now(), + } + if err := sink.Emit(context.Background(), reasoning); err != nil { + t.Fatalf("Emit reasoning_delta: %v", err) + } + + // Emit complete. + if err := sink.Emit(context.Background(), runtime.RuntimeEvent{ + RunID: runID, + Type: runtime.EventTypeComplete, + Message: "done", + Timestamp: time.Now(), + }); err != nil { + t.Fatalf("Emit complete: %v", err) + } + + got := out.String() + + // Reasoning delta should be preserved with [node-reasoning] prefix. + if !strings.Contains(got, "[node-reasoning] 사용자가 인사했습니다.") { + t.Fatalf("expected [node-reasoning] line 1, got %q", got) + } + if !strings.Contains(got, "[node-reasoning] 😀") { + if !strings.Contains(got, "[node-reasoning] ") { + t.Fatalf("expected [node-reasoning] prefix in stdout, got %q", got) + } + } + + // stdout should NOT contain the raw reasoning_delta fallback. + if strings.Contains(got, "reasoning_delta run_id=") && strings.Contains(got, "detail=\"\"") { + t.Fatalf("reasoning_delta should not fall through to default case, got %q", got) + } + + // [node-message] should not appear for reasoning-only output. + if strings.Contains(got, "[node-message]") { + t.Fatalf("[node-message] should not appear for reasoning-only output, got %q", got) + } + + // Verify proto event preserves reasoning_delta type and delta content. + if len(ms.sentEvents) != 3 { + t.Fatalf("expected 3 proto events, got %d", len(ms.sentEvents)) + } + if ms.sentEvents[1].GetType() != "reasoning_delta" { + t.Fatalf("proto event type = %q, want reasoning_delta", ms.sentEvents[1].GetType()) + } + if ms.sentEvents[1].Delta != "사용자가 인사했습니다.\n😀\n" { + t.Fatalf("proto event Delta = %q, expected multiline reasoning text", ms.sentEvents[1].Delta) + } +}