diff --git a/src/api/providers/__tests__/openrouter.spec.ts b/src/api/providers/__tests__/openrouter.spec.ts index 27a84363f6..cae83984ae 100644 --- a/src/api/providers/__tests__/openrouter.spec.ts +++ b/src/api/providers/__tests__/openrouter.spec.ts @@ -298,6 +298,32 @@ describe("OpenRouterHandler", () => { ) }) + it.each([ + { + modelId: "anthropic/claude-sonnet-4", + headers: { "x-anthropic-beta": "fine-grained-tool-streaming-2025-05-14" }, + }, + { modelId: "openai/gpt-4o", headers: undefined }, + ])("passes the abort signal with only the appropriate headers for $modelId", async ({ modelId, headers }) => { + const handler = new OpenRouterHandler({ ...mockOptions, openRouterModelId: modelId }) + const mockCreate = vitest.fn().mockResolvedValue(asyncStreamFrom([])) + Object.defineProperty(OpenAI.prototype, "chat", { + configurable: true, + value: { completions: { create: mockCreate } }, + }) + const controller = new AbortController() + + await collectStream( + handler.createMessage("test", [], { taskId: "test-task", abortSignal: controller.signal }), + ) + + expect(mockCreate).toHaveBeenCalledWith(expect.objectContaining({ model: modelId }), { + signal: controller.signal, + ...(headers ? { headers } : {}), + }) + expect(mockCreate.mock.calls[0][1].signal).toBe(controller.signal) + }) + it("adds cache control for supported models", async () => { const handler = new OpenRouterHandler( makeApiHandlerOptions({ diff --git a/src/api/providers/openrouter.ts b/src/api/providers/openrouter.ts index ed53c111b5..9c198a32a4 100644 --- a/src/api/providers/openrouter.ts +++ b/src/api/providers/openrouter.ts @@ -331,9 +331,15 @@ export class OpenRouterHandler extends BaseProvider implements SingleCompletionH } // Add Anthropic beta header for fine-grained tool streaming when using Anthropic models - const requestOptions = modelId.startsWith("anthropic/") - ? { headers: { "x-anthropic-beta": "fine-grained-tool-streaming-2025-05-14" } } - : undefined + const requestOptions = + modelId.startsWith("anthropic/") || metadata?.abortSignal + ? { + ...(modelId.startsWith("anthropic/") + ? { headers: { "x-anthropic-beta": "fine-grained-tool-streaming-2025-05-14" } } + : {}), + ...(metadata?.abortSignal ? { signal: metadata.abortSignal } : {}), + } + : undefined let stream try { diff --git a/src/core/task/ReasoningLoopDetector.ts b/src/core/task/ReasoningLoopDetector.ts new file mode 100644 index 0000000000..e479d62689 --- /dev/null +++ b/src/core/task/ReasoningLoopDetector.ts @@ -0,0 +1,56 @@ +const MIN_PATTERN_LENGTH = 80 +const MAX_PATTERN_LENGTH = 2_048 +const REQUIRED_REPETITIONS = 5 +const CHECK_INTERVAL = 64 +const BUFFER_LENGTH = MAX_PATTERN_LENGTH * (REQUIRED_REPETITIONS + 1) + +/** Detects sustained, exact repetition in streamed model reasoning. */ +export class ReasoningLoopDetector { + private buffer = "" + private uncheckedCharacters = 0 + + add(text: string): boolean { + if (!text) { + return false + } + + this.buffer = (this.buffer + text).slice(-BUFFER_LENGTH) + this.uncheckedCharacters += text.length + + if (this.uncheckedCharacters < CHECK_INTERVAL) { + return false + } + this.uncheckedCharacters = 0 + + const maxPatternLength = Math.min(MAX_PATTERN_LENGTH, Math.floor(this.buffer.length / REQUIRED_REPETITIONS)) + for (let patternLength = MIN_PATTERN_LENGTH; patternLength <= maxPatternLength; patternLength++) { + const patternStart = this.buffer.length - patternLength + const pattern = this.buffer.slice(patternStart) + let repeats = true + + for (let repetition = 2; repetition <= REQUIRED_REPETITIONS; repetition++) { + const start = this.buffer.length - patternLength * repetition + if (this.buffer.slice(start, start + patternLength) !== pattern) { + repeats = false + break + } + } + + if (repeats) { + return true + } + } + + return false + } +} + +export class RepetitiveReasoningError extends Error { + constructor() { + super( + "Repetitive reasoning detected. The model appears to be stuck in a loop, so Zoo Code stopped the request. Retry with less context or a different model.", + ) + this.name = "RepetitiveReasoningError" + Object.setPrototypeOf(this, RepetitiveReasoningError.prototype) + } +} diff --git a/src/core/task/Task.ts b/src/core/task/Task.ts index d5313f68cf..605a2bc94e 100644 --- a/src/core/task/Task.ts +++ b/src/core/task/Task.ts @@ -7,6 +7,7 @@ import EventEmitter from "events" import { AskIgnoredError } from "./AskIgnoredError" import { RateLimitClock, createRateLimitClock } from "./RateLimitClock" +import { ReasoningLoopDetector, RepetitiveReasoningError } from "./ReasoningLoopDetector" import { Anthropic } from "@anthropic-ai/sdk" import OpenAI from "openai" @@ -3267,6 +3268,7 @@ export class Task extends EventEmitter implements TaskLike { }) let assistantMessage = "" let reasoningMessage = "" + const reasoningLoopDetector = new ReasoningLoopDetector() const pendingGroundingSources: GroundingSource[] = [] this.isStreaming = true @@ -3303,6 +3305,11 @@ export class Task extends EventEmitter implements TaskLike { let item = await nextChunkWithAbort() while (!item.done) { const chunk = item.value + // Detect a loop before prefetching: the provider may stall on its next read. + if (chunk?.type === "reasoning" && reasoningLoopDetector.add(chunk.text)) { + this.cancelCurrentRequest() + throw new RepetitiveReasoningError() + } item = await nextChunkWithAbort() if (!chunk) { // Sometimes chunk is undefined, no idea that can cause @@ -3719,8 +3726,11 @@ export class Task extends EventEmitter implements TaskLike { // ??= keeps the first reason; a cancel can land during abortStream after cancelReason was already computed. this.abortReason ??= "user_cancelled" await this.abortTask() - } else if (error instanceof OutputTokenLimitError) { - // Truncation repeats on an identical request, so never auto-retry it + } else if ( + error instanceof OutputTokenLimitError || + error instanceof RepetitiveReasoningError + ) { + // These failures repeat on an identical request, so never auto-retry them // (even with auto-approval); let the user decide once. const { response } = await this.ask("api_req_failed", rawErrorMessage) @@ -3732,7 +3742,7 @@ export class Task extends EventEmitter implements TaskLike { stack.push({ userContent: currentUserContent, includeFileDetails: false, - retryAttempt: 0, + retryAttempt: (currentItem.retryAttempt ?? 0) + 1, }) continue } else { diff --git a/src/core/task/__tests__/ReasoningLoopDetector.spec.ts b/src/core/task/__tests__/ReasoningLoopDetector.spec.ts new file mode 100644 index 0000000000..3697f1bbee --- /dev/null +++ b/src/core/task/__tests__/ReasoningLoopDetector.spec.ts @@ -0,0 +1,27 @@ +import { ReasoningLoopDetector } from "../ReasoningLoopDetector" + +describe("ReasoningLoopDetector", () => { + it("detects the reported reasoning loop across arbitrary stream chunks", () => { + const detector = new ReasoningLoopDetector() + const cycle = "OK.\n\nHmm. Let me read them.\n\nOK.\n\nHmm. Let me just do it.\n\nLet me read the files.\n\n" + const output = `I will inspect the implementation first.\n${cycle.repeat(8)}` + let detected = false + + for (let offset = 0; offset < output.length; offset += 17) { + detected ||= detector.add(output.slice(offset, offset + 17)) + } + + expect(detected).toBe(true) + }) + + it("does not flag repeated short phrases or ordinary long reasoning", () => { + const detector = new ReasoningLoopDetector() + const reasoning = Array.from( + { length: 40 }, + (_, index) => + `Step ${index}: inspect file ${index}, compare its behavior, and record the distinct result. OK.\n`, + ).join("") + + expect(detector.add(reasoning)).toBe(false) + }) +}) diff --git a/src/core/task/__tests__/Task.spec.ts b/src/core/task/__tests__/Task.spec.ts index eae85cf9f2..b55823e181 100644 --- a/src/core/task/__tests__/Task.spec.ts +++ b/src/core/task/__tests__/Task.spec.ts @@ -558,23 +558,112 @@ describe("Cline", () => { expect(askSpy).toHaveBeenCalledWith("api_req_failed", expect.stringContaining("Output token limit reached")) }) - it("retries from a fresh attempt count when the user confirms", async () => { + it("advances the retry count without duplicating the user message when the user confirms", async () => { const task = await createAutoApprovedTask() + let retryHistory: ApiMessage[] | undefined vi.spyOn(task, "ask").mockResolvedValue({ response: "yesButtonClicked" } satisfies TaskAskResult) const attemptApiRequestSpy = vi .spyOn(task, "attemptApiRequest") .mockImplementationOnce(() => truncatedStream()) - .mockImplementationOnce(() => asyncStreamFrom([{ type: "text", text: "done" }])) + .mockImplementationOnce(() => { + retryHistory = structuredClone(task.apiConversationHistory) + return asyncStreamFrom([{ type: "text", text: "done" }]) + }) .mockImplementation(() => { throw new Error("stop after retry response") }) await task.recursivelyMakeClineRequests([{ type: "text", text: "long request" }]) - expect(attemptApiRequestSpy.mock.calls.slice(0, 2).map(([retryAttempt]) => retryAttempt)).toEqual([0, 0]) + expect(attemptApiRequestSpy.mock.calls.slice(0, 2).map(([retryAttempt]) => retryAttempt)).toEqual([0, 1]) + expect(retryHistory?.filter((message) => message.role === "user")).toHaveLength(1) }) }) + describe("repetitive reasoning mid-stream", () => { + it.each(["noButtonClicked", "yesButtonClicked"] as const)( + "cancels before a stalled read, then handles %s", + async (response) => { + vi.spyOn(mockProvider, "getState").mockResolvedValue( + providerStateWith({ autoApprovalEnabled: true, requestDelaySeconds: 0 }), + ) + const task = new Task({ + provider: mockProvider, + apiConfiguration: mockApiConfig, + task: "test task", + startTask: false, + }) + vi.spyOn(task.diffViewProvider, "reset").mockResolvedValue(undefined) + vi.spyOn(getTaskTestAccess(task), "safeEnsureModelFetched").mockResolvedValue(stubModelInfo) + vi.spyOn(getTaskTestAccess(task), "presentAssistantMessageSafe").mockImplementation(() => {}) + let answerFailure!: (answer: TaskAskResult) => void + const failureResponse = new Promise((resolve) => { + answerFailure = resolve + }) + const askSpy = vi + .spyOn(task, "ask") + .mockResolvedValue({ response: "noButtonClicked" } satisfies TaskAskResult) + .mockReturnValueOnce(failureResponse) + const cancelSpy = vi.spyOn(task, "cancelCurrentRequest") + const controller = new AbortController() + let releaseRead!: () => void + const stalledRead = new Promise((resolve) => { + releaseRead = resolve + }) + let retryHistory: ApiMessage[] | undefined + const cycle = + "OK.\n\nHmm. Let me read them.\n\nOK.\n\nHmm. Let me just do it.\n\nLet me read the files.\n\n" + const attemptApiRequestSpy = vi + .spyOn(task, "attemptApiRequest") + .mockImplementationOnce(async function* () { + task.currentRequestAbortController = controller + yield { type: "reasoning", text: cycle.repeat(8) } + await stalledRead + }) + .mockImplementationOnce(() => { + retryHistory = structuredClone(task.apiConversationHistory) + return asyncStreamFrom([{ type: "text", text: "done" }]) + }) + .mockImplementation(() => { + throw new Error("stop after retry response") + }) + + const request = task.recursivelyMakeClineRequests([{ type: "text", text: "long request" }]) + try { + // The provider's next read and the user's retry decision are still unreleased. + await vi.waitFor(() => expect(controller.signal.aborted).toBe(true)) + expect(cancelSpy).toHaveBeenCalledOnce() + expect(attemptApiRequestSpy).toHaveBeenCalledOnce() + await vi.waitFor(() => + expect(askSpy).toHaveBeenCalledWith( + "api_req_failed", + expect.stringContaining("Repetitive reasoning detected"), + ), + ) + } finally { + releaseRead() + answerFailure({ response }) + await request + } + + expect(cancelSpy).toHaveBeenCalledOnce() + if (response === "yesButtonClicked") { + expect(attemptApiRequestSpy.mock.calls.slice(0, 2).map(([retryAttempt]) => retryAttempt)).toEqual([ + 0, 1, + ]) + expect(retryHistory?.filter((message) => message.role === "user")).toHaveLength(1) + expect(retryHistory?.[0]).toMatchObject({ + role: "user", + content: expect.arrayContaining([{ type: "text", text: "long request" }]), + }) + } else { + expect(attemptApiRequestSpy).toHaveBeenCalledOnce() + expect(askSpy).toHaveBeenCalledOnce() + } + }, + ) + }) + describe("native tool-call request isolation", () => { it("keeps overlapping Task parser state scoped to each request", async () => { const firstTask = new Task({