diff --git a/arkruntime/selfhosted/session_tool_runner.go b/arkruntime/selfhosted/session_tool_runner.go index 9c9fc9e..2b3936d 100644 --- a/arkruntime/selfhosted/session_tool_runner.go +++ b/arkruntime/selfhosted/session_tool_runner.go @@ -171,6 +171,7 @@ func (r *SessionToolRunner) run() { defer close(r.done) defer close(r.events) defer r.drainInFlight() + defer r.cancel() r.err = normalizeRunnerErr(r.runLoop()) } @@ -195,6 +196,9 @@ func (r *SessionToolRunner) runLoop() error { pendingAsk: map[string]Event{}, confirmations: map[string]Event{}, externalTools: map[string]Event{}, + toolUseEvents: map[string]Event{}, + scheduled: map[string]bool{}, + executionDone: make(chan toolExecutionResult, sessionRunnerResultsBuffer), sessionToolUses: map[string]bool{}, toolUsesSinceStatus: map[string]bool{}, blockingEventIDs: map[string]bool{}, @@ -231,16 +235,22 @@ func (r *SessionToolRunner) runLoop() error { } type toolRunnerState struct { - runner *SessionToolRunner - page string - processed map[string]bool - seen map[string]bool - answered map[string]bool - pendingResults map[string]Event - recoveredResults map[string]bool - pendingAsk map[string]Event - confirmations map[string]Event - externalTools map[string]Event + runner *SessionToolRunner + page string + processed map[string]bool + seen map[string]bool + answered map[string]bool + pendingResults map[string]Event + recoveredResults map[string]bool + pendingAsk map[string]Event + confirmations map[string]Event + externalTools map[string]Event + toolUseEvents map[string]Event + // Tool 按 CMA runner 的语义串行执行,事件消费通过完成通道保持非阻塞。 + scheduled map[string]bool + executionQueue []pendingToolEvent + activeExecution *activeToolExecution + executionDone chan toolExecutionResult sessionToolUses map[string]bool toolUsesSinceStatus map[string]bool blockingEventIDs map[string]bool @@ -250,8 +260,19 @@ type toolRunnerState struct { } type pendingToolEvent struct { - event Event - custom bool + event Event + custom bool + confirmation string +} + +type activeToolExecution struct { + pendingToolEvent + cancel context.CancelFunc +} + +type toolExecutionResult struct { + pendingToolEvent + result toolset.Result } func (s *toolRunnerState) consumeStreamLoop(ctx context.Context, streamer EventStreamer) error { @@ -352,6 +373,10 @@ func (s *toolRunnerState) consumeStream(ctx context.Context, stream *EventStream select { case <-ctx.Done(): return ctx.Err() + case result := <-s.executionDone: + if err := s.finishToolExecution(ctx, result); err != nil { + return err + } case <-timerC: if s.idleExpired() { return ErrIdleTimeout @@ -441,6 +466,11 @@ func (s *toolRunnerState) sleepOrIdle(ctx context.Context, d time.Duration) erro case <-ctx.Done(): timer.Stop() return ctx.Err() + case result := <-s.executionDone: + timer.Stop() + if err := s.finishToolExecution(ctx, result); err != nil { + return err + } case <-timer.C: } } @@ -539,6 +569,10 @@ func (s *toolRunnerState) consumeList(ctx context.Context) error { select { case <-ctx.Done(): return ctx.Err() + case result := <-s.executionDone: + if err := s.finishToolExecution(ctx, result); err != nil { + return err + } case <-timerC: if s.idleExpired() { return ErrIdleTimeout @@ -551,6 +585,7 @@ func (s *toolRunnerState) consumeList(ctx context.Context) error { func (s *toolRunnerState) processListedEvents(ctx context.Context, events []Event, reconcile bool) error { var pending []pendingToolEvent pendingIDs := map[string]bool{} + replayedToolUses := map[string]bool{} touchedIdle := false lastWasEndTurn := false for _, event := range events { @@ -559,7 +594,7 @@ func (s *toolRunnerState) processListedEvents(ctx context.Context, events []Even continue } s.observeSessionState(event) - if event.Type != EventTypeUserToolConfirmation { + if seenNow && event.Type != EventTypeUserToolConfirmation { touchedIdle = true lastWasEndTurn = event.Type == EventTypeSessionStatusIdle && event.StopReasonType() == SessionStopReasonEndTurn @@ -567,16 +602,24 @@ func (s *toolRunnerState) processListedEvents(ctx context.Context, events []Even switch event.Type { case EventTypeUserToolConfirmation: s.recordConfirmation(event) + case EventTypeUserInterrupt: + if reconcile { + s.handleInterruptForCalls(event, replayedToolUses) + } else { + s.handleInterrupt(event) + } case EventTypeUserToolResult, EventTypeUserCustomToolResult: s.markAnswered(toolResultCallID(event)) case EventTypeAgentToolUse: callID := toolUseCallID(event) + replayedToolUses[callID] = true if !pendingIDs[callID] { pending = append(pending, pendingToolEvent{event: event}) pendingIDs[callID] = true } case EventTypeAgentCustomToolUse: callID := toolUseCallID(event) + replayedToolUses[callID] = true if !pendingIDs[callID] { pending = append(pending, pendingToolEvent{event: event, custom: true}) pendingIDs[callID] = true @@ -602,11 +645,7 @@ func (s *toolRunnerState) processListedEvents(ctx context.Context, events []Even return err } if touchedIdle && lastWasEndTurn { - if s.hasUnblockedOutstandingTool(pending) { - s.disarmIdle() - } else { - s.armIdle() - } + s.armIdle() } return nil } @@ -666,7 +705,7 @@ func (s *toolRunnerState) maybeArmPendingIdle() { } func (s *toolRunnerState) hasIdleBlockers() bool { - return len(s.pendingAsk) > 0 || len(s.pendingResults) > 0 || len(s.externalTools) > 0 + return len(s.pendingAsk) > 0 || len(s.pendingResults) > 0 || len(s.externalTools) > 0 || len(s.scheduled) > 0 } func (s *toolRunnerState) idleExpired() bool { @@ -678,6 +717,8 @@ func (s *toolRunnerState) handleEvent(ctx context.Context, event Event) error { case EventTypeUserToolConfirmation: s.recordConfirmation(event) return s.releaseConfirmedToolUses(ctx) + case EventTypeUserInterrupt: + s.handleInterrupt(event) case EventTypeUserToolResult, EventTypeUserCustomToolResult: s.markAnswered(toolResultCallID(event)) case EventTypeAgentToolUse: @@ -695,6 +736,9 @@ func (s *toolRunnerState) markEventSeen(event Event) bool { if key == "" { key = toolUseCallID(event) } + if key == "" && event.Type == EventTypeUserInterrupt { + key = fmt.Sprintf("interrupt:%s:%s", event.ProcessedAt, event.SessionThreadID) + } if key == "" { return true } @@ -715,6 +759,7 @@ func (s *toolRunnerState) markAnswered(callID string) { delete(s.recoveredResults, callID) delete(s.pendingAsk, callID) delete(s.externalTools, callID) + delete(s.toolUseEvents, callID) s.maybeArmPendingIdle() } @@ -745,26 +790,10 @@ func (s *toolRunnerState) releaseConfirmedToolUses(ctx context.Context) error { return nil } -func (s *toolRunnerState) hasUnblockedOutstandingTool(pending []pendingToolEvent) bool { - for _, toolEvent := range pending { - callID := toolUseCallID(toolEvent.event) - if callID == "" || s.isAnswered(callID) || !s.shouldHandleToolUse(callID) { - continue - } - if _, ok := s.pendingAsk[callID]; ok { - continue - } - if _, ok := s.pendingResults[callID]; ok { - continue - } - return true - } - return false -} - func (s *toolRunnerState) handleToolUse(ctx context.Context, event Event, custom bool) error { + s.ensureRecoveryMaps() callID := toolUseCallID(event) - if callID == "" || s.isAnswered(callID) { + if callID == "" || s.isAnswered(callID) || s.scheduled[callID] { return nil } if pending := s.pendingResults[callID]; pending.ID != "" { @@ -795,9 +824,40 @@ func (s *toolRunnerState) handleToolUse(ctx context.Context, event Event, custom return s.sendResult(ctx, callID, event, custom, "", decision.Result) } } - var result toolset.Result + s.scheduled[callID] = true + s.executionQueue = append(s.executionQueue, pendingToolEvent{ + event: event, + custom: custom, + confirmation: confirmation, + }) + s.startNextToolExecution(ctx) + return nil +} + +func (s *toolRunnerState) startNextToolExecution(ctx context.Context) { + if s.activeExecution != nil { + return + } + for len(s.executionQueue) > 0 { + pending := s.executionQueue[0] + s.executionQueue = s.executionQueue[1:] + callID := toolUseCallID(pending.event) + if s.isAnswered(callID) { + delete(s.scheduled, callID) + continue + } + toolCtx, cancel := context.WithCancel(ctx) + s.activeExecution = &activeToolExecution{pendingToolEvent: pending, cancel: cancel} + go s.executeTool(toolCtx, pending) + return + } +} + +func (s *toolRunnerState) executeTool(ctx context.Context, pending pendingToolEvent) { + event := pending.event input := json.RawMessage(event.Input) - if custom { + var result toolset.Result + if pending.custom { tool := s.runner.opts.CustomTools[event.Name] result = s.executeWithTimeout(ctx, event, func(toolCtx context.Context) toolset.Result { return tool.Execute(toolCtx, input) @@ -807,7 +867,71 @@ func (s *toolRunnerState) handleToolUse(ctx context.Context, event Event, custom return s.runner.opts.Tools.Execute(toolCtx, event.Name, input) }) } - return s.postResult(ctx, event, custom, callID, result, confirmation) + select { + case s.executionDone <- toolExecutionResult{pendingToolEvent: pending, result: result}: + case <-s.runner.ctx.Done(): + } +} + +func (s *toolRunnerState) finishToolExecution(ctx context.Context, result toolExecutionResult) error { + callID := toolUseCallID(result.event) + if active := s.activeExecution; active != nil && toolUseCallID(active.event) == callID { + active.cancel() + s.activeExecution = nil + } + delete(s.scheduled, callID) + if !s.isAnswered(callID) { + if err := s.postResult(ctx, result.event, result.custom, callID, result.result, result.confirmation); err != nil { + return err + } + } + s.startNextToolExecution(ctx) + return nil +} + +func (s *toolRunnerState) handleInterrupt(event Event) { + s.handleInterruptForCalls(event, nil) +} + +func (s *toolRunnerState) handleInterruptForCalls(event Event, eligible map[string]bool) { + threadID := event.SessionThreadID + matches := func(callID string, toolEvent Event) bool { + if eligible != nil && !eligible[callID] { + return false + } + return threadID == "" || toolEvent.SessionThreadID == threadID + } + if active := s.activeExecution; active != nil && matches(toolUseCallID(active.event), active.event) { + active.cancel() + s.settleInterruptedToolUse(toolUseCallID(active.event)) + } + queued := s.executionQueue[:0] + for _, pending := range s.executionQueue { + if matches(toolUseCallID(pending.event), pending.event) { + s.settleInterruptedToolUse(toolUseCallID(pending.event)) + continue + } + queued = append(queued, pending) + } + s.executionQueue = queued + for callID, toolEvent := range s.toolUseEvents { + if !s.isAnswered(callID) && matches(callID, toolEvent) { + s.settleInterruptedToolUse(callID) + } + } +} + +func (s *toolRunnerState) settleInterruptedToolUse(callID string) { + if callID == "" || s.isAnswered(callID) { + return + } + delete(s.scheduled, callID) + s.markAnswered(callID) + if discarder, ok := s.runner.opts.ResultStore.(ToolResultStoreDiscarder); ok { + if err := discarder.Discard(callID); err != nil { + s.runner.logger.Warn("discard interrupted tool result failed", "tool_use_id", callID, "err", err) + } + } } func (s *toolRunnerState) ownsTool(event Event, custom bool) bool { @@ -1012,6 +1136,9 @@ func (s *toolRunnerState) observeSessionState(event Event) { callID := toolUseCallID(event) if callID != "" { s.sessionToolUses[callID] = true + if !s.isAnswered(callID) { + s.toolUseEvents[callID] = event + } s.toolUsesSinceStatus[callID] = true } case EventTypeSessionStatusIdle: @@ -1043,6 +1170,15 @@ func (s *toolRunnerState) ensureRecoveryMaps() { if s.blockingEventIDs == nil { s.blockingEventIDs = map[string]bool{} } + if s.toolUseEvents == nil { + s.toolUseEvents = map[string]Event{} + } + if s.scheduled == nil { + s.scheduled = map[string]bool{} + } + if s.executionDone == nil { + s.executionDone = make(chan toolExecutionResult, sessionRunnerResultsBuffer) + } } func (s *toolRunnerState) reconcileRecoveredResults() { diff --git a/arkruntime/selfhosted/session_tool_runner_test.go b/arkruntime/selfhosted/session_tool_runner_test.go index c8a5ec2..5e37d7b 100644 --- a/arkruntime/selfhosted/session_tool_runner_test.go +++ b/arkruntime/selfhosted/session_tool_runner_test.go @@ -15,6 +15,8 @@ import ( "github.com/volcengine/ark-runtime-go/arkruntime/toolset" ) +const runnerTestCallID = "call-id" + type runnerTestAPI struct { listCalls int listEvent Event @@ -65,6 +67,35 @@ func (t *runnerTestTool) Execute(context.Context, json.RawMessage) toolset.Resul return toolset.TextResult("ok") } +type runnerBlockingTool struct { + started chan struct{} + canceled chan struct{} +} + +func (t *runnerBlockingTool) Name() string { return "blocking" } + +func (t *runnerBlockingTool) Execute(ctx context.Context, _ json.RawMessage) toolset.Result { + close(t.started) + <-ctx.Done() + close(t.canceled) + return toolset.ErrorResult(ctx.Err().Error()) +} + +type runnerNamedBlockingTool struct { + name string + started chan struct{} + canceled chan struct{} +} + +func (t *runnerNamedBlockingTool) Name() string { return t.name } + +func (t *runnerNamedBlockingTool) Execute(ctx context.Context, _ json.RawMessage) toolset.Result { + close(t.started) + <-ctx.Done() + close(t.canceled) + return toolset.ErrorResult(ctx.Err().Error()) +} + type runnerFailingMarkSentStore struct { markCalls int } @@ -115,7 +146,7 @@ func TestSessionToolRunnerReconcileRetriesWithoutLosingToolUse(t *testing.T) { ID: "event-id", Type: EventTypeAgentCustomToolUse, Name: tool.Name(), - ToolUseID: "call-id", + ToolUseID: runnerTestCallID, SessionThreadID: "thread-id", Input: RawJSON(`{}`), }} @@ -139,10 +170,11 @@ func TestSessionToolRunnerReconcileRetriesWithoutLosingToolUse(t *testing.T) { if err := state.reconcile(context.Background()); err != nil { t.Fatal(err) } + finishRunnerTestExecution(t, state) if api.listCalls != 2 || tool.calls != 1 || len(api.sent) != 1 { t.Fatalf("list_calls=%d tool_calls=%d sent=%d", api.listCalls, tool.calls, len(api.sent)) } - if api.sent[0].CustomToolUseID != "call-id" { + if api.sent[0].CustomToolUseID != runnerTestCallID { t.Fatalf("sent event=%+v", api.sent[0]) } } @@ -155,7 +187,7 @@ func TestSessionToolRunnerConvertsToolPanicToErrorResult(t *testing.T) { state := &toolRunnerState{runner: runner} result := state.executeWithTimeout(context.Background(), Event{ ID: "event-id", - ToolUseID: "call-id", + ToolUseID: runnerTestCallID, Name: "panicking-tool", }, func(context.Context) toolset.Result { panic("boom") @@ -223,20 +255,20 @@ func TestSessionToolRunnerMarkSentFailureDoesNotRetainDeliveredResult(t *testing Logger: log.New(io.Discard, "", 0), }) runner.events = make(chan ToolCallResult, 1) - out := NewUserToolResultEvent("call-id", []ContentBlock{{Type: "text", Text: "ok"}}, false, "thread-id") + out := NewUserToolResultEvent(runnerTestCallID, []ContentBlock{{Type: "text", Text: "ok"}}, false, "thread-id") state := &toolRunnerState{ runner: runner, processed: map[string]bool{}, answered: map[string]bool{}, - pendingResults: map[string]Event{"call-id": out}, + pendingResults: map[string]Event{runnerTestCallID: out}, } - if err := state.sendResult(context.Background(), "call-id", Event{Name: "bash"}, false, "", out); err != nil { + if err := state.sendResult(context.Background(), runnerTestCallID, Event{Name: "bash"}, false, "", out); err != nil { t.Fatal(err) } if len(api.sent) != 1 || store.markCalls != 1 { t.Fatalf("sent=%d mark_calls=%d", len(api.sent), store.markCalls) } - if !state.isAnswered("call-id") || len(state.pendingResults) != 0 { + if !state.isAnswered(runnerTestCallID) || len(state.pendingResults) != 0 { t.Fatalf("answered=%v pending=%v", state.answered, state.pendingResults) } } @@ -248,12 +280,12 @@ func TestSessionToolRunnerFlushMarkSentFailureDoesNotRetainDeliveredResult(t *te ResultStore: store, Logger: log.New(io.Discard, "", 0), }) - out := NewUserToolResultEvent("call-id", []ContentBlock{{Type: "text", Text: "ok"}}, false, "thread-id") + out := NewUserToolResultEvent(runnerTestCallID, []ContentBlock{{Type: "text", Text: "ok"}}, false, "thread-id") state := &toolRunnerState{ runner: runner, processed: map[string]bool{}, answered: map[string]bool{}, - pendingResults: map[string]Event{"call-id": out}, + pendingResults: map[string]Event{runnerTestCallID: out}, } if err := state.flushResults(context.Background()); err != nil { t.Fatal(err) @@ -261,7 +293,7 @@ func TestSessionToolRunnerFlushMarkSentFailureDoesNotRetainDeliveredResult(t *te if len(api.sent) != 1 || store.markCalls != 1 { t.Fatalf("sent=%d mark_calls=%d", len(api.sent), store.markCalls) } - if !state.isAnswered("call-id") || len(state.pendingResults) != 0 { + if !state.isAnswered(runnerTestCallID) || len(state.pendingResults) != 0 { t.Fatalf("answered=%v pending=%v", state.answered, state.pendingResults) } } @@ -291,13 +323,13 @@ func TestSessionToolRunnerPermissionDenyDoesNotPostToolResult(t *testing.T) { ID: "event-id", Type: EventTypeAgentToolUse, Name: "read", - ToolUseID: "call-id", + ToolUseID: runnerTestCallID, EvaluatedPermission: PermissionDeny, } if err := state.handleToolUse(context.Background(), event, false); err != nil { t.Fatal(err) } - if len(api.sent) != 0 || !state.isAnswered("call-id") { + if len(api.sent) != 0 || !state.isAnswered(runnerTestCallID) { t.Fatalf("sent=%d answered=%v", len(api.sent), state.answered) } result := <-runner.events @@ -347,6 +379,7 @@ func TestSessionToolRunnerFiltersRecoveredResultsAgainstCurrentBlockers(t *testi if err := state.processListedEvents(context.Background(), events, true); err != nil { t.Fatal(err) } + finishRunnerTestExecution(t, state) if tool.calls != 1 || len(api.sent) != 1 || api.sent[0].CustomToolUseID != "current-call" { t.Fatalf("tool_calls=%d sent=%+v", tool.calls, api.sent) } @@ -358,6 +391,332 @@ func TestSessionToolRunnerFiltersRecoveredResultsAgainstCurrentBlockers(t *testi } } +func TestSessionToolRunnerInterruptCancelsActiveToolWithoutPostingResult(t *testing.T) { + tool := &runnerBlockingTool{started: make(chan struct{}), canceled: make(chan struct{})} + api := &runnerTestAPI{} + store := &runnerRecordingStore{} + runner := NewSessionToolRunner(context.Background(), api, "session-id", SessionToolRunnerOptions{ + CustomTools: map[string]toolset.Tool{tool.Name(): tool}, + ResultStore: store, + Logger: log.New(io.Discard, "", 0), + ToolTimeout: time.Second, + }) + runner.events = make(chan ToolCallResult, 1) + state := newRunnerTestState(runner) + toolUse := Event{ + ID: runnerTestCallID, + Type: EventTypeAgentCustomToolUse, + Name: tool.Name(), + ToolUseID: runnerTestCallID, + SessionThreadID: "thread-id", + Input: RawJSON(`{}`), + } + if err := state.handleStreamEvent(context.Background(), toolUse); err != nil { + t.Fatal(err) + } + select { + case <-tool.started: + case <-time.After(time.Second): + t.Fatal("tool did not start") + } + if err := state.handleStreamEvent(context.Background(), Event{ + ID: "interrupt-id", + Type: EventTypeUserInterrupt, + SessionThreadID: "thread-id", + }); err != nil { + t.Fatal(err) + } + select { + case <-tool.canceled: + case <-time.After(time.Second): + t.Fatal("tool context was not canceled") + } + finishRunnerTestExecution(t, state) + + if len(api.sent) != 0 || !state.isAnswered(runnerTestCallID) { + t.Fatalf("sent=%d answered=%v", len(api.sent), state.answered) + } + if len(store.discarded) != 1 || store.discarded[0] != runnerTestCallID { + t.Fatalf("discarded=%v", store.discarded) + } +} + +func TestSessionToolRunnerInterruptOnlyCancelsTargetThread(t *testing.T) { + tool := &runnerBlockingTool{started: make(chan struct{}), canceled: make(chan struct{})} + runner := NewSessionToolRunner(context.Background(), &runnerTestAPI{}, "session-id", SessionToolRunnerOptions{ + CustomTools: map[string]toolset.Tool{tool.Name(): tool}, + Logger: log.New(io.Discard, "", 0), + ToolTimeout: time.Second, + }) + runner.events = make(chan ToolCallResult, 1) + state := newRunnerTestState(runner) + toolUse := Event{ + ID: runnerTestCallID, + Type: EventTypeAgentCustomToolUse, + Name: tool.Name(), + ToolUseID: runnerTestCallID, + SessionThreadID: "thread-a", + Input: RawJSON(`{}`), + } + if err := state.handleStreamEvent(context.Background(), toolUse); err != nil { + t.Fatal(err) + } + select { + case <-tool.started: + case <-time.After(time.Second): + t.Fatal("tool did not start") + } + if err := state.handleStreamEvent(context.Background(), Event{ + ID: "other-interrupt", + Type: EventTypeUserInterrupt, + SessionThreadID: "thread-b", + }); err != nil { + t.Fatal(err) + } + if state.isAnswered(runnerTestCallID) { + t.Fatal("interrupt for another thread settled the active call") + } + select { + case <-tool.canceled: + t.Fatal("interrupt for another thread canceled the active call") + default: + } + + state.handleInterrupt(Event{Type: EventTypeUserInterrupt, SessionThreadID: "thread-a"}) + finishRunnerTestExecution(t, state) +} + +func TestSessionToolRunnerReconcileDoesNotRedispatchInterruptedToolUse(t *testing.T) { + tool := &runnerTestTool{} + api := &runnerTestAPI{} + runner := NewSessionToolRunner(context.Background(), api, "session-id", SessionToolRunnerOptions{ + CustomTools: map[string]toolset.Tool{tool.Name(): tool}, + Logger: log.New(io.Discard, "", 0), + }) + runner.events = make(chan ToolCallResult, 1) + state := newRunnerTestState(runner) + events := []Event{ + { + ID: runnerTestCallID, + Type: EventTypeAgentCustomToolUse, + Name: tool.Name(), + ToolUseID: runnerTestCallID, + SessionThreadID: "thread-id", + Input: RawJSON(`{}`), + }, + {ID: "interrupt-id", Type: EventTypeUserInterrupt, SessionThreadID: "thread-id"}, + {ID: "idle-id", Type: EventTypeSessionStatusIdle, StopReason: &SessionStopReason{Type: SessionStopReasonEndTurn}}, + } + if err := state.processListedEvents(context.Background(), events, true); err != nil { + t.Fatal(err) + } + + if tool.calls != 0 || len(api.sent) != 0 || !state.isAnswered(runnerTestCallID) { + t.Fatalf("tool_calls=%d sent=%d answered=%v", tool.calls, len(api.sent), state.answered) + } +} + +func TestSessionToolRunnerListInterruptCancelsToolFromEarlierPoll(t *testing.T) { + tool := &runnerBlockingTool{started: make(chan struct{}), canceled: make(chan struct{})} + runner := NewSessionToolRunner(context.Background(), &runnerTestAPI{}, "session-id", SessionToolRunnerOptions{ + CustomTools: map[string]toolset.Tool{tool.Name(): tool}, + Logger: log.New(io.Discard, "", 0), + ToolTimeout: time.Second, + }) + runner.events = make(chan ToolCallResult, 1) + state := newRunnerTestState(runner) + toolUse := Event{ + ID: "cross-poll-call", Type: EventTypeAgentCustomToolUse, Name: tool.Name(), + ToolUseID: "cross-poll-call", SessionThreadID: "thread-id", Input: RawJSON(`{}`), + } + if err := state.processListedEvents(context.Background(), []Event{toolUse}, false); err != nil { + t.Fatal(err) + } + select { + case <-tool.started: + case <-time.After(time.Second): + t.Fatal("tool did not start") + } + interrupt := Event{ + Type: EventTypeUserInterrupt, ProcessedAt: "2026-09-20T00:00:00Z", SessionThreadID: "thread-id", + } + if err := state.processListedEvents(context.Background(), []Event{toolUse, interrupt}, false); err != nil { + t.Fatal(err) + } + select { + case <-tool.canceled: + case <-time.After(time.Second): + t.Fatal("cross-poll interrupt did not cancel tool") + } + finishRunnerTestExecution(t, state) + if !state.isAnswered("cross-poll-call") { + t.Fatal("cross-poll interrupt did not settle tool call") + } +} + +func TestSessionToolRunnerSerialQueueInterruptRemovesMatchingToolOnly(t *testing.T) { + activeTool := &runnerNamedBlockingTool{name: "active", started: make(chan struct{}), canceled: make(chan struct{})} + queuedTool := &runnerNamedBlockingTool{name: "queued", started: make(chan struct{}), canceled: make(chan struct{})} + runner := NewSessionToolRunner(context.Background(), &runnerTestAPI{}, "session-id", SessionToolRunnerOptions{ + CustomTools: map[string]toolset.Tool{activeTool.Name(): activeTool, queuedTool.Name(): queuedTool}, + Logger: log.New(io.Discard, "", 0), + ToolTimeout: time.Second, + }) + runner.events = make(chan ToolCallResult, 2) + state := newRunnerTestState(runner) + events := []Event{ + {ID: "active-call", Type: EventTypeAgentCustomToolUse, Name: activeTool.Name(), ToolUseID: "active-call", SessionThreadID: "thread-a", Input: RawJSON(`{}`)}, + {ID: "queued-call", Type: EventTypeAgentCustomToolUse, Name: queuedTool.Name(), ToolUseID: "queued-call", SessionThreadID: "thread-b", Input: RawJSON(`{}`)}, + } + if err := state.processListedEvents(context.Background(), events, false); err != nil { + t.Fatal(err) + } + select { + case <-activeTool.started: + case <-time.After(time.Second): + t.Fatal("first tool did not start") + } + select { + case <-queuedTool.started: + t.Fatal("queued tool started before active tool finished") + case <-time.After(50 * time.Millisecond): + } + state.handleInterrupt(Event{Type: EventTypeUserInterrupt, SessionThreadID: "thread-b"}) + if !state.isAnswered("queued-call") || state.isAnswered("active-call") { + t.Fatalf("answered=%v", state.answered) + } + select { + case <-activeTool.canceled: + t.Fatal("interrupt for queued thread canceled active tool") + default: + } + state.handleInterrupt(Event{Type: EventTypeUserInterrupt}) + finishRunnerTestExecution(t, state) + select { + case <-queuedTool.started: + t.Fatal("interrupted queued tool was dispatched") + default: + } +} + +func TestSessionToolRunnerListEndTurnArmsIdleAfterToolCompletes(t *testing.T) { + maxIdle := time.Minute + tool := &runnerTestTool{} + runner := NewSessionToolRunner(context.Background(), &runnerTestAPI{}, "session-id", SessionToolRunnerOptions{ + CustomTools: map[string]toolset.Tool{tool.Name(): tool}, + Logger: log.New(io.Discard, "", 0), + MaxIdle: &maxIdle, + }) + runner.events = make(chan ToolCallResult, 1) + state := newRunnerTestState(runner) + toolUse := Event{ID: "idle-call", Type: EventTypeAgentCustomToolUse, Name: tool.Name(), ToolUseID: "idle-call", Input: RawJSON(`{}`)} + if err := state.processListedEvents(context.Background(), []Event{toolUse}, false); err != nil { + t.Fatal(err) + } + idle := Event{ID: "idle-end-turn", Type: EventTypeSessionStatusIdle, StopReason: &SessionStopReason{Type: SessionStopReasonEndTurn}} + if err := state.processListedEvents(context.Background(), []Event{toolUse, idle}, false); err != nil { + t.Fatal(err) + } + if !state.idleArmPending || !state.idleArmedAt.IsZero() { + t.Fatalf("pending=%v armed_at=%v", state.idleArmPending, state.idleArmedAt) + } + finishRunnerTestExecution(t, state) + if state.idleArmPending || state.idleArmedAt.IsZero() { + t.Fatalf("pending=%v armed_at=%v", state.idleArmPending, state.idleArmedAt) + } +} + +func TestSessionToolRunnerListReplayOfAnonymousInterruptDoesNotCancelLaterToolUse(t *testing.T) { + tool := &runnerBlockingTool{started: make(chan struct{}), canceled: make(chan struct{})} + api := &runnerTestAPI{} + runner := NewSessionToolRunner(context.Background(), api, "session-id", SessionToolRunnerOptions{ + CustomTools: map[string]toolset.Tool{tool.Name(): tool}, + Logger: log.New(io.Discard, "", 0), + ToolTimeout: time.Second, + }) + runner.events = make(chan ToolCallResult, 1) + state := newRunnerTestState(runner) + events := []Event{ + { + ID: "old-call", + Type: EventTypeAgentCustomToolUse, + Name: tool.Name(), + ToolUseID: "old-call", + SessionThreadID: "thread-id", + Input: RawJSON(`{}`), + }, + {Type: EventTypeUserInterrupt, SessionThreadID: "thread-id"}, + { + ID: "new-call", + Type: EventTypeAgentCustomToolUse, + Name: tool.Name(), + ToolUseID: "new-call", + SessionThreadID: "thread-id", + Input: RawJSON(`{}`), + }, + } + if err := state.processListedEvents(context.Background(), events, false); err != nil { + t.Fatal(err) + } + select { + case <-tool.started: + case <-time.After(time.Second): + t.Fatal("later tool did not start") + } + if err := state.processListedEvents(context.Background(), events, false); err != nil { + t.Fatal(err) + } + if err := state.processListedEvents(context.Background(), events, true); err != nil { + t.Fatal(err) + } + if state.isAnswered("new-call") { + t.Fatal("replayed interrupt canceled a later tool call") + } + if _, retained := state.toolUseEvents["old-call"]; retained { + t.Fatal("answered tool event was retained") + } + select { + case <-tool.canceled: + t.Fatal("later tool context was canceled by replay") + default: + } + + state.handleInterrupt(Event{Type: EventTypeUserInterrupt}) + finishRunnerTestExecution(t, state) +} + +func newRunnerTestState(runner *SessionToolRunner) *toolRunnerState { + return &toolRunnerState{ + runner: runner, + processed: map[string]bool{}, + seen: map[string]bool{}, + answered: map[string]bool{}, + pendingResults: map[string]Event{}, + recoveredResults: map[string]bool{}, + pendingAsk: map[string]Event{}, + confirmations: map[string]Event{}, + externalTools: map[string]Event{}, + toolUseEvents: map[string]Event{}, + scheduled: map[string]bool{}, + executionDone: make(chan toolExecutionResult, sessionRunnerResultsBuffer), + sessionToolUses: map[string]bool{}, + toolUsesSinceStatus: map[string]bool{}, + blockingEventIDs: map[string]bool{}, + } +} + +func finishRunnerTestExecution(t *testing.T, state *toolRunnerState) { + t.Helper() + select { + case result := <-state.executionDone: + if err := state.finishToolExecution(context.Background(), result); err != nil { + t.Fatal(err) + } + case <-time.After(time.Second): + t.Fatal("tool execution did not finish") + } +} + func TestSessionToolRunnerResendsRecoveredCurrentBlockerWithoutExecuting(t *testing.T) { tool := &runnerTestTool{} api := &runnerTestAPI{} diff --git a/arkruntime/selfhosted/types.go b/arkruntime/selfhosted/types.go index 4b0cf4e..e6a56ed 100644 --- a/arkruntime/selfhosted/types.go +++ b/arkruntime/selfhosted/types.go @@ -23,6 +23,8 @@ const ( EventTypeAgentCustomToolUse = "agent.custom_tool_use" // EventTypeUserToolConfirmation 表示用户确认工具执行。 EventTypeUserToolConfirmation = "user.tool_confirmation" + // EventTypeUserInterrupt 表示用户中断当前 session 或指定 thread。 + EventTypeUserInterrupt = "user.interrupt" // EventTypeUserToolResult 表示 self-host worker 回写内置工具结果。 EventTypeUserToolResult = "user.tool_result" // EventTypeUserCustomToolResult 表示 self-host worker 回写自定义工具结果。 diff --git a/arkruntime/toolset/file.go b/arkruntime/toolset/file.go index a0b5836..8e1edb7 100644 --- a/arkruntime/toolset/file.go +++ b/arkruntime/toolset/file.go @@ -4,20 +4,40 @@ package toolset import ( "context" + "encoding/base64" "encoding/json" + "errors" "fmt" + "io" + "mime" + "net/http" "os" "path/filepath" "strings" "unicode/utf8" ) +const ( + readBlockTypeImage = "image" + readBlockTypeDocument = "document" + readMediaTypePDF = "application/pdf" +) + // ReadTool 实现 read 工具。 type ReadTool struct { resolver *Resolver limits Limits } +type readRequest struct { + FilePath string `json:"file_path"` + Path string `json:"path"` + File string `json:"file"` + ViewRange []int `json:"view_range,omitempty"` + Offset *int `json:"offset,omitempty"` + Limit *int `json:"limit,omitempty"` +} + // NewReadTool 创建 read 工具。 func NewReadTool(resolver *Resolver, limits Limits) *ReadTool { return &ReadTool{resolver: resolver, limits: limits} @@ -27,17 +47,19 @@ func NewReadTool(resolver *Resolver, limits Limits) *ReadTool { func (t *ReadTool) Name() string { return "read" } // Execute 执行 read。 -func (t *ReadTool) Execute(_ context.Context, input json.RawMessage) Result { - var req struct { - FilePath string `json:"file_path"` - ViewRange []int `json:"view_range,omitempty"` - Offset int `json:"offset,omitempty"` - Limit int `json:"limit,omitempty"` - } +func (t *ReadTool) Execute(ctx context.Context, input json.RawMessage) Result { + var req readRequest if err := decodeInput(input, &req); err != nil { return ErrorResult(err.Error()) } - host, err := t.resolver.ResolveExisting(req.FilePath) + path, err := req.resolvedPath() + if err != nil { + return ErrorResult(err.Error()) + } + if len(req.ViewRange) > 0 && (req.Offset != nil || req.Limit != nil) { + return ErrorResult("view_range cannot be combined with offset or limit") + } + host, err := t.resolver.ResolveExisting(path) if err != nil { return ErrorResult(err.Error()) } @@ -48,6 +70,56 @@ func (t *ReadTool) Execute(_ context.Context, input json.RawMessage) Result { if !info.Mode().IsRegular() { return ErrorResult("path is not a regular file") } + if err := ctx.Err(); err != nil { + return ErrorResult(err.Error()) + } + blockType, mediaType, err := detectReadMedia(host) + if err != nil { + return ErrorResult(err.Error()) + } + if blockType != "" { + if len(req.ViewRange) > 0 || req.Offset != nil || req.Limit != nil { + return ErrorResult("view_range, offset, and limit are only supported for text files") + } + limit := t.limits.MaxMediaFileBytes + if limit == 0 { + limit = t.limits.MaxInputFileBytes + } + if limit > 0 && info.Size() > limit { + return ErrorResult(fmt.Sprintf("media file too large: %d bytes", info.Size())) + } + reader, err := os.Open(host) + if err != nil { + return ErrorResult(err.Error()) + } + defer func() { _ = reader.Close() }() + var source io.Reader = reader + if limit > 0 { + source = io.LimitReader(reader, limit+1) + } + data, err := io.ReadAll(source) + if err != nil { + return ErrorResult(err.Error()) + } + if limit > 0 && int64(len(data)) > limit { + size := info.Size() + if int64(len(data)) > size { + size = int64(len(data)) + } + return ErrorResult(fmt.Sprintf("media file too large: %d bytes", size)) + } + if err := ctx.Err(); err != nil { + return ErrorResult(err.Error()) + } + return Result{Content: []ContentBlock{{ + Type: blockType, + Source: map[string]any{ + "type": "base64", + "media_type": mediaType, + "data": base64.StdEncoding.EncodeToString(data), + }, + }}} + } if t.limits.MaxInputFileBytes > 0 && info.Size() > t.limits.MaxInputFileBytes { return ErrorResult(fmt.Sprintf("file too large: %d bytes", info.Size())) } @@ -63,9 +135,9 @@ func (t *ReadTool) Execute(_ context.Context, input json.RawMessage) Result { return ErrorResult("view_range must be [start_line, end_line]") } lines := strings.Split(string(data), "\n") - start := req.ViewRange[0] - 1 - if start < 0 { - start = 0 + start := 0 + if req.ViewRange[0] > 1 { + start = req.ViewRange[0] - 1 } if start >= len(lines) { return TextResult("") @@ -79,36 +151,39 @@ func (t *ReadTool) Execute(_ context.Context, input json.RawMessage) Result { } return TextResult(strings.Join(lines[start:end], "\n")) } - if req.Offset == 0 && req.Limit == 0 { - return TextResult(string(data)) - } lines := strings.Split(string(data), "\n") - if len(lines) > 0 && lines[len(lines)-1] == "" { + if len(lines) > 0 && lines[len(lines)-1] == "" && strings.HasSuffix(string(data), "\n") { lines = lines[:len(lines)-1] } - offset := req.Offset - if offset < 0 { - offset = 0 + offset := 0 + if req.Offset != nil { + if *req.Offset < 1 { + return ErrorResult(fmt.Sprintf("offset is the 1-based start line and must be >= 1, got %d", *req.Offset)) + } + offset = *req.Offset - 1 } if offset > len(lines) { offset = len(lines) } - limit := req.Limit + limit := 0 + if req.Limit != nil { + limit = *req.Limit + } if limit <= 0 { limit = t.limits.ReadDefaultLines } if limit <= 0 { limit = 2000 } - end := offset + limit - if end > len(lines) { - end = len(lines) + end := len(lines) + if limit < len(lines)-offset { + end = offset + limit } var b strings.Builder for i := offset; i < end; i++ { line := lines[i] if t.limits.ReadMaxLineChars > 0 && utf8.RuneCountInString(line) > t.limits.ReadMaxLineChars { - line = string([]rune(line)[:t.limits.ReadMaxLineChars]) + line = string([]rune(line)[:t.limits.ReadMaxLineChars]) + " [line truncated]" } fmt.Fprintf(&b, "%6d\t%s\n", i+1, line) } @@ -118,6 +193,62 @@ func (t *ReadTool) Execute(_ context.Context, input json.RawMessage) Result { return TextResult(b.String()) } +func (r readRequest) resolvedPath() (string, error) { + path := "" + for _, candidate := range []string{r.FilePath, r.Path, r.File} { + if candidate == "" { + continue + } + if path != "" && path != candidate { + return "", errors.New("file_path, path, and file must not conflict") + } + path = candidate + } + if path == "" { + return "", errors.New("file_path is required") + } + return path, nil +} + +func detectReadMedia(path string) (string, string, error) { + f, err := os.Open(path) + if err != nil { + return "", "", err + } + defer func() { _ = f.Close() }() + + header := make([]byte, 512) + n, err := io.ReadFull(f, header) + if err != nil && err != io.EOF && err != io.ErrUnexpectedEOF { + return "", "", err + } + header = header[:n] + detected := http.DetectContentType(header) + if blockType := supportedReadMediaBlock(detected); blockType != "" { + return blockType, detected, nil + } + if detected != "application/octet-stream" { + return "", "", nil + } + extensionType := mime.TypeByExtension(strings.ToLower(filepath.Ext(path))) + if semi := strings.IndexByte(extensionType, ';'); semi >= 0 { + extensionType = extensionType[:semi] + } + extensionType = strings.ToLower(strings.TrimSpace(extensionType)) + return supportedReadMediaBlock(extensionType), extensionType, nil +} + +func supportedReadMediaBlock(mediaType string) string { + switch mediaType { + case "image/jpeg", "image/png", "image/gif", "image/webp": + return readBlockTypeImage + case readMediaTypePDF: + return readBlockTypeDocument + default: + return "" + } +} + // WriteTool 实现 write 工具。 type WriteTool struct { resolver *Resolver diff --git a/arkruntime/toolset/toolset_test.go b/arkruntime/toolset/toolset_test.go index 60e9bf8..60d02ca 100644 --- a/arkruntime/toolset/toolset_test.go +++ b/arkruntime/toolset/toolset_test.go @@ -4,6 +4,7 @@ package toolset import ( "context" + "encoding/base64" "os" "path/filepath" "strings" @@ -187,7 +188,7 @@ func TestReadSupportsViewRange(t *testing.T) { } } -func TestReadDefaultsToRawContent(t *testing.T) { +func TestReadDefaultsToNumberedContent(t *testing.T) { root := t.TempDir() if err := os.WriteFile(filepath.Join(root, "demo.txt"), []byte("a\nb\n"), 0o600); err != nil { t.Fatal(err) @@ -199,7 +200,165 @@ func TestReadDefaultsToRawContent(t *testing.T) { defer set.Close() res := set.Execute(context.Background(), "read", []byte(`{"file_path":"demo.txt"}`)) - if res.IsError || res.Content[0].Text != "a\nb\n" { + if res.IsError || res.Content[0].Text != " 1\ta\n 2\tb\n" { + t.Fatalf("read result = %+v", res) + } +} + +func TestReadMarksTruncatedLines(t *testing.T) { + root := t.TempDir() + if err := os.WriteFile(filepath.Join(root, "long.txt"), []byte(strings.Repeat("a", 2001)), 0o600); err != nil { + t.Fatal(err) + } + set, err := NewDefault(Options{Workdir: root}) + if err != nil { + t.Fatal(err) + } + defer set.Close() + + res := set.Execute(context.Background(), "read", []byte(`{"file_path":"long.txt"}`)) + want := " 1\t" + strings.Repeat("a", 2000) + " [line truncated]\n" + if res.IsError || res.Content[0].Text != want { + t.Fatalf("read result = %+v", res) + } +} + +func TestReadSupportsManagedAgentLineRange(t *testing.T) { + root := t.TempDir() + if err := os.WriteFile(filepath.Join(root, "demo.txt"), []byte("a\nb\nc\nd\n"), 0o600); err != nil { + t.Fatal(err) + } + set, err := NewDefault(Options{Workdir: root}) + if err != nil { + t.Fatal(err) + } + defer set.Close() + + res := set.Execute(context.Background(), "read", []byte(`{"file_path":"demo.txt","offset":2,"limit":2}`)) + if res.IsError || res.Content[0].Text != " 2\tb\n 3\tc\n\n[truncated: showing lines 2-3 of 4]\n" { + t.Fatalf("read result = %+v", res) + } + invalid := set.Execute(context.Background(), "read", []byte(`{"file_path":"demo.txt","offset":0}`)) + if !invalid.IsError || !strings.Contains(invalid.Content[0].Text, "must be >= 1") { + t.Fatalf("invalid offset result = %+v", invalid) + } +} + +func TestReadKeepsLegacyPathAliases(t *testing.T) { + root := t.TempDir() + if err := os.WriteFile(filepath.Join(root, "demo.txt"), []byte("a\n"), 0o600); err != nil { + t.Fatal(err) + } + set, err := NewDefault(Options{Workdir: root}) + if err != nil { + t.Fatal(err) + } + defer set.Close() + + for _, input := range []string{ + `{"path":"demo.txt"}`, + `{"file":"demo.txt"}`, + `{"file_path":"demo.txt","path":"demo.txt"}`, + } { + res := set.Execute(context.Background(), "read", []byte(input)) + if res.IsError { + t.Fatalf("read(%s) result = %+v", input, res) + } + } + conflict := set.Execute(context.Background(), "read", []byte(`{"file_path":"demo.txt","path":"other.txt"}`)) + if !conflict.IsError || !strings.Contains(conflict.Content[0].Text, "must not conflict") { + t.Fatalf("conflicting path result = %+v", conflict) + } + rangeConflict := set.Execute(context.Background(), "read", []byte( + `{"file_path":"demo.txt","view_range":[1,1],"offset":1}`, + )) + if !rangeConflict.IsError || !strings.Contains(rangeConflict.Content[0].Text, "cannot be combined") { + t.Fatalf("conflicting range result = %+v", rangeConflict) + } +} + +func TestReadReturnsImageAsBase64ContentBlock(t *testing.T) { + root := t.TempDir() + data := append([]byte("\x89PNG\r\n\x1a\n"), make([]byte, 300<<10)...) + if err := os.WriteFile(filepath.Join(root, "image.png"), data, 0o600); err != nil { + t.Fatal(err) + } + set, err := NewDefault(Options{Workdir: root}) + if err != nil { + t.Fatal(err) + } + defer set.Close() + + res := set.Execute(context.Background(), "read", []byte(`{"file_path":"image.png"}`)) + if res.IsError || len(res.Content) != 1 || res.Content[0].Type != "image" { + t.Fatalf("read result = %+v", res) + } + source, ok := res.Content[0].Source.(map[string]any) + if !ok || source["type"] != "base64" || source["media_type"] != "image/png" { + t.Fatalf("image source = %#v", res.Content[0].Source) + } + if source["data"] != base64.StdEncoding.EncodeToString(data) { + t.Fatal("image data was not base64 encoded") + } +} + +func TestReadReturnsPDFAsDocumentContentBlock(t *testing.T) { + root := t.TempDir() + data := []byte("%PDF-1.7\n%%EOF\n") + if err := os.WriteFile(filepath.Join(root, "report.pdf"), data, 0o600); err != nil { + t.Fatal(err) + } + set, err := NewDefault(Options{Workdir: root}) + if err != nil { + t.Fatal(err) + } + defer set.Close() + + res := set.Execute(context.Background(), "read", []byte(`{"file_path":"report.pdf"}`)) + if res.IsError || len(res.Content) != 1 || res.Content[0].Type != "document" { + t.Fatalf("read result = %+v", res) + } + source, ok := res.Content[0].Source.(map[string]any) + if !ok || source["media_type"] != "application/pdf" || source["data"] != base64.StdEncoding.EncodeToString(data) { + t.Fatalf("document source = %#v", res.Content[0].Source) + } +} + +func TestReadRejectsRangesForMedia(t *testing.T) { + root := t.TempDir() + if err := os.WriteFile(filepath.Join(root, "image.gif"), []byte("GIF89a"), 0o600); err != nil { + t.Fatal(err) + } + set, err := NewDefault(Options{Workdir: root}) + if err != nil { + t.Fatal(err) + } + defer set.Close() + + res := set.Execute(context.Background(), "read", []byte(`{"file_path":"image.gif","view_range":[1,2]}`)) + if !res.IsError || !strings.Contains(res.Content[0].Text, "only supported for text") { + t.Fatalf("read result = %+v", res) + } +} + +func TestReadRejectsMediaAboveInlineLimit(t *testing.T) { + root := t.TempDir() + path := filepath.Join(root, "image.png") + if err := os.WriteFile(path, []byte("\x89PNG\r\n\x1a\n"), 0o600); err != nil { + t.Fatal(err) + } + limits := DefaultLimits() + if err := os.Truncate(path, limits.MaxMediaFileBytes+1); err != nil { + t.Fatal(err) + } + set, err := NewDefault(Options{Workdir: root, Limits: limits}) + if err != nil { + t.Fatal(err) + } + defer set.Close() + + res := set.Execute(context.Background(), "read", []byte(`{"file_path":"image.png"}`)) + if !res.IsError || !strings.Contains(res.Content[0].Text, "media file too large") { t.Fatalf("read result = %+v", res) } } diff --git a/arkruntime/toolset/types.go b/arkruntime/toolset/types.go index d9fb03c..c29fc4c 100644 --- a/arkruntime/toolset/types.go +++ b/arkruntime/toolset/types.go @@ -56,6 +56,8 @@ type Options struct { type Limits struct { MaxOutputBytes int64 MaxInputFileBytes int64 + // MaxMediaFileBytes 限制 read 内联图片和 PDF 的原始文件大小;零值回落到 MaxInputFileBytes,负值不限制。 + MaxMediaFileBytes int64 ReadDefaultLines int ReadMaxLineChars int GlobMaxMatches int @@ -67,6 +69,7 @@ func DefaultLimits() Limits { return Limits{ MaxOutputBytes: 100 << 10, MaxInputFileBytes: 256 << 10, + MaxMediaFileBytes: 7 << 20, ReadDefaultLines: 2000, ReadMaxLineChars: 2000, GlobMaxMatches: 1000,