Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions CONTRIBUTING.md
Original file line number Diff line number Diff line change
Expand Up @@ -153,6 +153,10 @@ description and keep ownership on the side listed here.
conversation has received no new prompt for one hour. A new prompt for the
same `AgentStateKey` renews that idle window. Explicit cancellation, device
shutdown, and daemon shutdown still terminate processes immediately.
- Run terminal-state persistence must use a short independent context. A
dispatch deadline or cancellation may stop connector work, but it must not
prevent the server from recording the resulting completed, failed, or
cancelled state.
- When an engine supports resume, persist the upstream session id through
`agent_engine_sessions` and pass `AgentSessionID` plus `AgentStateKey` over
the daemon protocol. Do not keep resume ids only in adapter memory, files
Expand Down
1 change: 1 addition & 0 deletions apps/parsar-daemon/internal/agent/claudecode/options.go
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,7 @@ func BuildArgs(opts map[string]any, resumeSessionID string) (BuildResult, error)
args := []string{
"--output-format", "stream-json",
"--input-format", "stream-json",
"--include-partial-messages",
"--verbose",
"--permission-prompt-tool", "stdio",
}
Expand Down
3 changes: 3 additions & 0 deletions apps/parsar-daemon/internal/agent/claudecode/options_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,9 @@ func TestBuildArgsBaseHasStreamFlags(t *testing.T) {
if !slices.Contains(res.Args, "--verbose") {
t.Errorf("missing --verbose in %v", res.Args)
}
if !slices.Contains(res.Args, "--include-partial-messages") {
t.Errorf("missing --include-partial-messages in %v", res.Args)
}
}

