retry.test.ts 16 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440
  1. import { describe, expect, test } from "bun:test"
  2. import { LayerNode } from "@kirincode-ai/core/effect/layer-node"
  3. import { SessionV1 } from "@kirincode-ai/core/v1/session"
  4. import type { NamedError } from "@kirincode-ai/core/util/error"
  5. import { APICallError } from "ai"
  6. import { setTimeout as sleep } from "node:timers/promises"
  7. import { Effect, Schedule, Schema } from "effect"
  8. import { CrossSpawnSpawner } from "@kirincode-ai/core/cross-spawn-spawner"
  9. import { SessionRetry } from "../../src/session/retry"
  10. import { MessageV2 } from "../../src/session/message-v2"
  11. import { ProviderError } from "../../src/provider/error"
  12. import { SessionID } from "../../src/session/schema"
  13. import { SessionStatus } from "../../src/session/status"
  14. import { testEffect } from "../lib/effect"
  15. import { ProviderV2 } from "@kirincode-ai/core/provider"
  16. const providerID = ProviderV2.ID.make("test")
  17. const retryProvider = "test"
  18. const it = testEffect(LayerNode.compile(LayerNode.group([SessionStatus.node, CrossSpawnSpawner.node])))
  19. function apiError(headers?: Record<string, string>): SessionV1.APIError {
  20. return Schema.decodeUnknownSync(SessionV1.APIError.Schema)(
  21. new SessionV1.APIError({
  22. message: "boom",
  23. isRetryable: true,
  24. responseHeaders: headers,
  25. }).toObject(),
  26. )
  27. }
  28. function wrap(message: unknown): ReturnType<NamedError["toObject"]> {
  29. return { name: "", data: { message } }
  30. }
  31. describe("session.retry.delay", () => {
  32. test("caps delay at 30 seconds when headers missing", () => {
  33. const error = apiError()
  34. const delays = Array.from({ length: 10 }, (_, index) => SessionRetry.delay(index + 1, error))
  35. expect(delays).toStrictEqual([2000, 4000, 8000, 16000, 30000, 30000, 30000, 30000, 30000, 30000])
  36. })
  37. test("prefers retry-after-ms when shorter than exponential", () => {
  38. const error = apiError({ "retry-after-ms": "1500" })
  39. expect(SessionRetry.delay(4, error)).toBe(1500)
  40. })
  41. test("uses retry-after seconds when reasonable", () => {
  42. const error = apiError({ "retry-after": "30" })
  43. expect(SessionRetry.delay(3, error)).toBe(30000)
  44. })
  45. test("accepts http-date retry-after values", () => {
  46. const date = new Date(Date.now() + 20000).toUTCString()
  47. const error = apiError({ "retry-after": date })
  48. const d = SessionRetry.delay(1, error)
  49. expect(d).toBeGreaterThanOrEqual(19000)
  50. expect(d).toBeLessThanOrEqual(20000)
  51. })
  52. test("ignores invalid retry hints", () => {
  53. const error = apiError({ "retry-after": "not-a-number" })
  54. expect(SessionRetry.delay(1, error)).toBe(2000)
  55. })
  56. test("ignores malformed date retry hints", () => {
  57. const error = apiError({ "retry-after": "Invalid Date String" })
  58. expect(SessionRetry.delay(1, error)).toBe(2000)
  59. })
  60. test("ignores past date retry hints", () => {
  61. const pastDate = new Date(Date.now() - 5000).toUTCString()
  62. const error = apiError({ "retry-after": pastDate })
  63. expect(SessionRetry.delay(1, error)).toBe(2000)
  64. })
  65. test("uses retry-after values even when exceeding 10 minutes with headers", () => {
  66. const error = apiError({ "retry-after": "50" })
  67. expect(SessionRetry.delay(1, error)).toBe(50000)
  68. const longError = apiError({ "retry-after-ms": "700000" })
  69. expect(SessionRetry.delay(1, longError)).toBe(700000)
  70. })
  71. test("caps oversized header delays to the runtime timer limit", () => {
  72. const error = apiError({ "retry-after-ms": "999999999999" })
  73. expect(SessionRetry.delay(1, error)).toBe(SessionRetry.RETRY_MAX_DELAY)
  74. })
  75. it.instance("policy updates retry status and increments attempts", () =>
  76. Effect.gen(function* () {
  77. const sessionID = SessionID.make("session-retry-test")
  78. const error = apiError({ "retry-after-ms": "0" })
  79. const status = yield* SessionStatus.Service
  80. const step = yield* Schedule.toStepWithMetadata(
  81. SessionRetry.policy({
  82. provider: "test",
  83. parse: Schema.decodeUnknownSync(SessionV1.APIError.Schema),
  84. set: (info) =>
  85. status.set(sessionID, {
  86. type: "retry",
  87. attempt: info.attempt,
  88. message: info.message,
  89. next: info.next,
  90. }),
  91. }),
  92. )
  93. yield* step(error)
  94. yield* step(error)
  95. expect(yield* status.get(sessionID)).toMatchObject({
  96. type: "retry",
  97. attempt: 2,
  98. message: "boom",
  99. })
  100. }),
  101. )
  102. })
  103. describe("session.retry.retryable", () => {
  104. test("maps too_many_requests json messages", () => {
  105. const error = wrap(JSON.stringify({ type: "error", error: { type: "too_many_requests" } }))
  106. expect(SessionRetry.retryable(error, retryProvider)).toEqual({ message: "Too Many Requests" })
  107. })
  108. test("maps overloaded provider codes", () => {
  109. const error = wrap(JSON.stringify({ code: "resource_exhausted" }))
  110. expect(SessionRetry.retryable(error, retryProvider)).toEqual({ message: "Provider is overloaded" })
  111. })
  112. test("does not retry unknown json messages", () => {
  113. const error = wrap(JSON.stringify({ error: { message: "no_kv_space" } }))
  114. expect(SessionRetry.retryable(error, retryProvider)).toBeUndefined()
  115. })
  116. test("does not throw on numeric error codes", () => {
  117. const error = wrap(JSON.stringify({ type: "error", error: { code: 123 } }))
  118. const result = SessionRetry.retryable(error, retryProvider)
  119. expect(result).toBeUndefined()
  120. })
  121. test("returns undefined for non-json message", () => {
  122. const error = wrap("not-json")
  123. expect(SessionRetry.retryable(error, retryProvider)).toBeUndefined()
  124. })
  125. test("retries plain text rate limit errors from Alibaba", () => {
  126. const msg =
  127. "Upstream error from Alibaba: Request rate increased too quickly. To ensure system stability, please adjust your client logic to scale requests more smoothly over time."
  128. const error = wrap(msg)
  129. expect(SessionRetry.retryable(error, retryProvider)).toEqual({ message: msg })
  130. })
  131. test("retries plain text rate limit errors", () => {
  132. const msg = "Rate limit exceeded, please try again later"
  133. const error = wrap(msg)
  134. expect(SessionRetry.retryable(error, retryProvider)).toEqual({ message: msg })
  135. })
  136. test("retries too many requests in plain text", () => {
  137. const msg = "Too many requests, please slow down"
  138. const error = wrap(msg)
  139. expect(SessionRetry.retryable(error, retryProvider)).toEqual({ message: msg })
  140. })
  141. test("retries transport timeout errors", () => {
  142. const request = MessageV2.fromError(new ProviderError.HeaderTimeoutError(10000), { providerID })
  143. expect(SessionV1.APIError.isInstance(request)).toBe(true)
  144. expect(SessionRetry.retryable(request, retryProvider)).toEqual({
  145. message: "Provider response headers timed out after 10000ms",
  146. })
  147. })
  148. test("retries websocket stream transport errors", () => {
  149. const request = MessageV2.fromError(
  150. new ProviderError.ResponseStreamError("WebSocket closed before response.completed (code 1006: Connection ended)"),
  151. { providerID },
  152. )
  153. expect(SessionV1.APIError.isInstance(request)).toBe(true)
  154. expect(SessionRetry.retryable(request, retryProvider)).toEqual({
  155. message: "WebSocket closed before response.completed (code 1006: Connection ended)",
  156. })
  157. })
  158. test("does not retry context overflow errors", () => {
  159. const error = new SessionV1.ContextOverflowError({
  160. message: "Input exceeds context window of this model",
  161. responseBody: '{"error":{"code":"context_length_exceeded"}}',
  162. }).toObject()
  163. expect(SessionRetry.retryable(error, retryProvider)).toBeUndefined()
  164. })
  165. test("retries 500 errors even when isRetryable is false", () => {
  166. const error = Schema.decodeUnknownSync(SessionV1.APIError.Schema)(
  167. new SessionV1.APIError({
  168. message: "Internal server error",
  169. isRetryable: false,
  170. statusCode: 500,
  171. responseBody: '{"type":"api_error","message":"Internal server error"}',
  172. }).toObject(),
  173. )
  174. expect(SessionRetry.retryable(error, retryProvider)).toEqual({ message: "Internal server error" })
  175. })
  176. test("retries 502 bad gateway errors", () => {
  177. const error = Schema.decodeUnknownSync(SessionV1.APIError.Schema)(
  178. new SessionV1.APIError({
  179. message: "Bad gateway",
  180. isRetryable: false,
  181. statusCode: 502,
  182. }).toObject(),
  183. )
  184. expect(SessionRetry.retryable(error, retryProvider)).toEqual({ message: "Bad gateway" })
  185. })
  186. test("retries 503 service unavailable errors", () => {
  187. const error = Schema.decodeUnknownSync(SessionV1.APIError.Schema)(
  188. new SessionV1.APIError({
  189. message: "Service unavailable",
  190. isRetryable: false,
  191. statusCode: 503,
  192. }).toObject(),
  193. )
  194. expect(SessionRetry.retryable(error, retryProvider)).toEqual({ message: "Service unavailable" })
  195. })
  196. test("does not retry 4xx errors when isRetryable is false", () => {
  197. const error = Schema.decodeUnknownSync(SessionV1.APIError.Schema)(
  198. new SessionV1.APIError({
  199. message: "Bad request",
  200. isRetryable: false,
  201. statusCode: 400,
  202. }).toObject(),
  203. )
  204. expect(SessionRetry.retryable(error, retryProvider)).toBeUndefined()
  205. })
  206. test("retries ZlibError decompression failures", () => {
  207. const error = Schema.decodeUnknownSync(SessionV1.APIError.Schema)(
  208. new SessionV1.APIError({
  209. message: "Response decompression failed",
  210. isRetryable: true,
  211. metadata: { code: "ZlibError" },
  212. }).toObject(),
  213. )
  214. const retryable = SessionRetry.retryable(error, retryProvider)
  215. expect(retryable).toBeDefined()
  216. expect(retryable).toEqual({ message: "Response decompression failed" })
  217. })
  218. test("maps free limits to Go upsell action", () => {
  219. const error = Schema.decodeUnknownSync(SessionV1.APIError.Schema)(
  220. new SessionV1.APIError({
  221. message: "Free usage exceeded",
  222. isRetryable: true,
  223. statusCode: 429,
  224. responseBody: JSON.stringify({
  225. type: "error",
  226. error: { type: "FreeUsageLimitError", message: "Free usage exceeded" },
  227. }),
  228. }).toObject(),
  229. )
  230. expect(SessionRetry.retryable(error, "kirincode")).toEqual({
  231. message: SessionRetry.GO_UPSELL_MESSAGE,
  232. action: {
  233. reason: "free_tier_limit",
  234. provider: "kirincode",
  235. title: "Free limit reached",
  236. message: "Subscribe to KirinCode Go for reliable access to the best open-source models, starting at $5/month.",
  237. label: "subscribe",
  238. link: SessionRetry.GO_UPSELL_URL,
  239. },
  240. })
  241. })
  242. test("maps Go subscription limits to workspace PAYG upsell", () => {
  243. const error = Schema.decodeUnknownSync(SessionV1.APIError.Schema)(
  244. new SessionV1.APIError({
  245. message: "Subscription quota exceeded. You can continue using free models.",
  246. isRetryable: true,
  247. statusCode: 429,
  248. responseHeaders: {
  249. "retry-after": "19380",
  250. },
  251. responseBody: JSON.stringify({
  252. type: "error",
  253. error: {
  254. type: "GoUsageLimitError",
  255. message: "Subscription quota exceeded. You can continue using free models.",
  256. },
  257. metadata: {
  258. workspace: "wrk_01K6XGM22R6FM8JVABE9XDQXGH",
  259. limitName: "5 hour",
  260. },
  261. }),
  262. }).toObject(),
  263. )
  264. expect(SessionRetry.retryable(error, "kirincode-go")).toEqual({
  265. message:
  266. "5 hour usage limit reached. It will reset in 5 hours 23 minutes. To continue using this model now, enable usage from your available balance - https://kirincode.ai/workspace/wrk_01K6XGM22R6FM8JVABE9XDQXGH/go",
  267. action: {
  268. reason: "account_rate_limit",
  269. provider: "kirincode-go",
  270. title: "Go limit reached",
  271. message:
  272. "5 hour usage limit reached. It will reset in 5 hours 23 minutes. To continue using this model now, enable usage from your available balance",
  273. label: "open settings",
  274. link: "https://kirincode.ai/workspace/wrk_01K6XGM22R6FM8JVABE9XDQXGH/go",
  275. },
  276. })
  277. })
  278. test("maps Go subscription limits without limit metadata", () => {
  279. const error = Schema.decodeUnknownSync(SessionV1.APIError.Schema)(
  280. new SessionV1.APIError({
  281. message: "Subscription quota exceeded. You can continue using free models.",
  282. isRetryable: true,
  283. statusCode: 429,
  284. responseHeaders: {
  285. "retry-after": "900",
  286. },
  287. responseBody: JSON.stringify({
  288. type: "error",
  289. error: {
  290. type: "GoUsageLimitError",
  291. message: "Subscription quota exceeded. You can continue using free models.",
  292. },
  293. metadata: {
  294. workspace: "wrk_01K6XGM22R6FM8JVABE9XDQXGH",
  295. },
  296. }),
  297. }).toObject(),
  298. )
  299. expect(SessionRetry.retryable(error, "kirincode-go")?.action?.message).toBe(
  300. "Usage limit reached. It will reset in 15 minutes. To continue using this model now, enable usage from your available balance",
  301. )
  302. })
  303. })
  304. describe("session.message-v2.fromError", () => {
  305. test.concurrent(
  306. "converts ECONNRESET socket errors to retryable APIError",
  307. async () => {
  308. using server = Bun.serve({
  309. port: 0,
  310. idleTimeout: 8,
  311. async fetch(_req) {
  312. return new Response(
  313. new ReadableStream({
  314. async pull(controller) {
  315. controller.enqueue("Hello,")
  316. await sleep(10000)
  317. controller.enqueue(" World!")
  318. controller.close()
  319. },
  320. }),
  321. { headers: { "Content-Type": "text/plain" } },
  322. )
  323. },
  324. })
  325. const error = await fetch(new URL("/", server.url.origin))
  326. .then((res) => res.text())
  327. .catch((e) => e)
  328. const result = MessageV2.fromError(error, { providerID })
  329. expect(SessionV1.APIError.isInstance(result)).toBe(true)
  330. if (!SessionV1.APIError.isInstance(result)) throw new Error("expected APIError")
  331. expect(result.data.isRetryable).toBe(true)
  332. expect(result.data.message).toBe("Connection reset by server")
  333. expect(result.data.metadata?.code).toBe("ECONNRESET")
  334. expect(result.data.metadata?.message).toInclude("socket connection")
  335. },
  336. 15_000,
  337. )
  338. test("ECONNRESET socket error is retryable", () => {
  339. const error = Schema.decodeUnknownSync(SessionV1.APIError.Schema)(
  340. new SessionV1.APIError({
  341. message: "Connection reset by server",
  342. isRetryable: true,
  343. metadata: { code: "ECONNRESET", message: "The socket connection was closed unexpectedly" },
  344. }).toObject(),
  345. )
  346. const retryable = SessionRetry.retryable(error, retryProvider)
  347. expect(retryable).toBeDefined()
  348. expect(retryable).toEqual({ message: "Connection reset by server" })
  349. })
  350. test("marks OpenAI 404 status codes as retryable", () => {
  351. const error = new APICallError({
  352. message: "boom",
  353. url: "https://api.openai.com/v1/chat/completions",
  354. requestBodyValues: {},
  355. statusCode: 404,
  356. responseHeaders: { "content-type": "application/json" },
  357. responseBody: '{"error":"boom"}',
  358. isRetryable: false,
  359. })
  360. const result = MessageV2.fromError(error, { providerID: ProviderV2.ID.make("openai") })
  361. if (!SessionV1.APIError.isInstance(result)) throw new Error("expected APIError")
  362. expect(result.data.isRetryable).toBe(true)
  363. })
  364. test("converts OpenAI server_error stream chunks to retryable APIError", () => {
  365. const result = MessageV2.fromError(
  366. {
  367. message: JSON.stringify({
  368. type: "error",
  369. sequence_number: 2,
  370. error: {
  371. type: "server_error",
  372. code: "server_error",
  373. message: "An error occurred while processing your request.",
  374. param: null,
  375. },
  376. }),
  377. },
  378. { providerID: ProviderV2.ID.make("openai") },
  379. )
  380. expect(SessionV1.APIError.isInstance(result)).toBe(true)
  381. if (!SessionV1.APIError.isInstance(result)) throw new Error("expected APIError")
  382. expect(result.data.isRetryable).toBe(true)
  383. expect(SessionRetry.retryable(result, retryProvider)).toEqual({
  384. message: "An error occurred while processing your request.",
  385. })
  386. })
  387. })