diff --git a/packages/opencode/src/flag/flag.ts b/packages/opencode/src/flag/flag.ts index 1ece6b10cbe..fbebfb3b03a 100644 --- a/packages/opencode/src/flag/flag.ts +++ b/packages/opencode/src/flag/flag.ts @@ -35,7 +35,6 @@ export namespace Flag { export const KILO_SERVER_PASSWORD = process.env["KILO_SERVER_PASSWORD"] export const KILO_SERVER_USERNAME = process.env["KILO_SERVER_USERNAME"] export const KILO_ENABLE_QUESTION_TOOL = truthy("KILO_ENABLE_QUESTION_TOOL") - // Experimental export const KILO_EXPERIMENTAL = truthy("KILO_EXPERIMENTAL") export const KILO_EXPERIMENTAL_FILEWATCHER = truthy("KILO_EXPERIMENTAL_FILEWATCHER") @@ -58,6 +57,7 @@ export namespace Flag { export const KILO_MODELS_URL = process.env["KILO_MODELS_URL"] export const KILO_MODELS_PATH = process.env["KILO_MODELS_PATH"] export const KILO_SKIP_MIGRATIONS = truthy("KILO_SKIP_MIGRATIONS") + export declare const KILO_SESSION_RETRY_LIMIT: number | undefined function number(key: string) { const value = process.env[key] @@ -110,3 +110,17 @@ Object.defineProperty(Flag, "KILO_CLIENT", { enumerable: true, configurable: false, }) + +// Dynamic getter for KILO_SESSION_RETRY_LIMIT +// This must be evaluated at access time, not module load time, +// because tests and runtime tooling may set this env var later +Object.defineProperty(Flag, "KILO_SESSION_RETRY_LIMIT", { + get() { + const value = process.env["KILO_SESSION_RETRY_LIMIT"] + if (!value) return undefined + const parsed = Number(value) + return Number.isInteger(parsed) && parsed > 0 ? parsed : undefined + }, + enumerable: true, + configurable: false, +}) diff --git a/packages/opencode/src/session/processor.ts b/packages/opencode/src/session/processor.ts index 585facce3d4..81b50c7b244 100644 --- a/packages/opencode/src/session/processor.ts +++ b/packages/opencode/src/session/processor.ts @@ -16,6 +16,7 @@ import { SessionCompaction } from "./compaction" import { PermissionNext } from "@/permission/next" import { Question } from "@/question" import { Telemetry } from "@kilocode/kilo-telemetry" // kilocode_change +import { Flag } from "@/flag/flag" // kilocode_change export namespace SessionProcessor { const DOOM_LOOP_THRESHOLD = 3 @@ -387,7 +388,11 @@ export namespace SessionProcessor { }) } else { const retry = SessionRetry.retryable(error) - if (retry !== undefined) { + if ( + retry !== undefined && + (Flag.KILO_SESSION_RETRY_LIMIT === undefined || attempt < Flag.KILO_SESSION_RETRY_LIMIT) + ) { + // kilocode_change attempt++ const delay = SessionRetry.delay(attempt, error.name === "APIError" ? error : undefined) SessionStatus.set(input.sessionID, { diff --git a/packages/opencode/test/kilocode/session-processor-retry-limit.test.ts b/packages/opencode/test/kilocode/session-processor-retry-limit.test.ts new file mode 100644 index 00000000000..4f90bd8e0d6 --- /dev/null +++ b/packages/opencode/test/kilocode/session-processor-retry-limit.test.ts @@ -0,0 +1,196 @@ +import { afterEach, describe, expect, mock, spyOn, test } from "bun:test" + +mock.module("@/kilo-sessions/remote-sender", () => ({ + RemoteSender: { + create() { + return { + queue() {}, + flush: async () => undefined, + } + }, + }, +})) + +import { APICallError } from "ai" +import { Bus } from "../../src/bus" +import { Flag } from "../../src/flag/flag" +import { Identifier } from "../../src/id/id" +import { Instance } from "../../src/project/instance" +import type { Provider } from "../../src/provider/provider" +import { LLM } from "../../src/session/llm" +import type { LLM as LLMType } from "../../src/session/llm" +import { MessageV2 } from "../../src/session/message-v2" +import { SessionRetry } from "../../src/session/retry" +import { SessionStatus } from "../../src/session/status" +import { Log } from "../../src/util/log" +import { tmpdir } from "../fixture/fixture" + +Log.init({ print: false }) + +function createModel(): Provider.Model { + return { + id: "gpt-4", + providerID: "openai", + name: "GPT-4", + limit: { + context: 128000, + input: 0, + output: 4096, + }, + cost: { input: 0, output: 0, cache: { read: 0, write: 0 } }, + capabilities: { + toolcall: true, + attachment: false, + reasoning: false, + temperature: true, + input: { text: true, image: false, audio: false, video: false }, + output: { text: true, image: false, audio: false, video: false }, + }, + api: { id: "openai", url: "https://api.openai.com/v1", npm: "@ai-sdk/openai" }, + options: {}, + headers: {}, + } as Provider.Model +} + +function retryable429() { + return new APICallError({ + message: "429 status code (no body)", + url: "https://api.openai.com/v1/chat/completions", + requestBodyValues: {}, + statusCode: 429, + responseHeaders: { "content-type": "application/json" }, + isRetryable: true, + }) +} + +function sentinel() { + return new Error("unexpected extra llm call") +} + +async function seed(model: Provider.Model) { + const { Session } = await import("../../src/session") + const session = await Session.create({}) + const user = (await Session.updateMessage({ + id: Identifier.ascending("message"), + role: "user", + sessionID: session.id, + time: { created: Date.now() }, + agent: "code", + model: { providerID: model.providerID, modelID: model.id }, + tools: {}, + })) as MessageV2.User + const assistant = (await Session.updateMessage({ + id: Identifier.ascending("message"), + parentID: user.id, + role: "assistant", + mode: "code", + agent: "code", + path: { + cwd: Instance.directory, + root: Instance.worktree, + }, + cost: 0, + tokens: { + input: 0, + output: 0, + reasoning: 0, + cache: { read: 0, write: 0 }, + }, + modelID: model.id, + providerID: model.providerID, + time: { created: Date.now() }, + sessionID: session.id, + })) as MessageV2.Assistant + return { assistant, session, user } +} + +function input(model: Provider.Model, sessionID: string, user: MessageV2.User): LLMType.StreamInput { + return { + user, + sessionID, + model, + agent: { name: "code", mode: "primary", permission: [], options: {} } as any, + system: [], + abort: AbortSignal.any([]), + messages: [], + tools: {}, + } +} + +afterEach(() => { + delete process.env.KILO_SESSION_RETRY_LIMIT +}) + +describe("session processor retry limit", () => { + test("stops after two retries with the normalized retryable error", async () => { + await using tmp = await tmpdir({ git: true }) + process.env.KILO_SESSION_RETRY_LIMIT = "2" + + await Instance.provide({ + directory: tmp.path, + fn: async () => { + const { Session } = await import("../../src/session") + const { SessionProcessor } = await import("../../src/session/processor") + const model = createModel() + const seeded = await seed(model) + const retry: number[] = [] + const errors: Array = [] + const unsubStatus = Bus.subscribe(SessionStatus.Event.Status, (event) => { + if (event.properties.sessionID !== seeded.session.id) return + if (event.properties.status.type !== "retry") return + retry.push(event.properties.status.attempt) + }) + const unsubError = Bus.subscribe(Session.Event.Error, (event) => { + if (event.properties.sessionID !== seeded.session.id) return + errors.push(event.properties.error) + }) + const llm = spyOn(LLM, "stream") + .mockRejectedValueOnce(retryable429()) + .mockRejectedValueOnce(retryable429()) + .mockRejectedValueOnce(retryable429()) + .mockRejectedValue(sentinel()) + const sleep = spyOn(SessionRetry, "sleep").mockResolvedValue(undefined) + const processor = SessionProcessor.create({ + assistantMessage: seeded.assistant, + sessionID: seeded.session.id, + model, + abort: AbortSignal.any([]), + }) + + try { + const result = await processor.process(input(model, seeded.session.id, seeded.user)) + const expected = MessageV2.fromError(retryable429(), { providerID: "openai" }) + + expect(result).toBe("stop") + expect(llm).toHaveBeenCalledTimes(3) + expect(sleep).toHaveBeenCalledTimes(2) + expect(retry).toStrictEqual([1, 2]) + expect(processor.message.error).toStrictEqual(expected) + expect(errors).toStrictEqual([expected]) + } finally { + unsubStatus() + unsubError() + llm.mockRestore() + sleep.mockRestore() + } + }, + }) + }) + + test("only positive integers enable the limit", () => { + delete process.env.KILO_SESSION_RETRY_LIMIT + expect(Flag.KILO_SESSION_RETRY_LIMIT).toBeUndefined() + + process.env.KILO_SESSION_RETRY_LIMIT = "0" + expect(Flag.KILO_SESSION_RETRY_LIMIT).toBeUndefined() + + process.env.KILO_SESSION_RETRY_LIMIT = "-1" + expect(Flag.KILO_SESSION_RETRY_LIMIT).toBeUndefined() + + process.env.KILO_SESSION_RETRY_LIMIT = "abc" + expect(Flag.KILO_SESSION_RETRY_LIMIT).toBeUndefined() + + process.env.KILO_SESSION_RETRY_LIMIT = "2" + expect(Flag.KILO_SESSION_RETRY_LIMIT).toBe(2) + }) +})