mirror of
https://github.com/multica-ai/multica.git
synced 2026-07-26 12:35:35 +02:00
MUL-5231 fix(agent): parse Cursor top-level thinking and tool_call events (#5846)
Cursor stream-json emits reasoning and tool calls as top-level `thinking` / `tool_call` events; the parser only looked for them inside assistant messages, so transcripts showed a single step. Match the top-level events, taking the tool name from the nested `<name>ToolCall` key and normalizing the packed call_id. Subtypes are matched explicitly — only `started` opens a tool, only `completed` closes it, only `delta` carries reasoning — so an unknown/missing subtype is ignored rather than synthesizing a fake result (which would decrement the daemon in-flight tool count early and misfire the watchdog) or polluting reasoning. Covered by a recorded-stream fixture test, an unknown-subtype regression test, and an opt-in real-CLI smoke test.
This commit is contained in:
@@ -9,6 +9,7 @@ import (
|
||||
"log/slog"
|
||||
"os/exec"
|
||||
"regexp"
|
||||
"sort"
|
||||
"strings"
|
||||
"sync"
|
||||
"time"
|
||||
@@ -114,6 +115,11 @@ func (b *cursorBackend) Execute(ctx context.Context, prompt string, opts ExecOpt
|
||||
invalidEventCount := 0
|
||||
assistantEventCount := 0
|
||||
toolUseCount := 0
|
||||
// unknownSubtypeCount tracks thinking/tool_call events whose subtype we
|
||||
// don't recognize. They are ignored (never synthesized into a message)
|
||||
// and surfaced once as a bounded, content-free diagnostic so an upstream
|
||||
// protocol addition is visible instead of silent.
|
||||
unknownSubtypeCount := 0
|
||||
lastEventType := "none"
|
||||
// stepUsage accumulates per-step token counts from "step_finish" events.
|
||||
// resultUsage holds authoritative session totals from "result" events.
|
||||
@@ -122,6 +128,7 @@ func (b *cursorBackend) Execute(ctx context.Context, prompt string, opts ExecOpt
|
||||
stepUsage := make(map[string]TokenUsage)
|
||||
resultUsage := make(map[string]TokenUsage)
|
||||
hasResultUsage := false
|
||||
var thinking cursorThinkingStream
|
||||
|
||||
scanner := bufio.NewScanner(stdout)
|
||||
scanner.Buffer(make([]byte, 0, 1024*1024), 10*1024*1024)
|
||||
@@ -162,6 +169,55 @@ func (b *cursorBackend) Execute(ctx context.Context, prompt string, opts ExecOpt
|
||||
assistantEventCount++
|
||||
b.handleCursorAssistant(&evt, msgCh, &output)
|
||||
|
||||
case "thinking":
|
||||
// Reasoning is a top-level event streamed as deltas, not a
|
||||
// content block inside assistant messages. Match subtypes
|
||||
// explicitly: only `delta` carries reasoning content (forwarded
|
||||
// as it lands so the daemon's 500ms flush shows it mid-run),
|
||||
// `completed` closes the block. An unknown subtype is NOT folded
|
||||
// into reasoning — silently absorbing upstream additions is the
|
||||
// exact failure mode this fix exists to prevent (MUL-5231).
|
||||
switch evt.Subtype {
|
||||
case "delta":
|
||||
if content := thinking.delta(evt.Text); content != "" {
|
||||
trySend(msgCh, Message{Type: MessageThinking, Content: content})
|
||||
}
|
||||
case "completed":
|
||||
thinking.complete()
|
||||
default:
|
||||
unknownSubtypeCount++
|
||||
}
|
||||
|
||||
case "tool_call":
|
||||
// Only the two subtypes that define a call's boundaries drive the
|
||||
// transcript: `started` opens it, `completed` closes it. A
|
||||
// non-terminal (e.g. a future `progress`) or missing subtype must
|
||||
// NOT synthesize a result — that would decrement the daemon's
|
||||
// in-flight tool count early and drop a still-running long tool
|
||||
// from the tool watchdog onto the shorter idle watchdog, which
|
||||
// can force-stop it as falsely stuck (MUL-5231 review).
|
||||
switch evt.Subtype {
|
||||
case "started":
|
||||
call := parseCursorToolCall(&evt)
|
||||
toolUseCount++
|
||||
trySend(msgCh, Message{
|
||||
Type: MessageToolUse,
|
||||
Tool: call.Name,
|
||||
CallID: call.CallID,
|
||||
Input: call.Input,
|
||||
})
|
||||
case "completed":
|
||||
call := parseCursorToolCall(&evt)
|
||||
trySend(msgCh, Message{
|
||||
Type: MessageToolResult,
|
||||
Tool: call.Name,
|
||||
CallID: call.CallID,
|
||||
Output: call.Result,
|
||||
})
|
||||
default:
|
||||
unknownSubtypeCount++
|
||||
}
|
||||
|
||||
case "tool_use":
|
||||
toolUseCount++
|
||||
var params map[string]any
|
||||
@@ -323,6 +379,13 @@ func (b *cursorBackend) Execute(ctx context.Context, prompt string, opts ExecOpt
|
||||
lastEventType: lastEventType,
|
||||
})
|
||||
|
||||
if unknownSubtypeCount > 0 {
|
||||
// Content-free and emitted once per run: signals that the CLI sent a
|
||||
// thinking/tool_call subtype we chose to ignore, so a protocol
|
||||
// addition is diagnosable without the parser having guessed at it.
|
||||
b.cfg.Logger.Warn("cursor-agent ignored unknown event subtypes", "count", unknownSubtypeCount)
|
||||
}
|
||||
|
||||
b.cfg.Logger.Info("cursor-agent finished", "pid", cmd.Process.Pid, "status", finalStatus, "duration", duration.Round(time.Millisecond).String())
|
||||
|
||||
finalOutput := output.String()
|
||||
@@ -420,6 +483,127 @@ func (b *cursorBackend) handleCursorAssistant(evt *cursorStreamEvent, ch chan<-
|
||||
}
|
||||
}
|
||||
|
||||
// cursorThinkingStream turns the CLI's `thinking` event sequence into the
|
||||
// content we forward as MessageThinking. Cursor streams a reasoning block as
|
||||
// `subtype:"delta"` events carrying a text fragment each, terminated by a
|
||||
// `subtype:"completed"` event that carries no text of its own.
|
||||
//
|
||||
// Deltas are forwarded immediately (the daemon concatenates and flushes them on
|
||||
// a ticker, so reasoning shows up while the task runs) and blocks are separated
|
||||
// by a blank line, because consecutive blocks would otherwise be glued together
|
||||
// into one paragraph by that same concatenation.
|
||||
type cursorThinkingStream struct {
|
||||
blockOpen bool
|
||||
anySent bool
|
||||
}
|
||||
|
||||
// delta returns the content to forward for one `delta` event, or "" when the
|
||||
// fragment is empty. The first fragment of a new block is prefixed with a blank
|
||||
// line so the daemon's concatenation keeps blocks visually separated.
|
||||
func (t *cursorThinkingStream) delta(text string) string {
|
||||
if text == "" {
|
||||
return ""
|
||||
}
|
||||
if !t.blockOpen && t.anySent {
|
||||
text = "\n\n" + text
|
||||
}
|
||||
t.blockOpen = true
|
||||
t.anySent = true
|
||||
return text
|
||||
}
|
||||
|
||||
// complete closes the current reasoning block. The terminal event carries no
|
||||
// content of its own in the observed protocol, so nothing is forwarded; the
|
||||
// next block's first delta gets the separating blank line.
|
||||
func (t *cursorThinkingStream) complete() {
|
||||
t.blockOpen = false
|
||||
}
|
||||
|
||||
// cursorToolCall is a top-level `tool_call` event projected onto the fields the
|
||||
// daemon transcript needs.
|
||||
type cursorToolCall struct {
|
||||
Name string
|
||||
CallID string
|
||||
Input map[string]any
|
||||
Result string
|
||||
}
|
||||
|
||||
// cursorToolCallKeySuffix is how Cursor names the per-tool payload: the tool is
|
||||
// the key of a nested object rather than a `name` field, so `readToolCall`
|
||||
// identifies a read and holds that call's `args` and (once completed) `result`.
|
||||
const cursorToolCallKeySuffix = "ToolCall"
|
||||
|
||||
// parseCursorToolCall extracts tool identity, arguments and result from a
|
||||
// `tool_call` event:
|
||||
//
|
||||
// {"type":"tool_call","subtype":"started","call_id":"call-…",
|
||||
// "tool_call":{"readToolCall":{"args":{"path":"/x/a.txt"}},"toolCallId":"call-…"}}
|
||||
//
|
||||
// A payload we cannot decode still yields the call ID, so started/completed
|
||||
// events stay paired and the transcript keeps a row for the call.
|
||||
func parseCursorToolCall(evt *cursorStreamEvent) cursorToolCall {
|
||||
call := cursorToolCall{CallID: cursorCallID(evt.CallID)}
|
||||
|
||||
var envelope map[string]json.RawMessage
|
||||
if len(evt.ToolCall) == 0 || json.Unmarshal(evt.ToolCall, &envelope) != nil {
|
||||
return call
|
||||
}
|
||||
if call.CallID == "" {
|
||||
var nestedID string
|
||||
if err := json.Unmarshal(envelope["toolCallId"], &nestedID); err == nil {
|
||||
call.CallID = cursorCallID(nestedID)
|
||||
}
|
||||
}
|
||||
|
||||
key := cursorToolPayloadKey(envelope)
|
||||
if key == "" {
|
||||
return call
|
||||
}
|
||||
call.Name = strings.TrimSuffix(key, cursorToolCallKeySuffix)
|
||||
|
||||
var payload struct {
|
||||
Args map[string]any `json:"args"`
|
||||
Result json.RawMessage `json:"result"`
|
||||
}
|
||||
if err := json.Unmarshal(envelope[key], &payload); err != nil {
|
||||
return call
|
||||
}
|
||||
call.Input = payload.Args
|
||||
if len(payload.Result) > 0 {
|
||||
call.Result = string(payload.Result)
|
||||
}
|
||||
return call
|
||||
}
|
||||
|
||||
// cursorToolPayloadKey picks the `<name>ToolCall` key of a tool_call envelope.
|
||||
// Observed payloads carry exactly one; the sort keeps the choice deterministic
|
||||
// if a future CLI ever emits more than one.
|
||||
func cursorToolPayloadKey(envelope map[string]json.RawMessage) string {
|
||||
var keys []string
|
||||
for key := range envelope {
|
||||
if len(key) > len(cursorToolCallKeySuffix) && strings.HasSuffix(key, cursorToolCallKeySuffix) {
|
||||
keys = append(keys, key)
|
||||
}
|
||||
}
|
||||
if len(keys) == 0 {
|
||||
return ""
|
||||
}
|
||||
sort.Strings(keys)
|
||||
return keys[0]
|
||||
}
|
||||
|
||||
// cursorCallID normalizes the CLI's call identifier. Cursor packs two ids into
|
||||
// one newline-separated string ("call-…\nfc_…"); the first line alone is
|
||||
// unique, and keeping IDs single-line stops the embedded newline from breaking
|
||||
// up daemon log lines.
|
||||
func cursorCallID(raw string) string {
|
||||
id := strings.TrimSpace(raw)
|
||||
if idx := strings.IndexByte(id, '\n'); idx >= 0 {
|
||||
id = id[:idx]
|
||||
}
|
||||
return strings.TrimSpace(id)
|
||||
}
|
||||
|
||||
func cursorUsageModel(evtModel, configuredModel string) string {
|
||||
if model := strings.TrimSpace(evtModel); model != "" {
|
||||
return model
|
||||
@@ -466,6 +650,13 @@ type cursorStreamEvent struct {
|
||||
// assistant fields
|
||||
Message json.RawMessage `json:"message,omitempty"`
|
||||
|
||||
// thinking fields
|
||||
Text string `json:"text,omitempty"`
|
||||
|
||||
// tool_call fields (current CLI): the tool is a nested key under tool_call
|
||||
ToolCall json.RawMessage `json:"tool_call,omitempty"`
|
||||
CallID string `json:"call_id,omitempty"`
|
||||
|
||||
// tool_use fields
|
||||
ToolName string `json:"tool_name,omitempty"`
|
||||
ToolID string `json:"tool_id,omitempty"`
|
||||
|
||||
99
server/pkg/agent/cursor_integration_test.go
Normal file
99
server/pkg/agent/cursor_integration_test.go
Normal file
@@ -0,0 +1,99 @@
|
||||
//go:build agentintegration
|
||||
|
||||
package agent
|
||||
|
||||
import (
|
||||
"context"
|
||||
"log/slog"
|
||||
"os/exec"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
// TestCursorRealStreamObservability drives the real cursor-agent binary with a
|
||||
// task that must use tools, and asserts the run is observable: reasoning and
|
||||
// each tool call/result reach the message stream, not just the final answer.
|
||||
//
|
||||
// The fixture test pins the parser against a recorded stream; this one is what
|
||||
// catches the upstream protocol moving on (MUL-5231).
|
||||
func TestCursorRealStreamObservability(t *testing.T) {
|
||||
requireRealAgentSmoke(t)
|
||||
if testing.Short() {
|
||||
t.Skip("skipping real-binary smoke test in -short mode")
|
||||
}
|
||||
path, err := exec.LookPath("cursor-agent")
|
||||
if err != nil {
|
||||
t.Skip("cursor-agent not on PATH; skipping real-binary smoke test")
|
||||
}
|
||||
if version, err := exec.Command(path, "--version").CombinedOutput(); err == nil {
|
||||
t.Logf("cursor-agent CLI version: %s", strings.TrimSpace(string(version)))
|
||||
} else {
|
||||
t.Logf("cursor-agent CLI version unavailable: %v (%s)", err, strings.TrimSpace(string(version)))
|
||||
}
|
||||
|
||||
backend, err := New("cursor", Config{ExecutablePath: path, Logger: slog.Default()})
|
||||
if err != nil {
|
||||
t.Fatalf("new cursor backend: %v", err)
|
||||
}
|
||||
ctx, cancel := context.WithTimeout(context.Background(), 180*time.Second)
|
||||
defer cancel()
|
||||
|
||||
session, err := backend.Execute(ctx,
|
||||
"Create a file named ping.txt in this workspace containing exactly the word pong, then read it back and tell me what it says.",
|
||||
ExecOptions{Cwd: t.TempDir(), Timeout: 170 * time.Second})
|
||||
if err != nil {
|
||||
t.Fatalf("execute: %v", err)
|
||||
}
|
||||
|
||||
var thinking strings.Builder
|
||||
var toolUses, toolResults []Message
|
||||
done := make(chan struct{})
|
||||
go func() {
|
||||
defer close(done)
|
||||
for msg := range session.Messages {
|
||||
switch msg.Type {
|
||||
case MessageThinking:
|
||||
thinking.WriteString(msg.Content)
|
||||
case MessageToolUse:
|
||||
toolUses = append(toolUses, msg)
|
||||
case MessageToolResult:
|
||||
toolResults = append(toolResults, msg)
|
||||
}
|
||||
}
|
||||
}()
|
||||
|
||||
result := <-session.Result
|
||||
<-done
|
||||
|
||||
if result.Status != "completed" {
|
||||
t.Fatalf("real cursor run did not complete: status=%q error=%q", result.Status, result.Error)
|
||||
}
|
||||
if thinking.Len() == 0 {
|
||||
t.Error("no reasoning observed from the real cursor stream")
|
||||
}
|
||||
if len(toolUses) == 0 {
|
||||
t.Fatal("no tool calls observed from the real cursor stream")
|
||||
}
|
||||
if len(toolResults) != len(toolUses) {
|
||||
t.Errorf("tool results = %d, tool uses = %d; every call must report a result", len(toolResults), len(toolUses))
|
||||
}
|
||||
for _, msg := range toolUses {
|
||||
if msg.Tool == "" {
|
||||
t.Errorf("tool call without a name: %+v", msg)
|
||||
}
|
||||
if strings.ContainsAny(msg.CallID, "\r\n") {
|
||||
t.Errorf("tool call id spans multiple lines: %q", msg.CallID)
|
||||
}
|
||||
}
|
||||
t.Logf("real cursor smoke OK: thinking=%d bytes, tools=%v, output=%q",
|
||||
thinking.Len(), toolNames(toolUses), result.Output)
|
||||
}
|
||||
|
||||
func toolNames(msgs []Message) []string {
|
||||
names := make([]string, 0, len(msgs))
|
||||
for _, msg := range msgs {
|
||||
names = append(names, msg.Tool)
|
||||
}
|
||||
return names
|
||||
}
|
||||
251
server/pkg/agent/cursor_stream_fixture_unix_test.go
Normal file
251
server/pkg/agent/cursor_stream_fixture_unix_test.go
Normal file
@@ -0,0 +1,251 @@
|
||||
//go:build unix
|
||||
|
||||
package agent
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"log/slog"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"testing"
|
||||
"time"
|
||||
)
|
||||
|
||||
// TestCursorExecuteParsesRecordedStream replays a stream captured from a real
|
||||
// `cursor-agent -p --output-format stream-json --yolo` run (2026.07.20-8cc9c0b,
|
||||
// the exact invocation the daemon uses) and pins every observable step of it.
|
||||
//
|
||||
// The regression it guards: reasoning and tool calls arrive as TOP-LEVEL
|
||||
// `thinking` / `tool_call` events, not as content blocks inside assistant
|
||||
// messages. When the parser only understood the assistant-block shape it
|
||||
// dropped all of them silently and the task transcript showed a single step
|
||||
// (MUL-5231). A fixture test is what makes the next upstream event rename loud
|
||||
// instead of silent.
|
||||
func TestCursorExecuteParsesRecordedStream(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
fixture, err := filepath.Abs(filepath.Join("testdata", "cursor-agent-2026.07.20-stream-json.jsonl"))
|
||||
if err != nil {
|
||||
t.Fatalf("resolve fixture: %v", err)
|
||||
}
|
||||
script := fmt.Sprintf("#!/bin/sh\nexec cat %q\n", fixture)
|
||||
|
||||
fakePath := filepath.Join(t.TempDir(), "cursor-agent")
|
||||
writeTestExecutable(t, fakePath, []byte(script))
|
||||
|
||||
backend, err := New("cursor", Config{ExecutablePath: fakePath, Logger: slog.Default()})
|
||||
if err != nil {
|
||||
t.Fatalf("New(cursor): %v", err)
|
||||
}
|
||||
session, err := backend.Execute(t.Context(), "hello", ExecOptions{Timeout: 30 * time.Second})
|
||||
if err != nil {
|
||||
t.Fatalf("Execute: %v", err)
|
||||
}
|
||||
|
||||
var messages []Message
|
||||
done := make(chan struct{})
|
||||
go func() {
|
||||
defer close(done)
|
||||
for msg := range session.Messages {
|
||||
messages = append(messages, msg)
|
||||
}
|
||||
}()
|
||||
result := <-session.Result
|
||||
<-done
|
||||
|
||||
if result.Status != "completed" {
|
||||
t.Fatalf("status = %q, want completed; error=%q", result.Status, result.Error)
|
||||
}
|
||||
|
||||
// Reasoning: three blocks streamed as nine deltas, kept in order and
|
||||
// separated by a blank line so the daemon's concatenation stays readable.
|
||||
const wantThinking = "Reading notes.txt.\n\nRunning wc -l on notes.txt. Writing the line count to count.txt." +
|
||||
"\n\nnotes.txt contains 3 lines. I will write 3 into count.txt." +
|
||||
"\n\nFinished. The line count was written to count.txt."
|
||||
var gotThinking strings.Builder
|
||||
for _, msg := range messages {
|
||||
if msg.Type == MessageThinking {
|
||||
gotThinking.WriteString(msg.Content)
|
||||
}
|
||||
}
|
||||
if gotThinking.String() != wantThinking {
|
||||
t.Errorf("thinking =\n%q\nwant\n%q", gotThinking.String(), wantThinking)
|
||||
}
|
||||
|
||||
// Tool calls: name comes from the nested `<name>ToolCall` key, arguments
|
||||
// from that key's `args`, and the call id is normalized to a single line.
|
||||
wantUses := []struct {
|
||||
tool string
|
||||
callID string
|
||||
inputKey string
|
||||
inputVal string
|
||||
}{
|
||||
{tool: "read", callID: "call-bb11656a-e59e-4356-9866-5b206aedb390-0", inputKey: "path", inputVal: "/tmp/curcap/notes.txt"},
|
||||
{tool: "shell", callID: "call-bb11656a-e59e-4356-9866-5b206aedb390-1", inputKey: "command", inputVal: "wc -l notes.txt"},
|
||||
{tool: "edit", callID: "call-c52c0cd6-81ad-4a87-94c5-f5b0119f3ed4-2", inputKey: "path", inputVal: "/tmp/curcap/count.txt"},
|
||||
}
|
||||
var uses []Message
|
||||
for _, msg := range messages {
|
||||
if msg.Type == MessageToolUse {
|
||||
uses = append(uses, msg)
|
||||
}
|
||||
}
|
||||
if len(uses) != len(wantUses) {
|
||||
t.Fatalf("tool_use count = %d, want %d (%+v)", len(uses), len(wantUses), uses)
|
||||
}
|
||||
for i, want := range wantUses {
|
||||
got := uses[i]
|
||||
if got.Tool != want.tool {
|
||||
t.Errorf("tool_use[%d].Tool = %q, want %q", i, got.Tool, want.tool)
|
||||
}
|
||||
if got.CallID != want.callID {
|
||||
t.Errorf("tool_use[%d].CallID = %q, want %q", i, got.CallID, want.callID)
|
||||
}
|
||||
if fmt.Sprint(got.Input[want.inputKey]) != want.inputVal {
|
||||
t.Errorf("tool_use[%d].Input[%q] = %v, want %q", i, want.inputKey, got.Input[want.inputKey], want.inputVal)
|
||||
}
|
||||
}
|
||||
|
||||
// Results pair with their call by id, carry the tool name directly, and
|
||||
// keep the tool's own result payload as output.
|
||||
wantResults := []struct {
|
||||
tool string
|
||||
callID string
|
||||
contains string
|
||||
}{
|
||||
{tool: "read", callID: "call-bb11656a-e59e-4356-9866-5b206aedb390-0", contains: "alpha"},
|
||||
{tool: "shell", callID: "call-bb11656a-e59e-4356-9866-5b206aedb390-1", contains: "3 notes.txt"},
|
||||
{tool: "edit", callID: "call-c52c0cd6-81ad-4a87-94c5-f5b0119f3ed4-2", contains: "Wrote contents to /tmp/curcap/count.txt"},
|
||||
}
|
||||
var results []Message
|
||||
for _, msg := range messages {
|
||||
if msg.Type == MessageToolResult {
|
||||
results = append(results, msg)
|
||||
}
|
||||
}
|
||||
if len(results) != len(wantResults) {
|
||||
t.Fatalf("tool_result count = %d, want %d (%+v)", len(results), len(wantResults), results)
|
||||
}
|
||||
for i, want := range wantResults {
|
||||
got := results[i]
|
||||
if got.Tool != want.tool {
|
||||
t.Errorf("tool_result[%d].Tool = %q, want %q", i, got.Tool, want.tool)
|
||||
}
|
||||
if got.CallID != want.callID {
|
||||
t.Errorf("tool_result[%d].CallID = %q, want %q", i, got.CallID, want.callID)
|
||||
}
|
||||
if !strings.Contains(got.Output, want.contains) {
|
||||
t.Errorf("tool_result[%d].Output = %q, want it to contain %q", i, got.Output, want.contains)
|
||||
}
|
||||
}
|
||||
|
||||
// A result must never precede its own tool_use, otherwise the daemon's
|
||||
// in-flight tool counter goes negative and re-arms the watchdog early.
|
||||
seen := map[string]bool{}
|
||||
for _, msg := range messages {
|
||||
switch msg.Type {
|
||||
case MessageToolUse:
|
||||
seen[msg.CallID] = true
|
||||
case MessageToolResult:
|
||||
if !seen[msg.CallID] {
|
||||
t.Errorf("tool_result for %q arrived before its tool_use", msg.CallID)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
const wantText = "I'll read `notes.txt`, run `wc -l`, then write the line count to `count.txt`." +
|
||||
"`notes.txt` has 3 lines (`alpha`, `beta`, `gamma`). Wrote `3` to `count.txt`."
|
||||
if result.Output != wantText {
|
||||
t.Errorf("result output = %q, want %q", result.Output, wantText)
|
||||
}
|
||||
if result.SessionID == "" {
|
||||
t.Error("session id not captured from the recorded stream")
|
||||
}
|
||||
}
|
||||
|
||||
// TestCursorExecuteIgnoresUnknownSubtypes pins the fail-safe: a `tool_call` or
|
||||
// `thinking` event whose subtype we don't recognize (a future non-terminal
|
||||
// event, or one arriving without a subtype) must be dropped, not coerced into a
|
||||
// message. Synthesizing a tool_result from a non-terminal subtype would
|
||||
// decrement the daemon's in-flight tool count early and drop a still-running
|
||||
// long tool onto the shorter idle watchdog; folding unknown text into reasoning
|
||||
// corrupts the transcript. Both are the silent-misparse regression this parser
|
||||
// exists to prevent (MUL-5231 review).
|
||||
func TestCursorExecuteIgnoresUnknownSubtypes(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
// One real tool call (started + completed) bracketing an unknown
|
||||
// `progress` subtype and a subtype-less tool_call; one real thinking delta
|
||||
// alongside an unknown-subtype thinking event that carries text.
|
||||
lines := []string{
|
||||
`{"type":"system","subtype":"init","session_id":"sess-unknown"}`,
|
||||
`{"type":"thinking","subtype":"delta","text":"real reasoning"}`,
|
||||
`{"type":"thinking","subtype":"progress","text":"NOT reasoning"}`,
|
||||
`{"type":"thinking","text":"also NOT reasoning"}`,
|
||||
`{"type":"tool_call","subtype":"started","call_id":"call-x","tool_call":{"readToolCall":{"args":{"path":"/x"}},"toolCallId":"call-x"}}`,
|
||||
`{"type":"tool_call","subtype":"progress","call_id":"call-x","tool_call":{"readToolCall":{"args":{"path":"/x"}},"toolCallId":"call-x"}}`,
|
||||
`{"type":"tool_call","call_id":"call-x","tool_call":{"readToolCall":{"args":{"path":"/x"}},"toolCallId":"call-x"}}`,
|
||||
`{"type":"tool_call","subtype":"completed","call_id":"call-x","tool_call":{"readToolCall":{"args":{"path":"/x"},"result":{"success":{}}},"toolCallId":"call-x"}}`,
|
||||
`{"type":"result","subtype":"success","is_error":false,"result":"done","session_id":"sess-unknown"}`,
|
||||
}
|
||||
script := "#!/bin/sh\n"
|
||||
for _, line := range lines {
|
||||
script += fmt.Sprintf("printf '%%s\\n' %s\n", shellSingleQuote(line))
|
||||
}
|
||||
|
||||
fakePath := filepath.Join(t.TempDir(), "cursor-agent")
|
||||
writeTestExecutable(t, fakePath, []byte(script))
|
||||
backend, err := New("cursor", Config{ExecutablePath: fakePath, Logger: slog.Default()})
|
||||
if err != nil {
|
||||
t.Fatalf("New(cursor): %v", err)
|
||||
}
|
||||
session, err := backend.Execute(t.Context(), "hello", ExecOptions{Timeout: 30 * time.Second})
|
||||
if err != nil {
|
||||
t.Fatalf("Execute: %v", err)
|
||||
}
|
||||
|
||||
var messages []Message
|
||||
done := make(chan struct{})
|
||||
go func() {
|
||||
defer close(done)
|
||||
for msg := range session.Messages {
|
||||
messages = append(messages, msg)
|
||||
}
|
||||
}()
|
||||
result := <-session.Result
|
||||
<-done
|
||||
|
||||
if result.Status != "completed" {
|
||||
t.Fatalf("status = %q, want completed; error=%q", result.Status, result.Error)
|
||||
}
|
||||
|
||||
var thinking strings.Builder
|
||||
var toolUses, toolResults int
|
||||
for _, msg := range messages {
|
||||
switch msg.Type {
|
||||
case MessageThinking:
|
||||
thinking.WriteString(msg.Content)
|
||||
case MessageToolUse:
|
||||
toolUses++
|
||||
case MessageToolResult:
|
||||
toolResults++
|
||||
}
|
||||
}
|
||||
if thinking.String() != "real reasoning" {
|
||||
t.Errorf("thinking = %q, want only the delta content %q", thinking.String(), "real reasoning")
|
||||
}
|
||||
if toolUses != 1 {
|
||||
t.Errorf("tool_use count = %d, want 1 (only the started event)", toolUses)
|
||||
}
|
||||
// Exactly one result: the `completed` event. The `progress` and
|
||||
// subtype-less events must not each add a spurious result.
|
||||
if toolResults != 1 {
|
||||
t.Errorf("tool_result count = %d, want 1 (only the completed event)", toolResults)
|
||||
}
|
||||
}
|
||||
|
||||
// shellSingleQuote wraps s in single quotes for POSIX sh, escaping any embedded
|
||||
// single quote, so a JSON line reaches printf verbatim.
|
||||
func shellSingleQuote(s string) string {
|
||||
return "'" + strings.ReplaceAll(s, "'", `'\''`) + "'"
|
||||
}
|
||||
@@ -710,3 +710,132 @@ func TestCursorStreamEventUnmarshalLegacyUsage(t *testing.T) {
|
||||
t.Fatalf("accumulated usage = %+v, want input=800 output=400 cache_read=200 cache_write=100", u)
|
||||
}
|
||||
}
|
||||
|
||||
func TestCursorThinkingStreamForwardsDeltasAndSeparatesBlocks(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
var stream cursorThinkingStream
|
||||
if got := stream.delta("first "); got != "first " {
|
||||
t.Fatalf("first delta = %q, want %q", got, "first ")
|
||||
}
|
||||
if got := stream.delta("block"); got != "block" {
|
||||
t.Fatalf("mid-block delta = %q, want %q", got, "block")
|
||||
}
|
||||
stream.complete()
|
||||
// A new block must not be glued onto the previous one: the daemon
|
||||
// concatenates every thinking message into one transcript entry.
|
||||
if got := stream.delta("second block"); got != "\n\nsecond block" {
|
||||
t.Fatalf("first delta of next block = %q, want %q", got, "\n\nsecond block")
|
||||
}
|
||||
stream.complete()
|
||||
}
|
||||
|
||||
func TestCursorThinkingStreamDropsEmptyDeltas(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
var stream cursorThinkingStream
|
||||
if got := stream.delta(""); got != "" {
|
||||
t.Fatalf("empty delta = %q, want empty", got)
|
||||
}
|
||||
// An empty delta must not open a block, so the first real fragment is still
|
||||
// treated as the very first content (no leading separator).
|
||||
if got := stream.delta("real"); got != "real" {
|
||||
t.Fatalf("first real delta = %q, want %q", got, "real")
|
||||
}
|
||||
}
|
||||
|
||||
func TestParseCursorToolCall(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
tests := []struct {
|
||||
name string
|
||||
event string
|
||||
wantName string
|
||||
wantCallID string
|
||||
wantArg string
|
||||
wantResult string
|
||||
}{
|
||||
{
|
||||
name: "started names the tool from the nested key",
|
||||
event: `{"type":"tool_call","subtype":"started",
|
||||
"call_id":"call-1\nfc_9",
|
||||
"tool_call":{"readToolCall":{"args":{"path":"/tmp/a.txt"}},"toolCallId":"call-1\nfc_9"}}`,
|
||||
wantName: "read",
|
||||
wantCallID: "call-1",
|
||||
wantArg: "/tmp/a.txt",
|
||||
},
|
||||
{
|
||||
name: "completed carries the tool result payload",
|
||||
event: `{"type":"tool_call","subtype":"completed",
|
||||
"call_id":"call-2",
|
||||
"tool_call":{"shellToolCall":{"args":{"path":"x"},"result":{"success":{"exitCode":0}}}}}`,
|
||||
wantName: "shell",
|
||||
wantCallID: "call-2",
|
||||
wantArg: "x",
|
||||
wantResult: `{"success":{"exitCode":0}}`,
|
||||
},
|
||||
{
|
||||
name: "call id falls back to the nested tool call id",
|
||||
event: `{"type":"tool_call","subtype":"started",
|
||||
"tool_call":{"editToolCall":{"args":{"path":"y"}},"toolCallId":"call-3\nfc_1"}}`,
|
||||
wantName: "edit",
|
||||
wantCallID: "call-3",
|
||||
wantArg: "y",
|
||||
},
|
||||
{
|
||||
name: "unknown payload still yields a pairable call id",
|
||||
event: `{"type":"tool_call","subtype":"started","call_id":"call-4",
|
||||
"tool_call":{"toolCallId":"call-4","startedAtMs":"1"}}`,
|
||||
wantCallID: "call-4",
|
||||
},
|
||||
{
|
||||
name: "missing tool_call object degrades to the top-level id",
|
||||
event: `{"type":"tool_call","subtype":"completed","call_id":"call-5"}`,
|
||||
wantCallID: "call-5",
|
||||
},
|
||||
}
|
||||
|
||||
for _, tt := range tests {
|
||||
t.Run(tt.name, func(t *testing.T) {
|
||||
t.Parallel()
|
||||
var evt cursorStreamEvent
|
||||
if err := json.Unmarshal([]byte(tt.event), &evt); err != nil {
|
||||
t.Fatalf("unmarshal event: %v", err)
|
||||
}
|
||||
call := parseCursorToolCall(&evt)
|
||||
if call.Name != tt.wantName {
|
||||
t.Errorf("Name = %q, want %q", call.Name, tt.wantName)
|
||||
}
|
||||
if call.CallID != tt.wantCallID {
|
||||
t.Errorf("CallID = %q, want %q", call.CallID, tt.wantCallID)
|
||||
}
|
||||
if tt.wantArg != "" && call.Input["path"] != tt.wantArg {
|
||||
t.Errorf("Input[path] = %v, want %q", call.Input["path"], tt.wantArg)
|
||||
}
|
||||
if call.Result != tt.wantResult {
|
||||
t.Errorf("Result = %q, want %q", call.Result, tt.wantResult)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestCursorToolPayloadKeyIsDeterministic(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
envelope := map[string]json.RawMessage{
|
||||
"toolCallId": json.RawMessage(`"call-1"`),
|
||||
"startedAtMs": json.RawMessage(`"1"`),
|
||||
"shellToolCall": json.RawMessage(`{}`),
|
||||
"readToolCall": json.RawMessage(`{}`),
|
||||
"ToolCall": json.RawMessage(`{}`),
|
||||
"someOtherField": json.RawMessage(`{}`),
|
||||
}
|
||||
for i := 0; i < 20; i++ {
|
||||
if got := cursorToolPayloadKey(envelope); got != "readToolCall" {
|
||||
t.Fatalf("iteration %d: key = %q, want readToolCall", i, got)
|
||||
}
|
||||
}
|
||||
if got := cursorToolPayloadKey(map[string]json.RawMessage{"toolCallId": json.RawMessage(`"x"`)}); got != "" {
|
||||
t.Fatalf("key without a tool payload = %q, want empty", got)
|
||||
}
|
||||
}
|
||||
|
||||
23
server/pkg/agent/testdata/cursor-agent-2026.07.20-stream-json.jsonl
vendored
Normal file
23
server/pkg/agent/testdata/cursor-agent-2026.07.20-stream-json.jsonl
vendored
Normal file
@@ -0,0 +1,23 @@
|
||||
{"type":"system","subtype":"init","apiKeySource":"login","cwd":"/private/tmp/curcap","session_id":"ebb521c2-404d-4a4e-8c2f-1f8bdb141043","model":"Auto","permissionMode":"default"}
|
||||
{"type":"user","message":{"role":"user","content":[{"type":"text","text":"Read notes.txt, then run \"wc -l notes.txt\" in the shell, then write the number of lines into count.txt. Keep it short."}]},"session_id":"ebb521c2-404d-4a4e-8c2f-1f8bdb141043"}
|
||||
{"type":"thinking","subtype":"delta","text":"Reading notes.txt.","session_id":"ebb521c2-404d-4a4e-8c2f-1f8bdb141043","timestamp_ms":1784819288883}
|
||||
{"type":"thinking","subtype":"delta","text":"\n\nRunning wc -l on notes.txt.","session_id":"ebb521c2-404d-4a4e-8c2f-1f8bdb141043","timestamp_ms":1784819289514}
|
||||
{"type":"thinking","subtype":"delta","text":" Writing the line count","session_id":"ebb521c2-404d-4a4e-8c2f-1f8bdb141043","timestamp_ms":1784819289514}
|
||||
{"type":"thinking","subtype":"delta","text":" to count.txt.","session_id":"ebb521c2-404d-4a4e-8c2f-1f8bdb141043","timestamp_ms":1784819289514}
|
||||
{"type":"thinking","subtype":"completed","session_id":"ebb521c2-404d-4a4e-8c2f-1f8bdb141043","timestamp_ms":1784819289515}
|
||||
{"type":"assistant","message":{"role":"assistant","content":[{"type":"text","text":"I'll read `notes.txt`, run `wc -l`, then write the line count to `count.txt`."}]},"session_id":"ebb521c2-404d-4a4e-8c2f-1f8bdb141043","model_call_id":"59ff63a6-195e-463f-89ce-f26a820987bf-0-c94p","timestamp_ms":1784819289520}
|
||||
{"type":"tool_call","subtype":"started","call_id":"call-bb11656a-e59e-4356-9866-5b206aedb390-0\nfc_35bc3e26-1dfc-9c07-b668-4c50a744b8f9_0","tool_call":{"readToolCall":{"args":{"path":"/tmp/curcap/notes.txt"}},"hookAdditionalContexts":[],"toolCallId":"call-bb11656a-e59e-4356-9866-5b206aedb390-0\nfc_35bc3e26-1dfc-9c07-b668-4c50a744b8f9_0","startedAtMs":"1784819289745"},"model_call_id":"59ff63a6-195e-463f-89ce-f26a820987bf-0-c94p","session_id":"ebb521c2-404d-4a4e-8c2f-1f8bdb141043","timestamp_ms":1784819289520}
|
||||
{"type":"tool_call","subtype":"started","call_id":"call-bb11656a-e59e-4356-9866-5b206aedb390-1\nfc_35bc3e26-1dfc-9c07-b668-4c50a744b8f9_1","tool_call":{"shellToolCall":{"args":{"command":"wc -l notes.txt","workingDirectory":"","timeout":30000,"toolCallId":"call-bb11656a-e59e-4356-9866-5b206aedb390-1\nfc_35bc3e26-1dfc-9c07-b668-4c50a744b8f9_1","simpleCommands":["wc"],"hasInputRedirect":false,"hasOutputRedirect":false,"parsingResult":{"parsingFailed":false,"executableCommands":[{"name":"wc","args":[{"type":"word","value":"-l"},{"type":"word","value":"notes.txt"}],"fullText":"wc -l notes.txt"}],"hasRedirects":false,"hasCommandSubstitution":false,"redirects":[]},"fileOutputThresholdBytes":"40000","isBackground":false,"skipApproval":false,"timeoutBehavior":"TIMEOUT_BEHAVIOR_BACKGROUND","hardTimeout":86400000,"description":"Count lines in notes.txt","closeStdin":true,"conversationId":"ebb521c2-404d-4a4e-8c2f-1f8bdb141043"},"description":"Count lines in notes.txt"},"hookAdditionalContexts":[],"toolCallId":"call-bb11656a-e59e-4356-9866-5b206aedb390-1\nfc_35bc3e26-1dfc-9c07-b668-4c50a744b8f9_1","startedAtMs":"1784819289750"},"model_call_id":"59ff63a6-195e-463f-89ce-f26a820987bf-0-c94p","session_id":"ebb521c2-404d-4a4e-8c2f-1f8bdb141043","timestamp_ms":1784819289526}
|
||||
{"type":"tool_call","subtype":"completed","call_id":"call-bb11656a-e59e-4356-9866-5b206aedb390-0\nfc_35bc3e26-1dfc-9c07-b668-4c50a744b8f9_0","tool_call":{"readToolCall":{"args":{"path":"/tmp/curcap/notes.txt"},"result":{"success":{"content":"alpha\nbeta\ngamma\n","isEmpty":false,"exceededLimit":false,"totalLines":4,"fileSize":17,"path":"/tmp/curcap/notes.txt","readRange":{"startLine":1,"endLine":4},"relatedCursorRulePaths":[],"relatedCursorRules":[]}}},"hookAdditionalContexts":[],"toolCallId":"call-bb11656a-e59e-4356-9866-5b206aedb390-0\nfc_35bc3e26-1dfc-9c07-b668-4c50a744b8f9_0","startedAtMs":"1784819289745","completedAtMs":"1784819289937"},"model_call_id":"59ff63a6-195e-463f-89ce-f26a820987bf-0-c94p","session_id":"ebb521c2-404d-4a4e-8c2f-1f8bdb141043","timestamp_ms":1784819289719}
|
||||
{"type":"tool_call","subtype":"completed","call_id":"call-bb11656a-e59e-4356-9866-5b206aedb390-1\nfc_35bc3e26-1dfc-9c07-b668-4c50a744b8f9_1","tool_call":{"shellToolCall":{"args":{"command":"wc -l notes.txt","workingDirectory":"","timeout":30000,"toolCallId":"call-bb11656a-e59e-4356-9866-5b206aedb390-1\nfc_35bc3e26-1dfc-9c07-b668-4c50a744b8f9_1","simpleCommands":["wc"],"hasInputRedirect":false,"hasOutputRedirect":false,"parsingResult":{"parsingFailed":false,"executableCommands":[{"name":"wc","args":[{"type":"word","value":"-l"},{"type":"word","value":"notes.txt"}],"fullText":"wc -l notes.txt"}],"hasRedirects":false,"hasCommandSubstitution":false,"redirects":[]},"fileOutputThresholdBytes":"40000","isBackground":false,"skipApproval":false,"timeoutBehavior":"TIMEOUT_BEHAVIOR_BACKGROUND","hardTimeout":86400000,"description":"Count lines in notes.txt","closeStdin":true,"conversationId":"ebb521c2-404d-4a4e-8c2f-1f8bdb141043"},"result":{"success":{"command":"wc -l notes.txt","workingDirectory":"","exitCode":0,"signal":"","stdout":" 3 notes.txt\n","stderr":"","executionTime":1479,"interleavedOutput":" 3 notes.txt\n","localExecutionTimeMs":280},"isBackground":false},"description":"Count lines in notes.txt"},"hookAdditionalContexts":[],"toolCallId":"call-bb11656a-e59e-4356-9866-5b206aedb390-1\nfc_35bc3e26-1dfc-9c07-b668-4c50a744b8f9_1","startedAtMs":"1784819289750","completedAtMs":"1784819291229"},"model_call_id":"59ff63a6-195e-463f-89ce-f26a820987bf-0-c94p","session_id":"ebb521c2-404d-4a4e-8c2f-1f8bdb141043","timestamp_ms":1784819291011}
|
||||
{"type":"thinking","subtype":"delta","text":"notes.txt contains 3","session_id":"ebb521c2-404d-4a4e-8c2f-1f8bdb141043","timestamp_ms":1784819292716}
|
||||
{"type":"thinking","subtype":"delta","text":" lines. I will write","session_id":"ebb521c2-404d-4a4e-8c2f-1f8bdb141043","timestamp_ms":1784819292716}
|
||||
{"type":"thinking","subtype":"delta","text":" 3 into count.txt.","session_id":"ebb521c2-404d-4a4e-8c2f-1f8bdb141043","timestamp_ms":1784819292716}
|
||||
{"type":"thinking","subtype":"completed","session_id":"ebb521c2-404d-4a4e-8c2f-1f8bdb141043","timestamp_ms":1784819292716}
|
||||
{"type":"tool_call","subtype":"started","call_id":"call-c52c0cd6-81ad-4a87-94c5-f5b0119f3ed4-2\nfc_27b4c3fb-9abc-92ae-967f-147ee264096e_0","tool_call":{"editToolCall":{"args":{"path":"/tmp/curcap/count.txt","streamContent":"3\n"}},"hookAdditionalContexts":[],"toolCallId":"call-c52c0cd6-81ad-4a87-94c5-f5b0119f3ed4-2\nfc_27b4c3fb-9abc-92ae-967f-147ee264096e_0","startedAtMs":"1784819292927"},"model_call_id":"59ff63a6-195e-463f-89ce-f26a820987bf-1-g9hb","session_id":"ebb521c2-404d-4a4e-8c2f-1f8bdb141043","timestamp_ms":1784819292717}
|
||||
{"type":"tool_call","subtype":"completed","call_id":"call-c52c0cd6-81ad-4a87-94c5-f5b0119f3ed4-2\nfc_27b4c3fb-9abc-92ae-967f-147ee264096e_0","tool_call":{"editToolCall":{"args":{"path":"/tmp/curcap/count.txt","streamContent":"3\n"},"result":{"success":{"path":"/tmp/curcap/count.txt","linesAdded":1,"linesRemoved":1,"diffString":"--- /dev/null\n+++ b//tmp/curcap/count.txt\n@@ -1 +1 @@\n-\n+3","afterFullFileContent":"3\n","message":"Wrote contents to /tmp/curcap/count.txt"}}},"hookAdditionalContexts":[],"toolCallId":"call-c52c0cd6-81ad-4a87-94c5-f5b0119f3ed4-2\nfc_27b4c3fb-9abc-92ae-967f-147ee264096e_0","startedAtMs":"1784819292927","completedAtMs":"1784819293314"},"model_call_id":"59ff63a6-195e-463f-89ce-f26a820987bf-1-g9hb","session_id":"ebb521c2-404d-4a4e-8c2f-1f8bdb141043","timestamp_ms":1784819293094}
|
||||
{"type":"thinking","subtype":"delta","text":"Finished. The line count","session_id":"ebb521c2-404d-4a4e-8c2f-1f8bdb141043","timestamp_ms":1784819294133}
|
||||
{"type":"thinking","subtype":"delta","text":" was written to count.txt.","session_id":"ebb521c2-404d-4a4e-8c2f-1f8bdb141043","timestamp_ms":1784819294133}
|
||||
{"type":"thinking","subtype":"completed","session_id":"ebb521c2-404d-4a4e-8c2f-1f8bdb141043","timestamp_ms":1784819294133}
|
||||
{"type":"assistant","message":{"role":"assistant","content":[{"type":"text","text":"`notes.txt` has 3 lines (`alpha`, `beta`, `gamma`). Wrote `3` to `count.txt`."}]},"session_id":"ebb521c2-404d-4a4e-8c2f-1f8bdb141043"}
|
||||
{"type":"result","subtype":"success","duration_ms":9679,"duration_api_ms":9679,"is_error":false,"result":"I'll read `notes.txt`, run `wc -l`, then write the line count to `count.txt`.`notes.txt` has 3 lines (`alpha`, `beta`, `gamma`). Wrote `3` to `count.txt`.","session_id":"ebb521c2-404d-4a4e-8c2f-1f8bdb141043","request_id":"59ff63a6-195e-463f-89ce-f26a820987bf","usage":{"inputTokens":8830,"outputTokens":185,"cacheReadTokens":44800,"cacheWriteTokens":0}}
|
||||
Reference in New Issue
Block a user