diff --git a/apps/node/internal/adapters/cli/emitters.go b/apps/node/internal/adapters/cli/emitters.go index 85a1d69..cdebd42 100644 --- a/apps/node/internal/adapters/cli/emitters.go +++ b/apps/node/internal/adapters/cli/emitters.go @@ -172,18 +172,39 @@ func (e codexJSONEmitter) Emit(line string) ([]runtime.RuntimeEvent, error) { var ev struct { Type string `json:"type"` Item struct { - Type string `json:"type"` - Text string `json:"text"` + Type string `json:"type"` + Text string `json:"text"` + Delta string `json:"delta"` + Message struct { + Content []struct { + Type string `json:"type"` + Text string `json:"text"` + } `json:"content"` + } `json:"message"` } `json:"item"` + Delta string `json:"delta"` + Text string `json:"text"` Message string `json:"message"` Error struct { Message string `json:"message"` } `json:"error"` + Content []struct { + Type string `json:"type"` + Text string `json:"text"` + } `json:"content"` } if err := json.Unmarshal([]byte(line), &ev); err != nil { return nil, nil } switch ev.Type { + case "item.delta", "item.updated", "response.output_text.delta", "output_text.delta": + if delta := codexDeltaText(ev); delta != "" { + return []runtime.RuntimeEvent{{ + Type: runtime.EventTypeDelta, + Delta: delta, + }}, nil + } + return nil, nil case "item.completed": if ev.Item.Type == "agent_message" && ev.Item.Text != "" { return []runtime.RuntimeEvent{{ @@ -191,6 +212,12 @@ func (e codexJSONEmitter) Emit(line string) ([]runtime.RuntimeEvent, error) { Delta: ev.Item.Text, }}, nil } + if delta := codexDeltaText(ev); delta != "" { + return []runtime.RuntimeEvent{{ + Type: runtime.EventTypeDelta, + Delta: delta, + }}, nil + } return nil, nil case "error": return []runtime.RuntimeEvent{{ @@ -206,6 +233,52 @@ func (e codexJSONEmitter) Emit(line string) ([]runtime.RuntimeEvent, error) { return nil, nil } +func codexDeltaText(ev struct { + Type string `json:"type"` + Item struct { + Type string `json:"type"` + Text string `json:"text"` + Delta string `json:"delta"` + Message struct { + Content []struct { + Type string `json:"type"` + Text string `json:"text"` + } `json:"content"` + } `json:"message"` + } `json:"item"` + Delta string `json:"delta"` + Text string `json:"text"` + Message string `json:"message"` + Error struct { + Message string `json:"message"` + } `json:"error"` + Content []struct { + Type string `json:"type"` + Text string `json:"text"` + } `json:"content"` +}) string { + if ev.Item.Delta != "" { + return ev.Item.Delta + } + if ev.Delta != "" { + return ev.Delta + } + if ev.Text != "" { + return ev.Text + } + for _, part := range ev.Item.Message.Content { + if part.Type == "output_text" && part.Text != "" { + return part.Text + } + } + for _, part := range ev.Content { + if part.Type == "output_text" && part.Text != "" { + return part.Text + } + } + return "" +} + // --- opencode-json emitter --- type opencodeJSONEmitter struct{} @@ -214,7 +287,7 @@ func (opencodeJSONEmitter) Name() string { return "opencode-json" } func (e opencodeJSONEmitter) Emit(line string) ([]runtime.RuntimeEvent, error) { var ev struct { - Type string `json:"type"` + Type string `json:"type"` Error struct { Name string `json:"name"` Data struct { @@ -317,4 +390,4 @@ func (e clineJSONEmitter) Emit(line string) ([]runtime.RuntimeEvent, error) { return nil, nil } return nil, nil -} \ No newline at end of file +} diff --git a/apps/node/internal/adapters/cli/emitters_internal_test.go b/apps/node/internal/adapters/cli/emitters_internal_test.go index 7dcd62d..0bd300b 100644 --- a/apps/node/internal/adapters/cli/emitters_internal_test.go +++ b/apps/node/internal/adapters/cli/emitters_internal_test.go @@ -218,6 +218,42 @@ func TestCodexJSONEmitter_AgentMessageBecomesDelta(t *testing.T) { } } +func TestCodexJSONEmitter_ItemDeltaBecomesDelta(t *testing.T) { + e := codexJSONEmitter{} + line := `{"type":"item.delta","item":{"type":"agent_message","delta":"Doing."}}` + events, _ := e.Emit(line) + if len(events) != 1 { + t.Fatalf("expected 1 event, got %d", len(events)) + } + if events[0].Delta != "Doing." { + t.Fatalf("delta = %q, want %q", events[0].Delta, "Doing.") + } +} + +func TestCodexJSONEmitter_OutputTextDeltaBecomesDelta(t *testing.T) { + e := codexJSONEmitter{} + line := `{"type":"response.output_text.delta","delta":"chunk"}` + events, _ := e.Emit(line) + if len(events) != 1 { + t.Fatalf("expected 1 event, got %d", len(events)) + } + if events[0].Delta != "chunk" { + t.Fatalf("delta = %q, want %q", events[0].Delta, "chunk") + } +} + +func TestCodexJSONEmitter_ContentOutputTextFallsBackToDelta(t *testing.T) { + e := codexJSONEmitter{} + line := `{"type":"item.updated","item":{"type":"agent_message","message":{"content":[{"type":"output_text","text":"partial text"}]}}}` + events, _ := e.Emit(line) + if len(events) != 1 { + t.Fatalf("expected 1 event, got %d", len(events)) + } + if events[0].Delta != "partial text" { + t.Fatalf("delta = %q, want %q", events[0].Delta, "partial text") + } +} + func TestCodexJSONEmitter_TurnFailedBecomesError(t *testing.T) { e := codexJSONEmitter{} line := `{"type":"turn.failed","error":{"message":"quota exceeded"}}` @@ -512,4 +548,4 @@ func TestDriveJSONLines_OutputTokensCountedForDeltaOnly(t *testing.T) { if outputTokens != 3 { t.Fatalf("expected 3 outputTokens, got %d", outputTokens) } -} \ No newline at end of file +} diff --git a/apps/node/internal/node/node.go b/apps/node/internal/node/node.go index f736f94..1133119 100644 --- a/apps/node/internal/node/node.go +++ b/apps/node/internal/node/node.go @@ -282,7 +282,8 @@ type sessionSink struct { nodeID string sessionID string background bool - response strings.Builder + streaming bool + lineEnded bool } func (s *sessionSink) Emit(_ context.Context, event runtime.RuntimeEvent) error { @@ -314,16 +315,39 @@ func (s *sessionSink) printEvent(event runtime.RuntimeEvent) { } switch event.Type { case runtime.EventTypeStart: + s.streaming = false + s.lineEnded = true fmt.Fprintf(s.out, "[node-event] start run_id=%s\n", event.RunID) case runtime.EventTypeDelta: - s.response.WriteString(event.Delta) + if event.Delta == "" { + return + } + if !s.streaming { + fmt.Fprint(s.out, "[node-message] ") + s.streaming = true + } + fmt.Fprint(s.out, event.Delta) + s.lineEnded = strings.HasSuffix(event.Delta, "\n") case runtime.EventTypeComplete: + if s.streaming && !s.lineEnded { + fmt.Fprintln(s.out) + } + s.streaming = false + s.lineEnded = true fmt.Fprintf(s.out, "[node-event] complete run_id=%s detail=%q\n", event.RunID, event.Message) - printTaggedMessage(s.out, "node-message", s.response.String()) case runtime.EventTypeError: + if s.streaming && !s.lineEnded { + fmt.Fprintln(s.out) + } + s.streaming = false + s.lineEnded = true fmt.Fprintf(s.out, "[node-event] error run_id=%s detail=%q\n", event.RunID, event.Error) - printTaggedMessage(s.out, "node-message", s.response.String()) case runtime.EventTypeCancelled: + if s.streaming && !s.lineEnded { + fmt.Fprintln(s.out) + } + s.streaming = false + s.lineEnded = true fmt.Fprintf(s.out, "[node-event] cancelled run_id=%s\n", event.RunID) default: fmt.Fprintf(s.out, "[node-event] %s run_id=%s detail=%q\n", event.Type, event.RunID, event.Message) diff --git a/apps/node/internal/node/sink_test.go b/apps/node/internal/node/sink_test.go index 3a7f15b..60036cf 100644 --- a/apps/node/internal/node/sink_test.go +++ b/apps/node/internal/node/sink_test.go @@ -1,7 +1,9 @@ package node import ( + "bytes" "context" + "strings" "testing" "time" @@ -57,3 +59,68 @@ func TestSessionSinkEmitIncludesNodeID(t *testing.T) { t.Error("expected Background true") } } + +func TestSessionSinkPrintEventStreamsDeltaImmediately(t *testing.T) { + ms := &mockSender{} + var out bytes.Buffer + sink := &sessionSink{ + sess: ms, + out: &out, + nodeID: "test-node-123", + sessionID: "session-abc", + } + + start := runtime.RuntimeEvent{ + RunID: "run-xyz", + Type: runtime.EventTypeStart, + Timestamp: time.Now(), + } + delta1 := runtime.RuntimeEvent{ + RunID: "run-xyz", + Type: runtime.EventTypeDelta, + Delta: "Hello", + Timestamp: time.Now(), + } + delta2 := runtime.RuntimeEvent{ + RunID: "run-xyz", + Type: runtime.EventTypeDelta, + Delta: " world", + Timestamp: time.Now(), + } + complete := runtime.RuntimeEvent{ + RunID: "run-xyz", + Type: runtime.EventTypeComplete, + Message: "done", + Timestamp: time.Now(), + } + + if err := sink.Emit(context.Background(), start); err != nil { + t.Fatalf("Emit start failed: %v", err) + } + if err := sink.Emit(context.Background(), delta1); err != nil { + t.Fatalf("Emit delta1 failed: %v", err) + } + got := out.String() + if !strings.Contains(got, "[node-message] Hello") { + t.Fatalf("expected first delta to be printed immediately, got %q", got) + } + + if err := sink.Emit(context.Background(), delta2); err != nil { + t.Fatalf("Emit delta2 failed: %v", err) + } + got = out.String() + if !strings.Contains(got, "[node-message] Hello world") { + t.Fatalf("expected second delta to append to stream, got %q", got) + } + + if err := sink.Emit(context.Background(), complete); err != nil { + t.Fatalf("Emit complete failed: %v", err) + } + got = out.String() + if !strings.Contains(got, "[node-event] complete run_id=run-xyz detail=\"done\"") { + t.Fatalf("expected complete event to be printed, got %q", got) + } + if strings.Count(got, "[node-message]") != 1 { + t.Fatalf("expected node-message prefix once, got %q", got) + } +} diff --git a/configs/edge.yaml b/configs/edge.yaml index 108836f..4dfa242 100644 --- a/configs/edge.yaml +++ b/configs/edge.yaml @@ -13,7 +13,7 @@ metrics: console: adapter: "cli" - agent: "codex" + agent: "cline-dgx" session_id: "default" background: false timeout_sec: 300 @@ -51,6 +51,8 @@ nodes: - "yolo" - "--output-format" - "stream-json" + - "-m" + - "gemini-2.5-flash" - "-p" env: [] persistent: false @@ -81,7 +83,7 @@ nodes: - "--title" - "untitle" - "--model" - - "ollama-m1/qwen3.6:27b-coding-mxfp8" + - "ollama-dgx/qwen3.6:35b-a3b-bf16" - "--format" - "json" - "--dangerously-skip-permissions"