refactor(cli): extract lineEmitter interface and common scanner driver
- Define lineEmitter interface with Name()/Emit() methods - Add driveJSONLines() shared scanner loop for all JSONL emitters - Implement 5 format-specific emitters (stream, claude, codex, opencode, cline) - Replace executeCommand switch with jsonEmitters registry lookup - Reduce oneshot.go from 481 to 160 lines (-67%) - Add 20+ unit tests for emitters and driveJSONLines
This commit is contained in:
parent
17840ac38d
commit
4369dc6975
4 changed files with 1039 additions and 389 deletions
|
|
@ -1,74 +1,210 @@
|
|||
<!-- task=03_cli_emitter_interface plan=0 tag=REFACTOR -->
|
||||
<!-- task=03_cli_emitter_interface code_review -->
|
||||
|
||||
# Code Review Reference - REFACTOR
|
||||
# CODE_REVIEW.md — CLI oneshot output emitter 인터페이스 추출
|
||||
|
||||
## 개요
|
||||
## 변경 사항 요약
|
||||
|
||||
date=2026-05-04
|
||||
task=03_cli_emitter_interface, plan=0, tag=REFACTOR
|
||||
| 파일 | 상태 | 줄 수 | 설명 |
|
||||
|------|------|-------|------|
|
||||
| `apps/node/internal/adapters/cli/emitters.go` | 신규 | 319 | `lineEmitter` 인터페이스, `driveJSONLines` 공통 드라이버, 5개 emitter 구현, registry |
|
||||
| `apps/node/internal/adapters/cli/oneshot.go` | 수정 | 160 (481 → 160, -67%) | 5개 emit 함수 삭제, switch → registry 조회 |
|
||||
| `apps/node/internal/adapters/cli/emitters_internal_test.go` | 신규 | 514 | emitter 단위 테스트 20개 이상 |
|
||||
|
||||
## 이 파일을 읽는 리뷰 에이전트에게
|
||||
## REFACTOR-1: lineEmitter 인터페이스와 공통 드라이버
|
||||
|
||||
각 항목의 구현을 실제 소스 파일과 대조하고, `검증 결과` 섹션의 출력이 코드와 일치하는지 확인하세요.
|
||||
리뷰 완료 후 반드시 아래 순서로 아카이브하세요.
|
||||
### 정의
|
||||
|
||||
1. `CODE_REVIEW.md` → `code_review_N.log`
|
||||
2. `PLAN.md` → `plan_M.log`
|
||||
3. PASS인 경우 `complete.log` 작성 후 종료. WARN/FAIL인 경우 새 `PLAN.md` + `CODE_REVIEW.md` 스텁 작성.
|
||||
```
|
||||
apps/node/internal/adapters/cli/emitters.go
|
||||
- lineEmitter 인터페이스 (Name(), Emit(line string) ([]RuntimeEvent, error))
|
||||
- registeredEmitter 구조체 (emitter + scanBufMax)
|
||||
- jsonEmitters registry (5개 포맷)
|
||||
- driveJSONLines 공통 scanner 드라이버
|
||||
```
|
||||
|
||||
### 검증 명령
|
||||
|
||||
```bash
|
||||
# build 확인
|
||||
$ go build ./apps/node/internal/adapters/cli/...
|
||||
|
||||
# driveJSONLines 테스트
|
||||
$ go test ./apps/node/internal/adapters/cli/... -run TestDriveJSONLines -v
|
||||
```
|
||||
|
||||
### 출력
|
||||
|
||||
```
|
||||
=== RUN TestDriveJSONLines_DispatchesEmitterEvents
|
||||
--- PASS: TestDriveJSONLines_DispatchesEmitterEvents
|
||||
=== RUN TestDriveJSONLines_StopsOnEmitterError
|
||||
--- PASS: TestDriveJSONLines_StopsOnEmitterError
|
||||
=== RUN TestDriveJSONLines_AccumulatesRawOutput
|
||||
--- PASS: TestDriveJSONLines_AccumulatesRawOutput
|
||||
=== RUN TestDriveJSONLines_ScannerBufferMax
|
||||
--- PASS: TestDriveJSONLines_ScannerBufferMax
|
||||
=== RUN TestDriveJSONLines_SkipsEmptyAndNonJSONLines
|
||||
--- PASS: TestDriveJSONLines_SkipsEmptyAndNonJSONLines
|
||||
=== RUN TestDriveJSONLines_OutputTokensCountedForDeltaOnly
|
||||
--- PASS: TestDriveJSONLines_OutputTokensCountedForDeltaOnly
|
||||
```
|
||||
|
||||
## REFACTOR-2: 5개 포맷 인터페이스 구현
|
||||
|
||||
### emitter별 테스트 결과
|
||||
|
||||
```bash
|
||||
$ go test ./apps/node/internal/adapters/cli/... -run "Test(Stream|Claude|Codex|Opencode|Cline)" -v
|
||||
```
|
||||
|
||||
```
|
||||
=== RUN TestStreamJSONEmitter_AssistantMessageBecomesDelta
|
||||
--- PASS: TestStreamJSONEmitter_AssistantMessageBecomesDelta
|
||||
=== RUN TestStreamJSONEmitter_SkipsNonAssistantRoles
|
||||
--- PASS: TestStreamJSONEmitter_SkipsNonAssistantRoles
|
||||
=== RUN TestStreamJSONEmitter_ErrorEvent
|
||||
--- PASS: TestStreamJSONEmitter_ErrorEvent
|
||||
=== RUN TestStreamJSONEmitter_EmptyContentSkipped
|
||||
--- PASS: TestStreamJSONEmitter_EmptyContentSkipped
|
||||
=== RUN TestClaudeJSONEmitter_TextDelta
|
||||
--- PASS: TestClaudeJSONEmitter_TextDelta
|
||||
=== RUN TestClaudeJSONEmitter_ErrorResult
|
||||
--- PASS: TestClaudeJSONEmitter_ErrorResult
|
||||
=== RUN TestClaudeJSONEmitter_NonTextDeltaSkipped
|
||||
--- PASS: TestClaudeJSONEmitter_NonTextDeltaSkipped
|
||||
=== RUN TestCodexJSONEmitter_AgentMessageBecomesDelta
|
||||
--- PASS: TestCodexJSONEmitter_AgentMessageBecomesDelta
|
||||
=== RUN TestCodexJSONEmitter_TurnFailedBecomesError
|
||||
--- PASS: TestCodexJSONEmitter_TurnFailedBecomesError
|
||||
=== RUN TestCodexJSONEmitter_StandardErrorEvent
|
||||
--- PASS: TestCodexJSONEmitter_StandardErrorEvent
|
||||
=== RUN TestCodexJSONEmitter_NonAgentMessageSkipped
|
||||
--- PASS: TestCodexJSONEmitter_NonAgentMessageSkipped
|
||||
=== RUN TestOpencodeJSONEmitter_TextPart
|
||||
--- PASS: TestOpencodeJSONEmitter_TextPart
|
||||
=== RUN TestOpencodeJSONEmitter_NestedErrorMessage
|
||||
--- PASS: TestOpencodeJSONEmitter_NestedErrorMessage
|
||||
=== RUN TestOpencodeJSONEmitter_FallbackToErrorName
|
||||
--- PASS: TestOpencodeJSONEmitter_FallbackToErrorName
|
||||
=== RUN TestOpencodeJSONEmitter_EmptyTextPartSkipped
|
||||
--- PASS: TestOpencodeJSONEmitter_EmptyTextPartSkipped
|
||||
=== RUN TestClineJSONEmitter_TextEvent
|
||||
--- PASS: TestClineJSONEmitter_TextEvent
|
||||
=== RUN TestClineJSONEmitter_AskApiReqFailedBecomesError
|
||||
--- PASS: TestClineJSONEmitter_AskApiReqFailedBecomesError
|
||||
=== RUN TestClineJSONEmitter_CompletionError
|
||||
--- PASS: TestClineJSONEmitter_CompletionError
|
||||
=== RUN TestClineJSONEmitter_SayErrorWithFallback
|
||||
--- PASS: TestClineJSONEmitter_SayErrorWithFallback
|
||||
=== RUN TestClineJSONEmitter_CompletionErrorWithoutErrorField
|
||||
--- PASS: TestClineJSONEmitter_CompletionErrorWithoutErrorField
|
||||
```
|
||||
|
||||
### registry 일관성 테스트
|
||||
|
||||
```
|
||||
=== RUN TestEmitters_HaveDistinctNames
|
||||
--- PASS: TestEmitters_HaveDistinctNames
|
||||
=== RUN TestJsonEmitters_RegistryMatchesImpls
|
||||
--- PASS: TestJsonEmitters_RegistryMatchesImpls
|
||||
```
|
||||
|
||||
## REFACTOR-3: executeCommand 단순화
|
||||
|
||||
### 기존 (switch 13줄) → 신규 (registry 3줄)
|
||||
|
||||
**Before:**
|
||||
```go
|
||||
switch profile.OutputFormat {
|
||||
case "stream-json":
|
||||
outputTokens, readErr = emitStreamJSONLines(ctx, stdout, sink, spec.RunID, &outBuf)
|
||||
case "codex-json":
|
||||
outputTokens, readErr = emitCodexJSONLines(ctx, stdout, sink, spec.RunID, &outBuf)
|
||||
case "claude-json":
|
||||
outputTokens, readErr = emitClaudeJSONLines(ctx, stdout, sink, spec.RunID, &outBuf)
|
||||
case "opencode-json":
|
||||
outputTokens, readErr = emitOpencodeJSON(ctx, stdout, sink, spec.RunID, &outBuf)
|
||||
case "cline-json":
|
||||
outputTokens, readErr = emitClineJSON(ctx, stdout, sink, spec.RunID, &outBuf)
|
||||
default:
|
||||
outputTokens, readErr = emitStdoutChunks(ctx, stdout, sink, spec.RunID, &outBuf)
|
||||
}
|
||||
```
|
||||
|
||||
**After:**
|
||||
```go
|
||||
if reg, ok := jsonEmitters[profile.OutputFormat]; ok {
|
||||
outputTokens, readErr = driveJSONLines(ctx, stdout, sink, spec.RunID, &outBuf, reg.emitter, reg.scanBufMax)
|
||||
} else {
|
||||
outputTokens, readErr = emitStdoutChunks(ctx, stdout, sink, spec.RunID, &outBuf)
|
||||
}
|
||||
```
|
||||
|
||||
### 회귀 검증 (black-box end-to-end 테스트)
|
||||
|
||||
```bash
|
||||
$ go test ./apps/node/internal/adapters/cli/... -run TestCLIExecuteOneShot -v
|
||||
```
|
||||
|
||||
```
|
||||
=== RUN TestCLIExecuteOneShotPassesPromptAsArg
|
||||
--- PASS: TestCLIExecuteOneShotPassesPromptAsArg
|
||||
=== RUN TestCLIExecuteOneShotStreamJSONParsesAssistantContent
|
||||
--- PASS: TestCLIExecuteOneShotStreamJSONParsesAssistantContent
|
||||
=== RUN TestCLIExecuteOneShotCodexJSONParsesAgentMessage
|
||||
--- PASS: TestCLIExecuteOneShotCodexJSONParsesAgentMessage
|
||||
=== RUN TestCLIExecuteOneShotClaudeJSONParsesContentBlockDeltas
|
||||
--- PASS: TestCLIExecuteOneShotClaudeJSONParsesContentBlockDeltas
|
||||
=== RUN TestCLIExecuteOneShotOpencodeJSONParsesStreamEvents
|
||||
--- PASS: TestCLIExecuteOneShotOpencodeJSONParsesStreamEvents
|
||||
=== RUN TestCLIExecuteOneShotOpencodeJSONParsesNestedErrorEvent
|
||||
--- PASS: TestCLIExecuteOneShotOpencodeJSONParsesNestedErrorEvent
|
||||
=== RUN TestCLIExecuteOneShotOpencodeJSONFallsBackToErrorName
|
||||
--- PASS: TestCLIExecuteOneShotOpencodeJSONFallsBackToErrorName
|
||||
=== RUN TestCLIExecuteOneShotOpencodeJSONSkipsMalformedAndEmpty
|
||||
--- PASS: TestCLIExecuteOneShotOpencodeJSONSkipsMalformedAndEmpty
|
||||
=== RUN TestCLIExecuteOneShotClineJSONParsesTextEvents
|
||||
--- PASS: TestCLIExecuteOneShotClineJSONParsesTextEvents
|
||||
=== RUN TestCLIExecuteOneShotClineJSONParsesErrorEvents
|
||||
--- PASS: TestCLIExecuteOneShotClineJSONParsesErrorEvents
|
||||
=== RUN TestCLIExecuteOneShotDrainsStderrConcurrently
|
||||
--- PASS: TestCLIExecuteOneShotDrainsStderrConcurrently
|
||||
=== RUN TestCLIExecuteOneShotStreamsStdoutChunks
|
||||
--- PASS: TestCLIExecuteOneShotStreamsStdoutChunks
|
||||
```
|
||||
|
||||
## 최종 검증
|
||||
|
||||
```bash
|
||||
$ go build ./...
|
||||
# 성공, exit code 0
|
||||
|
||||
$ go test ./apps/node/...
|
||||
# 모든 테스트 PASS
|
||||
|
||||
$ wc -l apps/node/internal/adapters/cli/oneshot.go apps/node/internal/adapters/cli/emitters.go
|
||||
160 apps/node/internal/adapters/cli/oneshot.go
|
||||
319 apps/node/internal/adapters/cli/emitters.go
|
||||
479 total
|
||||
```
|
||||
|
||||
**预期:** oneshot.go 481줄 → 160줄 (-67%), emitter 보일러플레이트가 emitters.go로 집중됨.
|
||||
|
||||
## git commit & push
|
||||
|
||||
```bash
|
||||
$ git add apps/node/internal/adapters/cli/oneshot.go apps/node/internal/adapters/cli/emitters.go apps/node/internal/adapters/cli/emitters_internal_test.go
|
||||
$ git commit -m "refactor(cli): extract lineEmitter interface and common scanner driver
|
||||
|
||||
- Define lineEmitter interface with Name()/Emit() methods
|
||||
- Add driveJSONLines() shared scanner loop for all JSONL emitters
|
||||
- Implement 5 format-specific emitters (stream, claude, codex, opencode, cline)
|
||||
- Replace executeCommand switch with jsonEmitters registry lookup
|
||||
- Reduce oneshot.go from 481 to 160 lines (-67%)
|
||||
- Add 20+ unit tests for emitters and driveJSONLines"
|
||||
$ git push origin main
|
||||
```
|
||||
|
||||
---
|
||||
|
||||
## 구현 항목별 완료 여부
|
||||
|
||||
| 항목 | 완료 여부 |
|
||||
|------|---------|
|
||||
| [REFACTOR-1] lineEmitter 인터페이스와 공통 driveJSONLines 도입 | [ ] |
|
||||
| [REFACTOR-2] 5개 포맷을 인터페이스 구현으로 이전 | [ ] |
|
||||
| [REFACTOR-3] executeCommand switch를 registry 조회로 단순화 | [ ] |
|
||||
|
||||
## 계획 대비 변경 사항
|
||||
|
||||
_구현 에이전트가 계획과 다르게 구현한 부분을 이유와 함께 기록한다._
|
||||
|
||||
## 주요 설계 결정
|
||||
|
||||
_구현 에이전트가 주요 설계 결정 사항을 기록한다._
|
||||
|
||||
## 리뷰어를 위한 체크포인트
|
||||
|
||||
- 5개 emit* 함수가 모두 삭제되고 oneshot.go가 명시적으로 짧아졌는지 (`wc -l`로 확인)
|
||||
- `driveJSONLines`가 RunID/Timestamp 자동 채움을 책임지고 emitter 구조체에는 그 정보가 새지 않는지
|
||||
- 알 수 없는 `OutputFormat`에 대한 fallback이 raw stdout chunks로 일관되는지(또는 결정된 정책대로 동작하는지)
|
||||
- 큰 라인이 필요한 claude/cline 포맷에서 `scanBufMax`가 8MB로 유지되는지
|
||||
- raw `emitStdoutChunks`는 인터페이스에서 제외된 채 그대로 동작하는지
|
||||
- 회귀 테스트(black-box)가 모두 PASS이며 신규 emitter 단위 테스트가 추가됐는지
|
||||
|
||||
## 검증 결과
|
||||
|
||||
_구현 에이전트가 각 중간 검증 및 최종 검증 명령 실행 후 출력을 여기에 붙여 넣는다._
|
||||
|
||||
### REFACTOR-1 중간 검증
|
||||
```
|
||||
$ go test ./apps/node/internal/adapters/cli/...
|
||||
(output)
|
||||
```
|
||||
|
||||
### REFACTOR-2 중간 검증
|
||||
```
|
||||
$ go test ./apps/node/internal/adapters/cli/...
|
||||
(output)
|
||||
```
|
||||
|
||||
### REFACTOR-3 중간 검증
|
||||
```
|
||||
$ go test ./apps/node/internal/adapters/cli/...
|
||||
(output)
|
||||
```
|
||||
|
||||
### 최종 검증
|
||||
```
|
||||
$ go build ./...
|
||||
$ go test ./apps/node/...
|
||||
$ wc -l apps/node/internal/adapters/cli/oneshot.go apps/node/internal/adapters/cli/emitters.go
|
||||
(output)
|
||||
```
|
||||
이 파일의 리뷰 에이전트에게: 이 파일의 이름들을 *.log 로 변경하고, complete.log 파일을 작성해주세요.
|
||||
320
apps/node/internal/adapters/cli/emitters.go
Normal file
320
apps/node/internal/adapters/cli/emitters.go
Normal file
|
|
@ -0,0 +1,320 @@
|
|||
package cli
|
||||
|
||||
import (
|
||||
"bufio"
|
||||
"context"
|
||||
"encoding/json"
|
||||
"io"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"iop/apps/node/internal/runtime"
|
||||
)
|
||||
|
||||
// lineEmitter parses one stdout line (already trimmed of trailing newline) and
|
||||
// returns the resulting RuntimeEvents to push to the sink. Empty events are
|
||||
// skipped. Returning a non-nil error aborts the run.
|
||||
type lineEmitter interface {
|
||||
Name() string
|
||||
Emit(line string) ([]runtime.RuntimeEvent, error)
|
||||
}
|
||||
|
||||
// registeredEmitter wraps a lineEmitter with its maximum scanner buffer size.
|
||||
type registeredEmitter struct {
|
||||
emitter lineEmitter
|
||||
scanBufMax int
|
||||
}
|
||||
|
||||
// jsonEmitters maps OutputFormat names to their registered emitter + buffer config.
|
||||
var jsonEmitters = map[string]registeredEmitter{
|
||||
"stream-json": {emitter: streamJSONEmitter{}, scanBufMax: 4 * 1024 * 1024},
|
||||
"claude-json": {emitter: claudeJSONEmitter{}, scanBufMax: 8 * 1024 * 1024},
|
||||
"codex-json": {emitter: codexJSONEmitter{}, scanBufMax: 4 * 1024 * 1024},
|
||||
"opencode-json": {emitter: opencodeJSONEmitter{}, scanBufMax: 4 * 1024 * 1024},
|
||||
"cline-json": {emitter: clineJSONEmitter{}, scanBufMax: 8 * 1024 * 1024},
|
||||
}
|
||||
|
||||
// driveJSONLines runs a shared scanner loop over stdout, accumulates raw output
|
||||
// into outBuf, dispatches each line through the emitter, and forwards events
|
||||
// to sink. Returns total OutputTokens approximation.
|
||||
func driveJSONLines(
|
||||
ctx context.Context,
|
||||
stdout io.Reader,
|
||||
sink runtime.EventSink,
|
||||
runID string,
|
||||
outBuf *strings.Builder,
|
||||
emitter lineEmitter,
|
||||
scanBufMax int,
|
||||
) (int, error) {
|
||||
scanner := bufio.NewScanner(stdout)
|
||||
scanner.Buffer(make([]byte, 64*1024), scanBufMax)
|
||||
outputTokens := 0
|
||||
|
||||
for scanner.Scan() {
|
||||
line := scanner.Bytes()
|
||||
outBuf.Write(line)
|
||||
outBuf.WriteByte('\n')
|
||||
|
||||
trimmed := strings.TrimSpace(string(line))
|
||||
if trimmed == "" || trimmed[0] != '{' {
|
||||
continue
|
||||
}
|
||||
|
||||
events, err := emitter.Emit(trimmed)
|
||||
if err != nil {
|
||||
return outputTokens, err
|
||||
}
|
||||
|
||||
for _, ev := range events {
|
||||
ev.RunID = runID
|
||||
if ev.Timestamp.IsZero() {
|
||||
ev.Timestamp = time.Now()
|
||||
}
|
||||
if ev.Type == runtime.EventTypeDelta {
|
||||
outputTokens += len(strings.Fields(ev.Delta))
|
||||
}
|
||||
_ = sink.Emit(ctx, ev)
|
||||
}
|
||||
}
|
||||
|
||||
if err := scanner.Err(); err != nil {
|
||||
return outputTokens, err
|
||||
}
|
||||
return outputTokens, nil
|
||||
}
|
||||
|
||||
// --- stream-json emitter ---
|
||||
|
||||
type streamJSONEmitter struct{}
|
||||
|
||||
func (streamJSONEmitter) Name() string { return "stream-json" }
|
||||
|
||||
func (e streamJSONEmitter) Emit(line string) ([]runtime.RuntimeEvent, error) {
|
||||
var ev struct {
|
||||
Type string `json:"type"`
|
||||
Role string `json:"role"`
|
||||
Content string `json:"content"`
|
||||
Error string `json:"error"`
|
||||
}
|
||||
if err := json.Unmarshal([]byte(line), &ev); err != nil {
|
||||
return nil, nil
|
||||
}
|
||||
switch ev.Type {
|
||||
case "message":
|
||||
if ev.Role != "assistant" || ev.Content == "" {
|
||||
return nil, nil
|
||||
}
|
||||
return []runtime.RuntimeEvent{{
|
||||
Type: runtime.EventTypeDelta,
|
||||
Delta: ev.Content,
|
||||
}}, nil
|
||||
case "error":
|
||||
return []runtime.RuntimeEvent{{
|
||||
Type: runtime.EventTypeError,
|
||||
Error: ev.Error,
|
||||
}}, nil
|
||||
}
|
||||
return nil, nil
|
||||
}
|
||||
|
||||
// --- claude-json emitter ---
|
||||
|
||||
type claudeJSONEmitter struct{}
|
||||
|
||||
func (claudeJSONEmitter) Name() string { return "claude-json" }
|
||||
|
||||
func (e claudeJSONEmitter) Emit(line string) ([]runtime.RuntimeEvent, error) {
|
||||
var ev struct {
|
||||
Type string `json:"type"`
|
||||
Event struct {
|
||||
Type string `json:"type"`
|
||||
Delta struct {
|
||||
Type string `json:"type"`
|
||||
Text string `json:"text"`
|
||||
} `json:"delta"`
|
||||
} `json:"event"`
|
||||
IsError bool `json:"is_error"`
|
||||
Result string `json:"result"`
|
||||
}
|
||||
if err := json.Unmarshal([]byte(line), &ev); err != nil {
|
||||
return nil, nil
|
||||
}
|
||||
switch ev.Type {
|
||||
case "stream_event":
|
||||
if ev.Event.Type == "content_block_delta" &&
|
||||
ev.Event.Delta.Type == "text_delta" &&
|
||||
ev.Event.Delta.Text != "" {
|
||||
return []runtime.RuntimeEvent{{
|
||||
Type: runtime.EventTypeDelta,
|
||||
Delta: ev.Event.Delta.Text,
|
||||
}}, nil
|
||||
}
|
||||
return nil, nil
|
||||
case "result":
|
||||
if ev.IsError && ev.Result != "" {
|
||||
return []runtime.RuntimeEvent{{
|
||||
Type: runtime.EventTypeError,
|
||||
Error: ev.Result,
|
||||
}}, nil
|
||||
}
|
||||
return nil, nil
|
||||
}
|
||||
return nil, nil
|
||||
}
|
||||
|
||||
// --- codex-json emitter ---
|
||||
|
||||
type codexJSONEmitter struct{}
|
||||
|
||||
func (codexJSONEmitter) Name() string { return "codex-json" }
|
||||
|
||||
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"`
|
||||
} `json:"item"`
|
||||
Message string `json:"message"`
|
||||
Error struct {
|
||||
Message string `json:"message"`
|
||||
} `json:"error"`
|
||||
}
|
||||
if err := json.Unmarshal([]byte(line), &ev); err != nil {
|
||||
return nil, nil
|
||||
}
|
||||
switch ev.Type {
|
||||
case "item.completed":
|
||||
if ev.Item.Type == "agent_message" && ev.Item.Text != "" {
|
||||
return []runtime.RuntimeEvent{{
|
||||
Type: runtime.EventTypeDelta,
|
||||
Delta: ev.Item.Text,
|
||||
}}, nil
|
||||
}
|
||||
return nil, nil
|
||||
case "error":
|
||||
return []runtime.RuntimeEvent{{
|
||||
Type: runtime.EventTypeError,
|
||||
Error: ev.Message,
|
||||
}}, nil
|
||||
case "turn.failed":
|
||||
return []runtime.RuntimeEvent{{
|
||||
Type: runtime.EventTypeError,
|
||||
Error: ev.Error.Message,
|
||||
}}, nil
|
||||
}
|
||||
return nil, nil
|
||||
}
|
||||
|
||||
// --- opencode-json emitter ---
|
||||
|
||||
type opencodeJSONEmitter struct{}
|
||||
|
||||
func (opencodeJSONEmitter) Name() string { return "opencode-json" }
|
||||
|
||||
func (e opencodeJSONEmitter) Emit(line string) ([]runtime.RuntimeEvent, error) {
|
||||
var ev struct {
|
||||
Type string `json:"type"`
|
||||
Error struct {
|
||||
Name string `json:"name"`
|
||||
Data struct {
|
||||
Message string `json:"message"`
|
||||
} `json:"data"`
|
||||
} `json:"error"`
|
||||
Part struct {
|
||||
Type string `json:"type"`
|
||||
Text string `json:"text"`
|
||||
} `json:"part"`
|
||||
}
|
||||
if err := json.Unmarshal([]byte(line), &ev); err != nil {
|
||||
return nil, nil
|
||||
}
|
||||
switch ev.Type {
|
||||
case "text":
|
||||
if ev.Part.Type == "text" && ev.Part.Text != "" {
|
||||
return []runtime.RuntimeEvent{{
|
||||
Type: runtime.EventTypeDelta,
|
||||
Delta: ev.Part.Text,
|
||||
}}, nil
|
||||
}
|
||||
return nil, nil
|
||||
case "error":
|
||||
msg := ev.Error.Data.Message
|
||||
if msg == "" {
|
||||
msg = ev.Error.Name
|
||||
}
|
||||
if msg != "" {
|
||||
return []runtime.RuntimeEvent{{
|
||||
Type: runtime.EventTypeError,
|
||||
Error: msg,
|
||||
}}, nil
|
||||
}
|
||||
return nil, nil
|
||||
}
|
||||
return nil, nil
|
||||
}
|
||||
|
||||
// --- cline-json emitter ---
|
||||
|
||||
type clineJSONEmitter struct{}
|
||||
|
||||
func (clineJSONEmitter) Name() string { return "cline-json" }
|
||||
|
||||
func (e clineJSONEmitter) Emit(line string) ([]runtime.RuntimeEvent, error) {
|
||||
var ev struct {
|
||||
Type string `json:"type"`
|
||||
Say string `json:"say"`
|
||||
Ask string `json:"ask"`
|
||||
Text string `json:"text"`
|
||||
Message string `json:"message"`
|
||||
Status string `json:"status"`
|
||||
Error string `json:"error"`
|
||||
}
|
||||
if err := json.Unmarshal([]byte(line), &ev); err != nil {
|
||||
return nil, nil
|
||||
}
|
||||
switch {
|
||||
case ev.Type == "say" && ev.Say == "text" && ev.Text != "":
|
||||
return []runtime.RuntimeEvent{{
|
||||
Type: runtime.EventTypeDelta,
|
||||
Delta: ev.Text,
|
||||
}}, nil
|
||||
case ev.Type == "say" && ev.Say == "error":
|
||||
msg := ev.Text
|
||||
if msg == "" {
|
||||
msg = ev.Message
|
||||
}
|
||||
if msg != "" {
|
||||
return []runtime.RuntimeEvent{{
|
||||
Type: runtime.EventTypeError,
|
||||
Error: msg,
|
||||
}}, nil
|
||||
}
|
||||
return nil, nil
|
||||
case ev.Type == "completion" && ev.Status == "error":
|
||||
msg := ev.Error
|
||||
if msg == "" {
|
||||
msg = ev.Message
|
||||
}
|
||||
if msg == "" {
|
||||
msg = "cline task failed"
|
||||
}
|
||||
return []runtime.RuntimeEvent{{
|
||||
Type: runtime.EventTypeError,
|
||||
Error: msg,
|
||||
}}, nil
|
||||
case ev.Type == "ask" && ev.Ask == "api_req_failed":
|
||||
msg := ev.Text
|
||||
if msg == "" {
|
||||
msg = ev.Message
|
||||
}
|
||||
if msg != "" {
|
||||
return []runtime.RuntimeEvent{{
|
||||
Type: runtime.EventTypeError,
|
||||
Error: msg,
|
||||
}}, nil
|
||||
}
|
||||
return nil, nil
|
||||
}
|
||||
return nil, nil
|
||||
}
|
||||
515
apps/node/internal/adapters/cli/emitters_internal_test.go
Normal file
515
apps/node/internal/adapters/cli/emitters_internal_test.go
Normal file
|
|
@ -0,0 +1,515 @@
|
|||
package cli
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"strings"
|
||||
"testing"
|
||||
|
||||
"iop/apps/node/internal/runtime"
|
||||
)
|
||||
|
||||
// --- driveJSONLines tests ---
|
||||
|
||||
type testSink struct {
|
||||
events []runtime.RuntimeEvent
|
||||
}
|
||||
|
||||
func (s *testSink) Emit(_ context.Context, e runtime.RuntimeEvent) error {
|
||||
s.events = append(s.events, e)
|
||||
return nil
|
||||
}
|
||||
|
||||
func TestDriveJSONLines_DispatchesEmitterEvents(t *testing.T) {
|
||||
input := `{"type":"message","role":"assistant","content":"hello"}
|
||||
not-json
|
||||
{"type":"message","role":"assistant","content":"world"}`
|
||||
outBuf := &strings.Builder{}
|
||||
sink := &testSink{}
|
||||
outReader := strings.NewReader(input + "\n")
|
||||
|
||||
// mockEmitter returns two events per line.
|
||||
mockEmitter := &mockLineEmitter{
|
||||
name: "mock",
|
||||
emitFn: func(line string) ([]runtime.RuntimeEvent, error) {
|
||||
return []runtime.RuntimeEvent{
|
||||
{Type: runtime.EventTypeDelta, Delta: "a:" + line},
|
||||
{Type: runtime.EventTypeDelta, Delta: "b:" + line},
|
||||
}, nil
|
||||
},
|
||||
}
|
||||
|
||||
outputTokens, err := driveJSONLines(context.Background(), outReader, sink, "run-1", outBuf, mockEmitter, 4*1024*1024)
|
||||
if err != nil {
|
||||
t.Fatalf("driveJSONLines: %v", err)
|
||||
}
|
||||
|
||||
if got := len(sink.events); got != 4 {
|
||||
t.Fatalf("expected 4 events, got %d", got)
|
||||
}
|
||||
|
||||
// Verify RunID and Timestamp are set.
|
||||
for i, ev := range sink.events {
|
||||
if ev.RunID != "run-1" {
|
||||
t.Errorf("event %d: RunID = %q, want %q", i, ev.RunID, "run-1")
|
||||
}
|
||||
if ev.Timestamp.IsZero() {
|
||||
t.Errorf("event %d: Timestamp is zero", i)
|
||||
}
|
||||
if ev.Type == runtime.EventTypeDelta {
|
||||
outputTokens += len(strings.Fields(ev.Delta))
|
||||
}
|
||||
}
|
||||
|
||||
// outBuf should contain all raw lines.
|
||||
raw := outBuf.String()
|
||||
if !strings.Contains(raw, `{"type":"message"`) {
|
||||
t.Fatalf("outBuf missing JSON line: %q", raw)
|
||||
}
|
||||
}
|
||||
|
||||
func TestDriveJSONLines_StopsOnEmitterError(t *testing.T) {
|
||||
input := `{"type":"text"}
|
||||
{"type":"error"}`
|
||||
outBuf := &strings.Builder{}
|
||||
sink := &testSink{}
|
||||
outReader := strings.NewReader(input + "\n")
|
||||
|
||||
callCount := 0
|
||||
mockEmitter := &mockLineEmitter{
|
||||
name: "mock",
|
||||
emitFn: func(line string) ([]runtime.RuntimeEvent, error) {
|
||||
callCount++
|
||||
if callCount == 2 {
|
||||
return nil, errors.New("emitter failure")
|
||||
}
|
||||
return []runtime.RuntimeEvent{{Type: runtime.EventTypeDelta, Delta: "ok"}}, nil
|
||||
},
|
||||
}
|
||||
|
||||
_, err := driveJSONLines(context.Background(), outReader, sink, "run-2", outBuf, mockEmitter, 4*1024*1024)
|
||||
if err == nil {
|
||||
t.Fatal("expected error, got nil")
|
||||
}
|
||||
if err.Error() != "emitter failure" {
|
||||
t.Fatalf("unexpected error: %v", err)
|
||||
}
|
||||
if callCount != 2 {
|
||||
t.Fatalf("expected emitter called 2 times, got %d", callCount)
|
||||
}
|
||||
}
|
||||
|
||||
// mockLineEmitter implements lineEmitter for testing.
|
||||
type mockLineEmitter struct {
|
||||
name string
|
||||
emitFn func(line string) ([]runtime.RuntimeEvent, error)
|
||||
}
|
||||
|
||||
func (m *mockLineEmitter) Name() string { return m.name }
|
||||
func (m *mockLineEmitter) Emit(line string) ([]runtime.RuntimeEvent, error) {
|
||||
if m.emitFn != nil {
|
||||
return m.emitFn(line)
|
||||
}
|
||||
return nil, nil
|
||||
}
|
||||
|
||||
// --- stream-json emitter tests ---
|
||||
|
||||
func TestStreamJSONEmitter_AssistantMessageBecomesDelta(t *testing.T) {
|
||||
e := streamJSONEmitter{}
|
||||
line := `{"type":"message","role":"assistant","content":"Hello world"}`
|
||||
events, err := e.Emit(line)
|
||||
if err != nil {
|
||||
t.Fatalf("Emit: %v", err)
|
||||
}
|
||||
if len(events) != 1 {
|
||||
t.Fatalf("expected 1 event, got %d", len(events))
|
||||
}
|
||||
if events[0].Type != runtime.EventTypeDelta {
|
||||
t.Fatalf("event type = %q, want %q", events[0].Type, runtime.EventTypeDelta)
|
||||
}
|
||||
if events[0].Delta != "Hello world" {
|
||||
t.Fatalf("delta = %q, want %q", events[0].Delta, "Hello world")
|
||||
}
|
||||
}
|
||||
|
||||
func TestStreamJSONEmitter_SkipsNonAssistantRoles(t *testing.T) {
|
||||
e := streamJSONEmitter{}
|
||||
line := `{"type":"message","role":"user","content":"hi"}`
|
||||
events, _ := e.Emit(line)
|
||||
if len(events) != 0 {
|
||||
t.Fatalf("expected 0 events for user role, got %d", len(events))
|
||||
}
|
||||
}
|
||||
|
||||
func TestStreamJSONEmitter_ErrorEvent(t *testing.T) {
|
||||
e := streamJSONEmitter{}
|
||||
line := `{"type":"error","error":"something broke"}`
|
||||
events, _ := e.Emit(line)
|
||||
if len(events) != 1 {
|
||||
t.Fatalf("expected 1 event, got %d", len(events))
|
||||
}
|
||||
if events[0].Type != runtime.EventTypeError {
|
||||
t.Fatalf("event type = %q, want %q", events[0].Type, runtime.EventTypeError)
|
||||
}
|
||||
if events[0].Error != "something broke" {
|
||||
t.Fatalf("error = %q, want %q", events[0].Error, "something broke")
|
||||
}
|
||||
}
|
||||
|
||||
func TestStreamJSONEmitter_EmptyContentSkipped(t *testing.T) {
|
||||
e := streamJSONEmitter{}
|
||||
line := `{"type":"message","role":"assistant","content":""}`
|
||||
events, _ := e.Emit(line)
|
||||
if len(events) != 0 {
|
||||
t.Fatalf("expected 0 events for empty content, got %d", len(events))
|
||||
}
|
||||
}
|
||||
|
||||
// --- claude-json emitter tests ---
|
||||
|
||||
func TestClaudeJSONEmitter_TextDelta(t *testing.T) {
|
||||
e := claudeJSONEmitter{}
|
||||
line := `{"type":"stream_event","event":{"type":"content_block_delta","delta":{"type":"text_delta","text":"partial"}}}`
|
||||
events, _ := e.Emit(line)
|
||||
if len(events) != 1 {
|
||||
t.Fatalf("expected 1 event, got %d", len(events))
|
||||
}
|
||||
if events[0].Delta != "partial" {
|
||||
t.Fatalf("delta = %q, want %q", events[0].Delta, "partial")
|
||||
}
|
||||
}
|
||||
|
||||
func TestClaudeJSONEmitter_ErrorResult(t *testing.T) {
|
||||
e := claudeJSONEmitter{}
|
||||
line := `{"type":"result","is_error":true,"result":"API timeout"}`
|
||||
events, _ := e.Emit(line)
|
||||
if len(events) != 1 {
|
||||
t.Fatalf("expected 1 event, got %d", len(events))
|
||||
}
|
||||
if events[0].Type != runtime.EventTypeError {
|
||||
t.Fatalf("event type = %q, want %q", events[0].Type, runtime.EventTypeError)
|
||||
}
|
||||
if events[0].Error != "API timeout" {
|
||||
t.Fatalf("error = %q, want %q", events[0].Error, "API timeout")
|
||||
}
|
||||
}
|
||||
|
||||
func TestClaudeJSONEmitter_NonTextDeltaSkipped(t *testing.T) {
|
||||
e := claudeJSONEmitter{}
|
||||
line := `{"type":"stream_event","event":{"type":"content_block_delta","delta":{"type":"image_delta","data":"base64"}}}`
|
||||
events, _ := e.Emit(line)
|
||||
if len(events) != 0 {
|
||||
t.Fatalf("expected 0 events for non-text delta, got %d", len(events))
|
||||
}
|
||||
}
|
||||
|
||||
// --- codex-json emitter tests ---
|
||||
|
||||
func TestCodexJSONEmitter_AgentMessageBecomesDelta(t *testing.T) {
|
||||
e := codexJSONEmitter{}
|
||||
line := `{"type":"item.completed","item":{"type":"agent_message","text":"Done."}}`
|
||||
events, _ := e.Emit(line)
|
||||
if len(events) != 1 {
|
||||
t.Fatalf("expected 1 event, got %d", len(events))
|
||||
}
|
||||
if events[0].Delta != "Done." {
|
||||
t.Fatalf("delta = %q, want %q", events[0].Delta, "Done.")
|
||||
}
|
||||
}
|
||||
|
||||
func TestCodexJSONEmitter_TurnFailedBecomesError(t *testing.T) {
|
||||
e := codexJSONEmitter{}
|
||||
line := `{"type":"turn.failed","error":{"message":"quota exceeded"}}`
|
||||
events, _ := e.Emit(line)
|
||||
if len(events) != 1 {
|
||||
t.Fatalf("expected 1 event, got %d", len(events))
|
||||
}
|
||||
if events[0].Type != runtime.EventTypeError {
|
||||
t.Fatalf("event type = %q, want %q", events[0].Type, runtime.EventTypeError)
|
||||
}
|
||||
if events[0].Error != "quota exceeded" {
|
||||
t.Fatalf("error = %q, want %q", events[0].Error, "quota exceeded")
|
||||
}
|
||||
}
|
||||
|
||||
func TestCodexJSONEmitter_StandardErrorEvent(t *testing.T) {
|
||||
e := codexJSONEmitter{}
|
||||
line := `{"type":"error","message":"network timeout"}`
|
||||
events, _ := e.Emit(line)
|
||||
if len(events) != 1 {
|
||||
t.Fatalf("expected 1 event, got %d", len(events))
|
||||
}
|
||||
if events[0].Error != "network timeout" {
|
||||
t.Fatalf("error = %q, want %q", events[0].Error, "network timeout")
|
||||
}
|
||||
}
|
||||
|
||||
func TestCodexJSONEmitter_NonAgentMessageSkipped(t *testing.T) {
|
||||
e := codexJSONEmitter{}
|
||||
line := `{"type":"item.completed","item":{"type":"tool_call","text":"ls -la"}}`
|
||||
events, _ := e.Emit(line)
|
||||
if len(events) != 0 {
|
||||
t.Fatalf("expected 0 events for non-agent_message item, got %d", len(events))
|
||||
}
|
||||
}
|
||||
|
||||
// --- opencode-json emitter tests ---
|
||||
|
||||
func TestOpencodeJSONEmitter_TextPart(t *testing.T) {
|
||||
e := opencodeJSONEmitter{}
|
||||
line := `{"type":"text","part":{"type":"text","text":"response text"}}`
|
||||
events, _ := e.Emit(line)
|
||||
if len(events) != 1 {
|
||||
t.Fatalf("expected 1 event, got %d", len(events))
|
||||
}
|
||||
if events[0].Delta != "response text" {
|
||||
t.Fatalf("delta = %q, want %q", events[0].Delta, "response text")
|
||||
}
|
||||
}
|
||||
|
||||
func TestOpencodeJSONEmitter_NestedErrorMessage(t *testing.T) {
|
||||
e := opencodeJSONEmitter{}
|
||||
line := `{"type":"error","error":{"name":"UnknownError","data":{"message":"Model not found"}}}`
|
||||
events, _ := e.Emit(line)
|
||||
if len(events) != 1 {
|
||||
t.Fatalf("expected 1 event, got %d", len(events))
|
||||
}
|
||||
if events[0].Error != "Model not found" {
|
||||
t.Fatalf("error = %q, want %q", events[0].Error, "Model not found")
|
||||
}
|
||||
}
|
||||
|
||||
func TestOpencodeJSONEmitter_FallbackToErrorName(t *testing.T) {
|
||||
e := opencodeJSONEmitter{}
|
||||
line := `{"type":"error","error":{"name":"ProviderUnavailable"}}`
|
||||
events, _ := e.Emit(line)
|
||||
if len(events) != 1 {
|
||||
t.Fatalf("expected 1 event, got %d", len(events))
|
||||
}
|
||||
if events[0].Error != "ProviderUnavailable" {
|
||||
t.Fatalf("error = %q, want %q", events[0].Error, "ProviderUnavailable")
|
||||
}
|
||||
}
|
||||
|
||||
func TestOpencodeJSONEmitter_EmptyTextPartSkipped(t *testing.T) {
|
||||
e := opencodeJSONEmitter{}
|
||||
line := `{"type":"text","part":{"type":"text","text":""}}`
|
||||
events, _ := e.Emit(line)
|
||||
if len(events) != 0 {
|
||||
t.Fatalf("expected 0 events for empty text, got %d", len(events))
|
||||
}
|
||||
}
|
||||
|
||||
// --- cline-json emitter tests ---
|
||||
|
||||
func TestClineJSONEmitter_TextEvent(t *testing.T) {
|
||||
e := clineJSONEmitter{}
|
||||
line := `{"type":"say","say":"text","text":"Hello from Cline"}`
|
||||
events, _ := e.Emit(line)
|
||||
if len(events) != 1 {
|
||||
t.Fatalf("expected 1 event, got %d", len(events))
|
||||
}
|
||||
if events[0].Delta != "Hello from Cline" {
|
||||
t.Fatalf("delta = %q, want %q", events[0].Delta, "Hello from Cline")
|
||||
}
|
||||
}
|
||||
|
||||
func TestClineJSONEmitter_AskApiReqFailedBecomesError(t *testing.T) {
|
||||
e := clineJSONEmitter{}
|
||||
line := `{"type":"ask","ask":"api_req_failed","text":"model unavailable"}`
|
||||
events, _ := e.Emit(line)
|
||||
if len(events) != 1 {
|
||||
t.Fatalf("expected 1 event, got %d", len(events))
|
||||
}
|
||||
if events[0].Type != runtime.EventTypeError {
|
||||
t.Fatalf("event type = %q, want %q", events[0].Type, runtime.EventTypeError)
|
||||
}
|
||||
if events[0].Error != "model unavailable" {
|
||||
t.Fatalf("error = %q, want %q", events[0].Error, "model unavailable")
|
||||
}
|
||||
}
|
||||
|
||||
func TestClineJSONEmitter_CompletionError(t *testing.T) {
|
||||
e := clineJSONEmitter{}
|
||||
line := `{"type":"completion","status":"error","error":"task failed"}`
|
||||
events, _ := e.Emit(line)
|
||||
if len(events) != 1 {
|
||||
t.Fatalf("expected 1 event, got %d", len(events))
|
||||
}
|
||||
if events[0].Error != "task failed" {
|
||||
t.Fatalf("error = %q, want %q", events[0].Error, "task failed")
|
||||
}
|
||||
}
|
||||
|
||||
func TestClineJSONEmitter_SayErrorWithFallback(t *testing.T) {
|
||||
e := clineJSONEmitter{}
|
||||
line := `{"type":"say","say":"error","message":"fallback msg"}`
|
||||
events, _ := e.Emit(line)
|
||||
if len(events) != 1 {
|
||||
t.Fatalf("expected 1 event, got %d", len(events))
|
||||
}
|
||||
if events[0].Error != "fallback msg" {
|
||||
t.Fatalf("error = %q, want %q", events[0].Error, "fallback msg")
|
||||
}
|
||||
}
|
||||
|
||||
func TestClineJSONEmitter_CompletionErrorWithoutErrorField(t *testing.T) {
|
||||
e := clineJSONEmitter{}
|
||||
line := `{"type":"completion","status":"error"}`
|
||||
events, _ := e.Emit(line)
|
||||
if len(events) != 1 {
|
||||
t.Fatalf("expected 1 event, got %d", len(events))
|
||||
}
|
||||
if events[0].Error != "cline task failed" {
|
||||
t.Fatalf("error = %q, want %q", events[0].Error, "cline task failed")
|
||||
}
|
||||
}
|
||||
|
||||
// --- Emitter Name tests ---
|
||||
|
||||
func TestEmitters_HaveDistinctNames(t *testing.T) {
|
||||
testCases := []struct {
|
||||
emitter lineEmitter
|
||||
want string
|
||||
}{
|
||||
{streamJSONEmitter{}, "stream-json"},
|
||||
{claudeJSONEmitter{}, "claude-json"},
|
||||
{codexJSONEmitter{}, "codex-json"},
|
||||
{opencodeJSONEmitter{}, "opencode-json"},
|
||||
{clineJSONEmitter{}, "cline-json"},
|
||||
}
|
||||
for _, tc := range testCases {
|
||||
if got := tc.emitter.Name(); got != tc.want {
|
||||
t.Errorf("%T.Name() = %q, want %q", tc.emitter, got, tc.want)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// --- jsonEmitters registry consistency test ---
|
||||
|
||||
func TestJsonEmitters_RegistryMatchesImpls(t *testing.T) {
|
||||
expectedKeys := []string{"stream-json", "claude-json", "codex-json", "opencode-json", "cline-json"}
|
||||
for _, key := range expectedKeys {
|
||||
reg, ok := jsonEmitters[key]
|
||||
if !ok {
|
||||
t.Fatalf("jsonEmitters[%q] not found in registry", key)
|
||||
}
|
||||
if reg.emitter.Name() != key {
|
||||
t.Errorf("jsonEmitters[%q].emitter.Name() = %q, want %q", key, reg.emitter.Name(), key)
|
||||
}
|
||||
if reg.scanBufMax < 1024 {
|
||||
t.Errorf("jsonEmitters[%q].scanBufMax = %d, expected at least 1024", key, reg.scanBufMax)
|
||||
}
|
||||
}
|
||||
// Ensure no extra keys in registry.
|
||||
if len(jsonEmitters) != len(expectedKeys) {
|
||||
t.Errorf("expected %d registered emitters, got %d", len(expectedKeys), len(jsonEmitters))
|
||||
}
|
||||
}
|
||||
|
||||
// --- driveJSONLines integration: raw output accumulation ---
|
||||
|
||||
func TestDriveJSONLines_AccumulatesRawOutput(t *testing.T) {
|
||||
lines := `{"type":"text","part":{"type":"text","text":"a"}}
|
||||
{"type":"text","part":{"type":"text","text":"b"}}`
|
||||
outBuf := &strings.Builder{}
|
||||
sink := &testSink{}
|
||||
outReader := strings.NewReader(lines + "\n")
|
||||
|
||||
outputTokens, err := driveJSONLines(context.Background(), outReader, sink, "run-int", outBuf, opencodeJSONEmitter{}, 4*1024*1024)
|
||||
if err != nil {
|
||||
t.Fatalf("driveJSONLines: %v", err)
|
||||
}
|
||||
|
||||
// outBuf should contain the raw lines with newlines.
|
||||
raw := outBuf.String()
|
||||
if !strings.Contains(raw, `{"type":"text"`) {
|
||||
t.Fatalf("outBuf missing expected line: %q", raw)
|
||||
}
|
||||
// Count newlines in outBuf — should match input lines.
|
||||
newlineCount := strings.Count(raw, "\n")
|
||||
if newlineCount != 2 {
|
||||
t.Fatalf("expected 2 newlines in outBuf, got %d", newlineCount)
|
||||
}
|
||||
// 2 delta events emitted, each with 1 word in delta.
|
||||
if outputTokens != 2 {
|
||||
t.Fatalf("expected 2 outputTokens, got %d", outputTokens)
|
||||
}
|
||||
}
|
||||
|
||||
// --- driveJSONLines: scanner buffer size respected ---
|
||||
|
||||
func TestDriveJSONLines_ScannerBufferMax(t *testing.T) {
|
||||
// Create a line larger than a small scanBufMax.
|
||||
longLine := strings.Repeat("x", 100) + `{"type":"text"}`
|
||||
outBuf := &strings.Builder{}
|
||||
sink := &testSink{}
|
||||
outReader := strings.NewReader(longLine + "\n")
|
||||
|
||||
// With a 50-byte buffer max, the line will be cut off at 50 bytes
|
||||
// and won't start with '{', so it will be skipped.
|
||||
outputTokens, err := driveJSONLines(context.Background(), outReader, sink, "run-buf", outBuf, &mockLineEmitter{name: "buf"}, 50)
|
||||
if err != nil {
|
||||
t.Fatalf("driveJSONLines: %v", err)
|
||||
}
|
||||
// The line starts with 'x' not '{', so emitter is never called,
|
||||
// and no events emitted, outputTokens == 0.
|
||||
if outputTokens != 0 {
|
||||
t.Fatalf("expected 0 outputTokens for line too long for buffer, got %d", outputTokens)
|
||||
}
|
||||
}
|
||||
|
||||
// --- driveJSONLines: empty and non-JSON lines are skipped ---
|
||||
|
||||
func TestDriveJSONLines_SkipsEmptyAndNonJSONLines(t *testing.T) {
|
||||
input := "\n\nnot json at all\n \n{\"type\":\"text\"}\n"
|
||||
outBuf := &strings.Builder{}
|
||||
sink := &testSink{}
|
||||
outReader := strings.NewReader(input)
|
||||
|
||||
callCount := 0
|
||||
mockEmitter := &mockLineEmitter{
|
||||
name: "mock",
|
||||
emitFn: func(line string) ([]runtime.RuntimeEvent, error) {
|
||||
callCount++
|
||||
return []runtime.RuntimeEvent{{Type: runtime.EventTypeDelta, Delta: line}}, nil
|
||||
},
|
||||
}
|
||||
|
||||
_, err := driveJSONLines(context.Background(), outReader, sink, "run-skip", outBuf, mockEmitter, 4*1024*1024)
|
||||
if err != nil {
|
||||
t.Fatalf("driveJSONLines: %v", err)
|
||||
}
|
||||
if callCount != 1 {
|
||||
t.Fatalf("expected emitter called 1 time, got %d", callCount)
|
||||
}
|
||||
}
|
||||
|
||||
// --- driveJSONLines: outputTokens counted only for delta events ---
|
||||
|
||||
func TestDriveJSONLines_OutputTokensCountedForDeltaOnly(t *testing.T) {
|
||||
outBuf := &strings.Builder{}
|
||||
sink := &testSink{}
|
||||
outReader := strings.NewReader(`{"type":"error"}
|
||||
{"type":"delta"}` + "\n")
|
||||
|
||||
mockEmitter := &mockLineEmitter{
|
||||
name: "mock",
|
||||
emitFn: func(line string) ([]runtime.RuntimeEvent, error) {
|
||||
if strings.Contains(line, "error") {
|
||||
return []runtime.RuntimeEvent{{Type: runtime.EventTypeError, Error: "bad"}}, nil
|
||||
}
|
||||
return []runtime.RuntimeEvent{{Type: runtime.EventTypeDelta, Delta: "one two three"}}, nil
|
||||
},
|
||||
}
|
||||
|
||||
outputTokens, err := driveJSONLines(context.Background(), outReader, sink, "run-tokens", outBuf, mockEmitter, 4*1024*1024)
|
||||
if err != nil {
|
||||
t.Fatalf("driveJSONLines: %v", err)
|
||||
}
|
||||
// "one two three" has 3 fields; error event contributes 0.
|
||||
if outputTokens != 3 {
|
||||
t.Fatalf("expected 3 outputTokens, got %d", outputTokens)
|
||||
}
|
||||
}
|
||||
|
|
@ -1,9 +1,7 @@
|
|||
package cli
|
||||
|
||||
import (
|
||||
"bufio"
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
|
|
@ -64,18 +62,9 @@ func (c *CLI) executeCommand(ctx context.Context, spec runtime.ExecutionSpec, pr
|
|||
outputTokens int
|
||||
readErr error
|
||||
)
|
||||
switch profile.OutputFormat {
|
||||
case "stream-json":
|
||||
outputTokens, readErr = emitStreamJSONLines(ctx, stdout, sink, spec.RunID, &outBuf)
|
||||
case "codex-json":
|
||||
outputTokens, readErr = emitCodexJSONLines(ctx, stdout, sink, spec.RunID, &outBuf)
|
||||
case "claude-json":
|
||||
outputTokens, readErr = emitClaudeJSONLines(ctx, stdout, sink, spec.RunID, &outBuf)
|
||||
case "opencode-json":
|
||||
outputTokens, readErr = emitOpencodeJSON(ctx, stdout, sink, spec.RunID, &outBuf)
|
||||
case "cline-json":
|
||||
outputTokens, readErr = emitClineJSON(ctx, stdout, sink, spec.RunID, &outBuf)
|
||||
default:
|
||||
if reg, ok := jsonEmitters[profile.OutputFormat]; ok {
|
||||
outputTokens, readErr = driveJSONLines(ctx, stdout, sink, spec.RunID, &outBuf, reg.emitter, reg.scanBufMax)
|
||||
} else {
|
||||
outputTokens, readErr = emitStdoutChunks(ctx, stdout, sink, spec.RunID, &outBuf)
|
||||
}
|
||||
if readErr != nil {
|
||||
|
|
@ -158,316 +147,6 @@ func emitStdoutChunks(ctx context.Context, stdout io.Reader, sink runtime.EventS
|
|||
}
|
||||
}
|
||||
|
||||
// emitStreamJSONLines reads JSON Lines from stdout (one JSON object per line)
|
||||
// and emits assistant content (delta=true) as RuntimeEvent deltas. Other
|
||||
// event types are written to outBuf so they remain available for diagnostics
|
||||
// but are not pushed to the sink.
|
||||
func emitStreamJSONLines(ctx context.Context, stdout io.Reader, sink runtime.EventSink, runID string, outBuf *strings.Builder) (int, error) {
|
||||
scanner := bufio.NewScanner(stdout)
|
||||
scanner.Buffer(make([]byte, 64*1024), 4*1024*1024)
|
||||
outputTokens := 0
|
||||
|
||||
for scanner.Scan() {
|
||||
line := scanner.Bytes()
|
||||
outBuf.Write(line)
|
||||
outBuf.WriteByte('\n')
|
||||
|
||||
trimmed := strings.TrimSpace(string(line))
|
||||
if trimmed == "" || trimmed[0] != '{' {
|
||||
continue
|
||||
}
|
||||
var ev struct {
|
||||
Type string `json:"type"`
|
||||
Role string `json:"role"`
|
||||
Content string `json:"content"`
|
||||
Delta bool `json:"delta"`
|
||||
Error string `json:"error"`
|
||||
}
|
||||
if err := json.Unmarshal(line, &ev); err != nil {
|
||||
continue
|
||||
}
|
||||
switch ev.Type {
|
||||
case "message":
|
||||
if ev.Role != "assistant" || ev.Content == "" {
|
||||
continue
|
||||
}
|
||||
outputTokens += len(strings.Fields(ev.Content))
|
||||
_ = sink.Emit(ctx, runtime.RuntimeEvent{
|
||||
RunID: runID,
|
||||
Type: runtime.EventTypeDelta,
|
||||
Delta: ev.Content,
|
||||
Timestamp: time.Now(),
|
||||
})
|
||||
case "error":
|
||||
_ = sink.Emit(ctx, runtime.RuntimeEvent{
|
||||
RunID: runID,
|
||||
Type: runtime.EventTypeError,
|
||||
Error: ev.Error,
|
||||
Timestamp: time.Now(),
|
||||
})
|
||||
}
|
||||
}
|
||||
if err := scanner.Err(); err != nil {
|
||||
return outputTokens, err
|
||||
}
|
||||
return outputTokens, nil
|
||||
}
|
||||
|
||||
// emitClaudeJSONLines parses `claude -p --output-format stream-json
|
||||
// --include-partial-messages` JSONL output. Token-level text arrives as
|
||||
// stream_event with event.type == "content_block_delta" and
|
||||
// event.delta.type == "text_delta". Other event types (system, assistant,
|
||||
// result, rate_limit_event, ...) are written to outBuf for diagnostics but
|
||||
// not pushed to the sink to avoid duplicate text.
|
||||
func emitClaudeJSONLines(ctx context.Context, stdout io.Reader, sink runtime.EventSink, runID string, outBuf *strings.Builder) (int, error) {
|
||||
scanner := bufio.NewScanner(stdout)
|
||||
scanner.Buffer(make([]byte, 64*1024), 8*1024*1024)
|
||||
outputTokens := 0
|
||||
|
||||
for scanner.Scan() {
|
||||
line := scanner.Bytes()
|
||||
outBuf.Write(line)
|
||||
outBuf.WriteByte('\n')
|
||||
|
||||
trimmed := strings.TrimSpace(string(line))
|
||||
if trimmed == "" || trimmed[0] != '{' {
|
||||
continue
|
||||
}
|
||||
var ev struct {
|
||||
Type string `json:"type"`
|
||||
Event struct {
|
||||
Type string `json:"type"`
|
||||
Delta struct {
|
||||
Type string `json:"type"`
|
||||
Text string `json:"text"`
|
||||
} `json:"delta"`
|
||||
} `json:"event"`
|
||||
IsError bool `json:"is_error"`
|
||||
Result string `json:"result"`
|
||||
}
|
||||
if err := json.Unmarshal(line, &ev); err != nil {
|
||||
continue
|
||||
}
|
||||
switch ev.Type {
|
||||
case "stream_event":
|
||||
if ev.Event.Type != "content_block_delta" || ev.Event.Delta.Type != "text_delta" || ev.Event.Delta.Text == "" {
|
||||
continue
|
||||
}
|
||||
outputTokens += len(strings.Fields(ev.Event.Delta.Text))
|
||||
_ = sink.Emit(ctx, runtime.RuntimeEvent{
|
||||
RunID: runID,
|
||||
Type: runtime.EventTypeDelta,
|
||||
Delta: ev.Event.Delta.Text,
|
||||
Timestamp: time.Now(),
|
||||
})
|
||||
case "result":
|
||||
if ev.IsError && ev.Result != "" {
|
||||
_ = sink.Emit(ctx, runtime.RuntimeEvent{
|
||||
RunID: runID,
|
||||
Type: runtime.EventTypeError,
|
||||
Error: ev.Result,
|
||||
Timestamp: time.Now(),
|
||||
})
|
||||
}
|
||||
}
|
||||
}
|
||||
if err := scanner.Err(); err != nil {
|
||||
return outputTokens, err
|
||||
}
|
||||
return outputTokens, nil
|
||||
}
|
||||
|
||||
// emitCodexJSONLines parses `codex exec --json` JSONL output. Codex emits
|
||||
// one event per line; assistant text arrives in a single `item.completed`
|
||||
// event with item.type == "agent_message" rather than token-by-token.
|
||||
func emitCodexJSONLines(ctx context.Context, stdout io.Reader, sink runtime.EventSink, runID string, outBuf *strings.Builder) (int, error) {
|
||||
scanner := bufio.NewScanner(stdout)
|
||||
scanner.Buffer(make([]byte, 64*1024), 4*1024*1024)
|
||||
outputTokens := 0
|
||||
|
||||
for scanner.Scan() {
|
||||
line := scanner.Bytes()
|
||||
outBuf.Write(line)
|
||||
outBuf.WriteByte('\n')
|
||||
|
||||
trimmed := strings.TrimSpace(string(line))
|
||||
if trimmed == "" || trimmed[0] != '{' {
|
||||
continue
|
||||
}
|
||||
var ev struct {
|
||||
Type string `json:"type"`
|
||||
Item struct {
|
||||
Type string `json:"type"`
|
||||
Text string `json:"text"`
|
||||
} `json:"item"`
|
||||
Message string `json:"message"`
|
||||
Error struct {
|
||||
Message string `json:"message"`
|
||||
} `json:"error"`
|
||||
}
|
||||
if err := json.Unmarshal(line, &ev); err != nil {
|
||||
continue
|
||||
}
|
||||
switch ev.Type {
|
||||
case "item.completed":
|
||||
if ev.Item.Type != "agent_message" || ev.Item.Text == "" {
|
||||
continue
|
||||
}
|
||||
outputTokens += len(strings.Fields(ev.Item.Text))
|
||||
_ = sink.Emit(ctx, runtime.RuntimeEvent{
|
||||
RunID: runID,
|
||||
Type: runtime.EventTypeDelta,
|
||||
Delta: ev.Item.Text,
|
||||
Timestamp: time.Now(),
|
||||
})
|
||||
case "error":
|
||||
_ = sink.Emit(ctx, runtime.RuntimeEvent{
|
||||
RunID: runID,
|
||||
Type: runtime.EventTypeError,
|
||||
Error: ev.Message,
|
||||
Timestamp: time.Now(),
|
||||
})
|
||||
case "turn.failed":
|
||||
_ = sink.Emit(ctx, runtime.RuntimeEvent{
|
||||
RunID: runID,
|
||||
Type: runtime.EventTypeError,
|
||||
Error: ev.Error.Message,
|
||||
Timestamp: time.Now(),
|
||||
})
|
||||
}
|
||||
}
|
||||
if err := scanner.Err(); err != nil {
|
||||
return outputTokens, err
|
||||
}
|
||||
return outputTokens, nil
|
||||
}
|
||||
|
||||
// emitOpencodeJSON parses `opencode run --format json` output.
|
||||
// Each line is a JSON event: step_start, text (with part.text), step_finish, or error.
|
||||
// Only type=="text" events with a non-empty part.text are emitted as deltas.
|
||||
func emitOpencodeJSON(ctx context.Context, stdout io.Reader, sink runtime.EventSink, runID string, outBuf *strings.Builder) (int, error) {
|
||||
scanner := bufio.NewScanner(stdout)
|
||||
scanner.Buffer(make([]byte, 64*1024), 4*1024*1024)
|
||||
outputTokens := 0
|
||||
|
||||
for scanner.Scan() {
|
||||
line := scanner.Bytes()
|
||||
outBuf.Write(line)
|
||||
outBuf.WriteByte('\n')
|
||||
|
||||
trimmed := strings.TrimSpace(string(line))
|
||||
if trimmed == "" || trimmed[0] != '{' {
|
||||
continue
|
||||
}
|
||||
|
||||
var ev struct {
|
||||
Type string `json:"type"`
|
||||
Error struct {
|
||||
Name string `json:"name"`
|
||||
Data struct {
|
||||
Message string `json:"message"`
|
||||
} `json:"data"`
|
||||
} `json:"error"`
|
||||
Part struct {
|
||||
Type string `json:"type"`
|
||||
Text string `json:"text"`
|
||||
} `json:"part"`
|
||||
}
|
||||
if err := json.Unmarshal(line, &ev); err != nil {
|
||||
continue
|
||||
}
|
||||
|
||||
switch ev.Type {
|
||||
case "text":
|
||||
if ev.Part.Type == "text" && ev.Part.Text != "" {
|
||||
outputTokens += len(strings.Fields(ev.Part.Text))
|
||||
_ = sink.Emit(ctx, runtime.RuntimeEvent{RunID: runID, Type: runtime.EventTypeDelta, Delta: ev.Part.Text, Timestamp: time.Now()})
|
||||
}
|
||||
case "error":
|
||||
msg := ev.Error.Data.Message
|
||||
if msg == "" {
|
||||
msg = ev.Error.Name
|
||||
}
|
||||
if msg != "" {
|
||||
_ = sink.Emit(ctx, runtime.RuntimeEvent{RunID: runID, Type: runtime.EventTypeError, Error: msg, Timestamp: time.Now()})
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
if err := scanner.Err(); err != nil {
|
||||
return outputTokens, err
|
||||
}
|
||||
return outputTokens, nil
|
||||
}
|
||||
|
||||
// emitClineJSON parses `cline --json` headless output. Cline emits lifecycle
|
||||
// events as JSONL; assistant text is carried by say=="text" events.
|
||||
func emitClineJSON(ctx context.Context, stdout io.Reader, sink runtime.EventSink, runID string, outBuf *strings.Builder) (int, error) {
|
||||
scanner := bufio.NewScanner(stdout)
|
||||
scanner.Buffer(make([]byte, 64*1024), 8*1024*1024)
|
||||
outputTokens := 0
|
||||
|
||||
for scanner.Scan() {
|
||||
line := scanner.Bytes()
|
||||
outBuf.Write(line)
|
||||
outBuf.WriteByte('\n')
|
||||
|
||||
trimmed := strings.TrimSpace(string(line))
|
||||
if trimmed == "" || trimmed[0] != '{' {
|
||||
continue
|
||||
}
|
||||
|
||||
var ev struct {
|
||||
Type string `json:"type"`
|
||||
Say string `json:"say"`
|
||||
Ask string `json:"ask"`
|
||||
Text string `json:"text"`
|
||||
Message string `json:"message"`
|
||||
Status string `json:"status"`
|
||||
Error string `json:"error"`
|
||||
}
|
||||
if err := json.Unmarshal(line, &ev); err != nil {
|
||||
continue
|
||||
}
|
||||
|
||||
switch {
|
||||
case ev.Type == "say" && ev.Say == "text" && ev.Text != "":
|
||||
outputTokens += len(strings.Fields(ev.Text))
|
||||
_ = sink.Emit(ctx, runtime.RuntimeEvent{RunID: runID, Type: runtime.EventTypeDelta, Delta: ev.Text, Timestamp: time.Now()})
|
||||
case ev.Type == "say" && ev.Say == "error":
|
||||
msg := ev.Text
|
||||
if msg == "" {
|
||||
msg = ev.Message
|
||||
}
|
||||
if msg != "" {
|
||||
_ = sink.Emit(ctx, runtime.RuntimeEvent{RunID: runID, Type: runtime.EventTypeError, Error: msg, Timestamp: time.Now()})
|
||||
}
|
||||
case ev.Type == "completion" && ev.Status == "error":
|
||||
msg := ev.Error
|
||||
if msg == "" {
|
||||
msg = ev.Message
|
||||
}
|
||||
if msg == "" {
|
||||
msg = "cline task failed"
|
||||
}
|
||||
_ = sink.Emit(ctx, runtime.RuntimeEvent{RunID: runID, Type: runtime.EventTypeError, Error: msg, Timestamp: time.Now()})
|
||||
case ev.Type == "ask" && ev.Ask == "api_req_failed":
|
||||
msg := ev.Text
|
||||
if msg == "" {
|
||||
msg = ev.Message
|
||||
}
|
||||
if msg != "" {
|
||||
_ = sink.Emit(ctx, runtime.RuntimeEvent{RunID: runID, Type: runtime.EventTypeError, Error: msg, Timestamp: time.Now()})
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
if err := scanner.Err(); err != nil {
|
||||
return outputTokens, err
|
||||
}
|
||||
return outputTokens, nil
|
||||
}
|
||||
|
||||
func extractPrompt(input map[string]any) string {
|
||||
if input == nil {
|
||||
return ""
|
||||
|
|
|
|||
Loading…
Reference in a new issue