// IS_SANDBOX=1 must be in every env passthrough. Without it Claude
Expand Down
99 changes: 91 additions & 8 deletions apps/parsar-daemon/internal/agent/claudecode/parser.go
Original file line number Diff line number Diff line change
Expand Up @@ -58,12 +58,13 @@ func defaultAskIDMinter() string {
// translator converts one NDJSON line from claude stdout into zero or
// more proto.Envelope frames. One translator lives per session.
type translator struct {
runID string
pending pendingRecorder
askPending askRecorder
seq atomic.Uint64
mint permIDMinter
askMint askIDMinter
runID string
pending pendingRecorder
askPending askRecorder
seq atomic.Uint64
mint permIDMinter
askMint askIDMinter
partialBlocks map[int]string
}

func newTranslator(runID string, pending pendingRecorder, askPending askRecorder, mint permIDMinter, askMint askIDMinter) *translator {
Expand Down Expand Up @@ -116,6 +117,8 @@ func (t *translator) Translate(line []byte) (translation, error) {
return t.translateSystem(line)
case "assistant":
return t.translateAssistant(line)
case "stream_event":
return t.translateStreamEvent(line)
case "user":
return t.translateUser(line)
case "control_request":
Expand All @@ -129,6 +132,67 @@ func (t *translator) Translate(line []byte) (translation, error) {
}
}

// translateStreamEvent handles the raw Anthropic events emitted by Claude
// Code with --include-partial-messages. Claude still emits a complete
// assistant frame after these events, so translateAssistant suppresses its
// text/thinking copies once a corresponding partial delta has been observed.
func (t *translator) translateStreamEvent(line []byte) (translation, error) {
var msg struct {
Event struct {
Type string `json:"type"`
Index int `json:"index"`
Delta struct {
Type string `json:"type"`
Text string `json:"text"`
Thinking string `json:"thinking"`
} `json:"delta"`
} `json:"event"`
}
if err := json.Unmarshal(line, &msg); err != nil {
return translation{}, fmt.Errorf("claudecode: parse stream_event frame: %w", err)
}
if msg.Event.Type != "content_block_delta" {
return translation{}, nil
}

switch msg.Event.Delta.Type {
case "text_delta":
if msg.Event.Delta.Text == "" {
return translation{}, nil
}
if t.partialBlocks == nil {
t.partialBlocks = make(map[int]string)
}
t.partialBlocks[msg.Event.Index] = "text"
env, err := proto.NewEnvelope(proto.TypeDelta, t.runID, proto.DeltaPayload{
Delta: msg.Event.Delta.Text,
Sequence: t.seq.Add(1),
})
if err != nil {
return translation{}, err
}
return translation{Envelopes: []proto.Envelope{env}}, nil
case "thinking_delta":
if msg.Event.Delta.Thinking == "" {
return translation{}, nil
}
if t.partialBlocks == nil {
t.partialBlocks = make(map[int]string)
}
t.partialBlocks[msg.Event.Index] = "thinking"
env, err := proto.NewEnvelope(proto.TypeThinking, t.runID, proto.ThinkingPayload{
Text: msg.Event.Delta.Thinking,
Sequence: t.seq.Add(1),
})
if err != nil {
return translation{}, err
}
return translation{Envelopes: []proto.Envelope{env}}, nil
default:
return translation{}, nil
}
}

func (t *translator) translateSystem(line []byte) (translation, error) {
var msg struct {
SessionID string `json:"session_id"`
Expand All @@ -150,7 +214,7 @@ func (t *translator) translateAssistant(line []byte) (translation, error) {
}

var envs []proto.Envelope
for _, raw := range msg.Message.Content {
for index, raw := range msg.Message.Content {
var head struct {
Type string `json:"type"`
}
Expand All @@ -159,6 +223,9 @@ func (t *translator) translateAssistant(line []byte) (translation, error) {
}
switch head.Type {
case "text":
if t.partialBlocks[index] == "text" {
continue
}
var item struct {
Text string `json:"text"`
}
Expand All @@ -174,6 +241,9 @@ func (t *translator) translateAssistant(line []byte) (translation, error) {
}
envs = append(envs, env)
case "thinking":
if t.partialBlocks[index] == "thinking" {
continue
}
var item struct {
Thinking string `json:"thinking"`
}
Expand Down Expand Up @@ -220,6 +290,10 @@ func (t *translator) translateAssistant(line []byte) (translation, error) {
envs = append(envs, env)
}
}
// Partial flags apply only to the complete assistant frame that follows
// those stream_event deltas. Reset them so a later assistant turn that is
// delivered without partial frames is not accidentally suppressed.
clear(t.partialBlocks)
return translation{Envelopes: envs}, nil
}

Expand Down Expand Up @@ -370,6 +444,8 @@ type resultUsage struct {
}

func (t *translator) translateResult(line []byte, subtype string) (translation, error) {
defer clear(t.partialBlocks)

var msg struct {
IsError bool `json:"is_error"`
Result string `json:"result"`
Expand Down Expand Up @@ -424,7 +500,14 @@ func (t *translator) translateResult(line []byte, subtype string) (translation,
// else (error_during_execution, error_max_turns, ...) is a failure.
isError := msg.IsError || (subtype != "" && subtype != "success" && strings.HasPrefix(subtype, "error"))
if isError {
errMsg := msg.Error
errMsg := strings.TrimSpace(msg.Error)
// Claude Code sometimes reports provider/API failures with
// subtype="success" and is_error=true, placing the useful error in
// result instead of error. Preserve that message rather than emitting
// the misleading fallback "claude_code: success".
if errMsg == "" && msg.IsError {
errMsg = strings.TrimSpace(msg.Result)
}
if errMsg == "" {
if subtype != "" {
errMsg = "claude_code: " + subtype
Expand Down
109 changes: 109 additions & 0 deletions apps/parsar-daemon/internal/agent/claudecode/parser_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -110,6 +110,100 @@ func TestTranslateAssistantThinkingEmitsThinking(t *testing.T) {
}
}

func TestTranslatePartialMessagesStreamIncrementallyWithoutAssistantDuplicates(t *testing.T) {
tr := claudecode.NewTranslatorForTest("run_partial", nil, counterMinter())
frames := [][]byte{
[]byte(`{"type":"stream_event","event":{"type":"content_block_delta","index":0,"delta":{"type":"text_delta","text":"hello "}}}`),
[]byte(`{"type":"stream_event","event":{"type":"content_block_delta","index":0,"delta":{"type":"text_delta","text":"world"}}}`),
[]byte(`{"type":"stream_event","event":{"type":"content_block_delta","index":1,"delta":{"type":"thinking_delta","thinking":"checking"}}}`),
}

var got []proto.Envelope
for _, frame := range frames {
out, err := tr.Translate(frame)
if err != nil {
t.Fatalf("Translate partial frame: %v", err)
}
got = append(got, out.Envelopes...)
}
if len(got) != 3 {
t.Fatalf("want 3 partial envelopes, got %d", len(got))
}
if got[0].Type != proto.TypeDelta || got[1].Type != proto.TypeDelta || got[2].Type != proto.TypeThinking {
t.Fatalf("partial envelope types = %q, %q, %q", got[0].Type, got[1].Type, got[2].Type)
}

full, err := tr.Translate([]byte(`{"type":"assistant","message":{"role":"assistant","content":[{"type":"text","text":"hello world"},{"type":"thinking","thinking":"checking"}]}}`))
if err != nil {
t.Fatalf("Translate complete assistant frame: %v", err)
}
if len(full.Envelopes) != 0 {
t.Fatalf("complete assistant frame duplicated partial content: %#v", full.Envelopes)
}

next, err := tr.Translate([]byte(`{"type":"assistant","message":{"role":"assistant","content":[{"type":"text","text":"next turn without partial frames"}]}}`))
if err != nil {
t.Fatalf("Translate next complete assistant frame: %v", err)
}
if len(next.Envelopes) != 1 || next.Envelopes[0].Type != proto.TypeDelta {
t.Fatalf("next assistant frame was suppressed by stale partial state: %#v", next.Envelopes)
}
}

func TestTranslateStreamEventIgnoresNonTextDeltas(t *testing.T) {
tr := claudecode.NewTranslatorForTest("run_partial", nil, counterMinter())
out, err := tr.Translate([]byte(`{"type":"stream_event","event":{"type":"content_block_delta","index":0,"delta":{"type":"input_json_delta","partial_json":"{\"path\":"}}}`))
if err != nil {
t.Fatalf("Translate: %v", err)
}
if len(out.Envelopes) != 0 {
t.Fatalf("input JSON delta should not emit an envelope: %#v", out.Envelopes)
}
}

func TestTranslatePartialMessagesSuppressOnlyMatchingContentBlock(t *testing.T) {
tr := claudecode.NewTranslatorForTest("run_partial_blocks", nil, counterMinter())
_, err := tr.Translate([]byte(`{"type":"stream_event","event":{"type":"content_block_delta","index":0,"delta":{"type":"text_delta","text":"streamed"}}}`))
if err != nil {
t.Fatalf("Translate partial frame: %v", err)
}

out, err := tr.Translate([]byte(`{"type":"assistant","message":{"role":"assistant","content":[{"type":"text","text":"streamed"},{"type":"text","text":"complete-only"}]}}`))
if err != nil {
t.Fatalf("Translate complete frame: %v", err)
}
if len(out.Envelopes) != 1 || out.Envelopes[0].Type != proto.TypeDelta {
t.Fatalf("complete-only block should be preserved: %#v", out.Envelopes)
}
got := mustDecode[struct {
Delta string `json:"delta"`
}](t, out.Envelopes[0].Payload)
if got.Delta != "complete-only" {
t.Fatalf("delta = %q, want complete-only", got.Delta)
}
}

func TestTranslateResultClearsPartialBlocksBeforeNextTurn(t *testing.T) {
tr := claudecode.NewTranslatorForTest("run_partial_error", nil, counterMinter())
_, err := tr.Translate([]byte(`{"type":"stream_event","event":{"type":"content_block_delta","index":0,"delta":{"type":"text_delta","text":"partial before failure"}}}`))
if err != nil {
t.Fatalf("Translate partial frame: %v", err)
}

_, err = tr.Translate([]byte(`{"type":"result","subtype":"error_during_execution","is_error":true,"error":"provider failed"}`))
if err != nil {
t.Fatalf("Translate error result: %v", err)
}

next, err := tr.Translate([]byte(`{"type":"assistant","message":{"role":"assistant","content":[{"type":"text","text":"next turn"}]}}`))
if err != nil {
t.Fatalf("Translate next assistant frame: %v", err)
}
if len(next.Envelopes) != 1 || next.Envelopes[0].Type != proto.TypeDelta {
t.Fatalf("next assistant frame was suppressed by result state: %#v", next.Envelopes)
}
}

func TestTranslateAssistantToolUseEmitsBeforeStage(t *testing.T) {
tr := claudecode.NewTranslatorForTest("run_t", nil, counterMinter())
line := []byte(`{"type":"assistant","message":{"role":"assistant","content":[
Expand Down Expand Up @@ -389,6 +483,21 @@ func TestTranslateResultErrorWithoutMessageFallsBackToSubtype(t *testing.T) {
}
}

func TestTranslateResultIsErrorSuccessSubtypeUsesResultMessage(t *testing.T) {
tr := claudecode.NewTranslatorForTest("run_99", nil, counterMinter())
line := []byte(`{"type":"result","subtype":"success","is_error":true,"result":"API Error: 400 content rejected"}`)
out, err := tr.Translate(line)
if err != nil {
t.Fatalf("Translate: %v", err)
}
got := mustDecode[struct {
Error string `json:"error"`
}](t, out.Envelopes[0].Payload)
if got.Error != "API Error: 400 content rejected" {
t.Errorf("error text = %q", got.Error)
}
}

func TestTranslateUnknownTypeIsNoOp(t *testing.T) {
tr := claudecode.NewTranslatorForTest("run_99", nil, counterMinter())
out, err := tr.Translate([]byte(`{"type":"some_future_thing","foo":1}`))
Expand Down
Loading
Loading