- 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
320 lines
No EOL
7.4 KiB
Go
320 lines
No EOL
7.4 KiB
Go
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
|
|
} |