processor-effect.test.ts 35 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067
  1. import { SessionV1 } from "@kirincode-ai/core/v1/session"
  2. import { Database } from "@kirincode-ai/core/database/database"
  3. import { LayerNode } from "@kirincode-ai/core/effect/layer-node"
  4. import { EventV2Bridge } from "@/event-v2-bridge"
  5. import { expect } from "bun:test"
  6. import { tool } from "ai"
  7. import { Cause, Effect, Exit, Fiber, Layer, Stream } from "effect"
  8. import path from "path"
  9. import z from "zod"
  10. import type { Agent } from "../../src/agent/agent"
  11. import { Provider } from "@/provider/provider"
  12. import { Session } from "@/session/session"
  13. import { LLM } from "../../src/session/llm"
  14. import { MessageV2 } from "../../src/session/message-v2"
  15. import { SessionProcessor } from "../../src/session/processor"
  16. import { MessageID, PartID, SessionID } from "../../src/session/schema"
  17. import { SessionStatus } from "../../src/session/status"
  18. import { SessionSummary } from "../../src/session/summary"
  19. import { CrossSpawnSpawner } from "@kirincode-ai/core/cross-spawn-spawner"
  20. import { provideTmpdirInstance, provideTmpdirServer } from "../fixture/fixture"
  21. import { testEffect } from "../lib/effect"
  22. import { raw, reply, TestLLMServer } from "../lib/llm-server"
  23. import { RuntimeFlags } from "@/effect/runtime-flags"
  24. import { ProviderV2 } from "@kirincode-ai/core/provider"
  25. import { ModelV2 } from "@kirincode-ai/core/model"
  26. import { SessionProjector } from "@kirincode-ai/core/session/projector"
  27. import { LLMEvent } from "@kirincode-ai/llm"
  28. const summary = Layer.succeed(
  29. SessionSummary.Service,
  30. SessionSummary.Service.of({
  31. summarize: () => Effect.void,
  32. diff: () => Effect.succeed([]),
  33. computeDiff: () => Effect.succeed([]),
  34. }),
  35. )
  36. const ref = {
  37. providerID: ProviderV2.ID.make("test"),
  38. modelID: ModelV2.ID.make("test-model"),
  39. }
  40. const cfg = {
  41. provider: {
  42. test: {
  43. name: "Test",
  44. id: "test",
  45. env: [],
  46. npm: "@ai-sdk/openai-compatible",
  47. models: {
  48. "test-model": {
  49. id: "test-model",
  50. name: "Test Model",
  51. attachment: false,
  52. reasoning: false,
  53. temperature: false,
  54. tool_call: true,
  55. release_date: "2025-01-01",
  56. limit: { context: 100000, output: 10000 },
  57. cost: { input: 0, output: 0 },
  58. options: {},
  59. },
  60. },
  61. options: {
  62. apiKey: "test-key",
  63. baseURL: "http://localhost:1/v1",
  64. },
  65. },
  66. },
  67. }
  68. function providerCfg(url: string) {
  69. return {
  70. ...cfg,
  71. provider: {
  72. ...cfg.provider,
  73. test: {
  74. ...cfg.provider.test,
  75. options: {
  76. ...cfg.provider.test.options,
  77. baseURL: url,
  78. },
  79. },
  80. },
  81. }
  82. }
  83. function agent(): Agent.Info {
  84. return {
  85. name: "build",
  86. mode: "primary",
  87. options: {},
  88. permission: [{ permission: "*", pattern: "*", action: "allow" }],
  89. }
  90. }
  91. function defer<T>() {
  92. let resolve!: (value: T | PromiseLike<T>) => void
  93. const promise = new Promise<T>((done) => {
  94. resolve = done
  95. })
  96. return { promise, resolve }
  97. }
  98. const waitFor = <A>(check: Effect.Effect<A | undefined>, message: string) =>
  99. Effect.gen(function* () {
  100. const stop = Date.now() + 500
  101. while (Date.now() < stop) {
  102. const value = yield* check
  103. if (value !== undefined) return value
  104. yield* Effect.sleep("10 millis")
  105. }
  106. return yield* Effect.fail(new Error(message))
  107. })
  108. const user = Effect.fn("TestSession.user")(function* (sessionID: SessionID, text: string) {
  109. const session = yield* Session.Service
  110. const msg = yield* session.updateMessage({
  111. id: MessageID.ascending(),
  112. role: "user",
  113. sessionID,
  114. agent: "build",
  115. model: ref,
  116. time: { created: Date.now() },
  117. })
  118. yield* session.updatePart({
  119. id: PartID.ascending(),
  120. messageID: msg.id,
  121. sessionID,
  122. type: "text",
  123. text,
  124. })
  125. return msg
  126. })
  127. const assistant = Effect.fn("TestSession.assistant")(function* (
  128. sessionID: SessionID,
  129. parentID: MessageID,
  130. root: string,
  131. ) {
  132. const session = yield* Session.Service
  133. const msg: SessionV1.Assistant = {
  134. id: MessageID.ascending(),
  135. role: "assistant",
  136. sessionID,
  137. mode: "build",
  138. agent: "build",
  139. path: { cwd: root, root },
  140. cost: 0,
  141. tokens: {
  142. total: 0,
  143. input: 0,
  144. output: 0,
  145. reasoning: 0,
  146. cache: { read: 0, write: 0 },
  147. },
  148. modelID: ref.modelID,
  149. providerID: ref.providerID,
  150. parentID,
  151. time: { created: Date.now() },
  152. finish: "end_turn",
  153. }
  154. yield* session.updateMessage(msg)
  155. return msg
  156. })
  157. const root = LayerNode.group([
  158. SessionProcessor.node,
  159. Session.node,
  160. SessionProjector.node,
  161. Provider.node,
  162. Database.node,
  163. EventV2Bridge.node,
  164. SessionStatus.node,
  165. CrossSpawnSpawner.node,
  166. ])
  167. const replacements = [
  168. [SessionSummary.node, summary],
  169. [RuntimeFlags.node, RuntimeFlags.layer({ experimentalEventSystem: true })],
  170. ] as const
  171. const env = LayerNode.compile(
  172. LayerNode.group([root, LayerNode.make({ service: TestLLMServer, layer: TestLLMServer.layer, deps: [] })]),
  173. replacements,
  174. )
  175. const it = testEffect(env)
  176. const providerErrorLLM = Layer.succeed(
  177. LLM.Service,
  178. LLM.Service.of({
  179. stream: () =>
  180. Stream.make(
  181. LLMEvent.stepStart({ index: 0 }),
  182. LLMEvent.toolInputStart({ id: "call-1", name: "lookup" }),
  183. LLMEvent.toolInputEnd({ id: "call-1", name: "lookup" }),
  184. LLMEvent.toolCall({ id: "call-1", name: "lookup", input: {}, providerExecuted: true }),
  185. LLMEvent.toolResult({
  186. id: "call-1",
  187. name: "lookup",
  188. result: { type: "error", value: "provider boom" },
  189. providerExecuted: true,
  190. }),
  191. LLMEvent.stepFinish({ index: 0, reason: "stop" }),
  192. LLMEvent.finish({ reason: "stop" }),
  193. ),
  194. }),
  195. )
  196. const providerErrorEnv = LayerNode.compile(root, [...replacements, [LLM.node, providerErrorLLM]])
  197. const itProviderError = testEffect(providerErrorEnv)
  198. const fragmentFailureLLM = Layer.succeed(
  199. LLM.Service,
  200. LLM.Service.of({
  201. stream: () =>
  202. Stream.make(
  203. LLMEvent.stepStart({ index: 0 }),
  204. LLMEvent.reasoningStart({ id: "reasoning-1" }),
  205. LLMEvent.reasoningDelta({ id: "reasoning-1", text: "thinking" }),
  206. LLMEvent.textStart({ id: "text-1" }),
  207. LLMEvent.textDelta({ id: "text-1", text: "partial" }),
  208. LLMEvent.providerError({ message: "provider boom" }),
  209. ),
  210. }),
  211. )
  212. const fragmentFailureEnv = LayerNode.compile(root, [...replacements, [LLM.node, fragmentFailureLLM]])
  213. const itFragmentFailure = testEffect(fragmentFailureEnv)
  214. const boot = Effect.fn("test.boot")(function* () {
  215. const processors = yield* SessionProcessor.Service
  216. const session = yield* Session.Service
  217. const provider = yield* Provider.Service
  218. return { processors, session, provider }
  219. })
  220. // ---------------------------------------------------------------------------
  221. // Tests
  222. // ---------------------------------------------------------------------------
  223. it.live("session.processor effect tests capture llm input cleanly", () =>
  224. provideTmpdirServer(
  225. ({ dir, llm }) =>
  226. Effect.gen(function* () {
  227. const database = yield* Database.Service
  228. const { processors, session, provider } = yield* boot()
  229. yield* llm.text("hello")
  230. const chat = yield* session.create({})
  231. const parent = yield* user(chat.id, "hi")
  232. const msg = yield* assistant(chat.id, parent.id, path.resolve(dir))
  233. const mdl = yield* provider.getModel(ref.providerID, ref.modelID)
  234. const handle = yield* processors.create({
  235. assistantMessage: msg,
  236. sessionID: chat.id,
  237. model: mdl,
  238. })
  239. const input = {
  240. user: {
  241. id: parent.id,
  242. sessionID: chat.id,
  243. role: "user",
  244. time: parent.time,
  245. agent: parent.agent,
  246. model: { providerID: ref.providerID, modelID: ref.modelID },
  247. } satisfies SessionV1.User,
  248. sessionID: chat.id,
  249. model: mdl,
  250. agent: agent(),
  251. system: [],
  252. messages: [{ role: "user", content: "hi" }],
  253. tools: {},
  254. } satisfies LLM.StreamInput
  255. const value = yield* handle.process(input)
  256. const parts = yield* MessageV2.parts(msg.id)
  257. const calls = yield* llm.calls
  258. expect(value).toBe("continue")
  259. expect(calls).toBe(1)
  260. expect(parts.some((part) => part.type === "text" && part.text === "hello")).toBe(true)
  261. }),
  262. { config: (url) => providerCfg(url) },
  263. ),
  264. )
  265. it.live("session.processor effect tests preserve text start time", () =>
  266. provideTmpdirServer(
  267. ({ dir, llm }) =>
  268. Effect.gen(function* () {
  269. const database = yield* Database.Service
  270. const gate = defer<void>()
  271. const { processors, session, provider } = yield* boot()
  272. yield* llm.push(
  273. raw({
  274. head: [
  275. {
  276. id: "chatcmpl-test",
  277. object: "chat.completion.chunk",
  278. choices: [{ delta: { role: "assistant" } }],
  279. },
  280. {
  281. id: "chatcmpl-test",
  282. object: "chat.completion.chunk",
  283. choices: [{ delta: { content: "hello" } }],
  284. },
  285. ],
  286. wait: gate.promise,
  287. tail: [
  288. {
  289. id: "chatcmpl-test",
  290. object: "chat.completion.chunk",
  291. choices: [{ delta: {}, finish_reason: "stop" }],
  292. },
  293. ],
  294. }),
  295. )
  296. const chat = yield* session.create({})
  297. const parent = yield* user(chat.id, "hi")
  298. const msg = yield* assistant(chat.id, parent.id, path.resolve(dir))
  299. const mdl = yield* provider.getModel(ref.providerID, ref.modelID)
  300. const handle = yield* processors.create({
  301. assistantMessage: msg,
  302. sessionID: chat.id,
  303. model: mdl,
  304. })
  305. const run = yield* handle
  306. .process({
  307. user: {
  308. id: parent.id,
  309. sessionID: chat.id,
  310. role: "user",
  311. time: parent.time,
  312. agent: parent.agent,
  313. model: { providerID: ref.providerID, modelID: ref.modelID },
  314. } satisfies SessionV1.User,
  315. sessionID: chat.id,
  316. model: mdl,
  317. agent: agent(),
  318. system: [],
  319. messages: [{ role: "user", content: "hi" }],
  320. tools: {},
  321. })
  322. .pipe(Effect.forkChild)
  323. yield* waitFor(
  324. MessageV2.parts(msg.id).pipe(
  325. Effect.map((parts) => parts.find((part): part is SessionV1.TextPart => part.type === "text")),
  326. Effect.provideService(Database.Service, database),
  327. ),
  328. "timed out waiting for text part",
  329. )
  330. yield* Effect.sleep("20 millis")
  331. gate.resolve()
  332. const exit = yield* Fiber.await(run)
  333. const text = (yield* MessageV2.parts(msg.id)).find((part): part is SessionV1.TextPart => part.type === "text")
  334. expect(Exit.isSuccess(exit)).toBe(true)
  335. expect(text?.text).toBe("hello")
  336. expect(text?.time?.start).toBeDefined()
  337. expect(text?.time?.end).toBeDefined()
  338. if (!text?.time?.start || !text.time.end) return
  339. expect(text.time.start).toBeLessThan(text.time.end)
  340. }),
  341. { config: (url) => providerCfg(url) },
  342. ),
  343. )
  344. it.live("session.processor effect tests stop after token overflow requests compaction", () =>
  345. provideTmpdirServer(
  346. ({ dir, llm }) =>
  347. Effect.gen(function* () {
  348. const database = yield* Database.Service
  349. const { processors, session, provider } = yield* boot()
  350. yield* llm.text("after", { usage: { input: 100, output: 0 } })
  351. const chat = yield* session.create({})
  352. const parent = yield* user(chat.id, "compact")
  353. const msg = yield* assistant(chat.id, parent.id, path.resolve(dir))
  354. const base = yield* provider.getModel(ref.providerID, ref.modelID)
  355. const mdl = { ...base, limit: { context: 20, output: 10 } }
  356. const handle = yield* processors.create({
  357. assistantMessage: msg,
  358. sessionID: chat.id,
  359. model: mdl,
  360. })
  361. const value = yield* handle.process({
  362. user: {
  363. id: parent.id,
  364. sessionID: chat.id,
  365. role: "user",
  366. time: parent.time,
  367. agent: parent.agent,
  368. model: { providerID: ref.providerID, modelID: ref.modelID },
  369. } satisfies SessionV1.User,
  370. sessionID: chat.id,
  371. model: mdl,
  372. agent: agent(),
  373. system: [],
  374. messages: [{ role: "user", content: "compact" }],
  375. tools: {},
  376. })
  377. const parts = yield* MessageV2.parts(msg.id)
  378. expect(value).toBe("compact")
  379. expect(parts.some((part) => part.type === "text" && part.text === "after")).toBe(true)
  380. expect(parts.some((part) => part.type === "step-finish")).toBe(true)
  381. }),
  382. { config: (url) => providerCfg(url) },
  383. ),
  384. )
  385. it.live("session.processor effect tests capture reasoning from http mock", () =>
  386. provideTmpdirServer(
  387. ({ dir, llm }) =>
  388. Effect.gen(function* () {
  389. const database = yield* Database.Service
  390. const { processors, session, provider } = yield* boot()
  391. yield* llm.push(reply().reason("think").text("done").stop())
  392. const chat = yield* session.create({})
  393. const parent = yield* user(chat.id, "reason")
  394. const msg = yield* assistant(chat.id, parent.id, path.resolve(dir))
  395. const mdl = yield* provider.getModel(ref.providerID, ref.modelID)
  396. const handle = yield* processors.create({
  397. assistantMessage: msg,
  398. sessionID: chat.id,
  399. model: mdl,
  400. })
  401. const value = yield* handle.process({
  402. user: {
  403. id: parent.id,
  404. sessionID: chat.id,
  405. role: "user",
  406. time: parent.time,
  407. agent: parent.agent,
  408. model: { providerID: ref.providerID, modelID: ref.modelID },
  409. } satisfies SessionV1.User,
  410. sessionID: chat.id,
  411. model: mdl,
  412. agent: agent(),
  413. system: [],
  414. messages: [{ role: "user", content: "reason" }],
  415. tools: {},
  416. })
  417. const parts = yield* MessageV2.parts(msg.id)
  418. const reasoning = parts.find((part): part is SessionV1.ReasoningPart => part.type === "reasoning")
  419. const text = parts.find((part): part is SessionV1.TextPart => part.type === "text")
  420. expect(value).toBe("continue")
  421. expect(yield* llm.calls).toBe(1)
  422. expect(reasoning?.text).toBe("think")
  423. expect(text?.text).toBe("done")
  424. }),
  425. { config: (url) => providerCfg(url) },
  426. ),
  427. )
  428. it.live("session.processor effect tests reset reasoning state across retries", () =>
  429. provideTmpdirServer(
  430. ({ dir, llm }) =>
  431. Effect.gen(function* () {
  432. const { processors, session, provider } = yield* boot()
  433. yield* llm.push(reply().reason("one").reset(), reply().reason("two").stop())
  434. const chat = yield* session.create({})
  435. const parent = yield* user(chat.id, "reason")
  436. const msg = yield* assistant(chat.id, parent.id, path.resolve(dir))
  437. const mdl = yield* provider.getModel(ref.providerID, ref.modelID)
  438. const handle = yield* processors.create({
  439. assistantMessage: msg,
  440. sessionID: chat.id,
  441. model: mdl,
  442. })
  443. const value = yield* handle.process({
  444. user: {
  445. id: parent.id,
  446. sessionID: chat.id,
  447. role: "user",
  448. time: parent.time,
  449. agent: parent.agent,
  450. model: { providerID: ref.providerID, modelID: ref.modelID },
  451. } satisfies SessionV1.User,
  452. sessionID: chat.id,
  453. model: mdl,
  454. agent: agent(),
  455. system: [],
  456. messages: [{ role: "user", content: "reason" }],
  457. tools: {},
  458. })
  459. const parts = yield* MessageV2.parts(msg.id)
  460. const reasoning = parts.filter((part): part is SessionV1.ReasoningPart => part.type === "reasoning")
  461. expect(value).toBe("continue")
  462. expect(yield* llm.calls).toBe(2)
  463. expect(reasoning.some((part) => part.text === "two")).toBe(true)
  464. expect(reasoning.some((part) => part.text === "onetwo")).toBe(false)
  465. }),
  466. { config: (url) => providerCfg(url) },
  467. ),
  468. )
  469. it.live("session.processor effect tests do not retry unknown json errors", () =>
  470. provideTmpdirServer(
  471. ({ dir, llm }) =>
  472. Effect.gen(function* () {
  473. const { processors, session, provider } = yield* boot()
  474. yield* llm.error(400, { error: { message: "no_kv_space" } })
  475. const chat = yield* session.create({})
  476. const parent = yield* user(chat.id, "json")
  477. const msg = yield* assistant(chat.id, parent.id, path.resolve(dir))
  478. const mdl = yield* provider.getModel(ref.providerID, ref.modelID)
  479. const handle = yield* processors.create({
  480. assistantMessage: msg,
  481. sessionID: chat.id,
  482. model: mdl,
  483. })
  484. const value = yield* handle.process({
  485. user: {
  486. id: parent.id,
  487. sessionID: chat.id,
  488. role: "user",
  489. time: parent.time,
  490. agent: parent.agent,
  491. model: { providerID: ref.providerID, modelID: ref.modelID },
  492. } satisfies SessionV1.User,
  493. sessionID: chat.id,
  494. model: mdl,
  495. agent: agent(),
  496. system: [],
  497. messages: [{ role: "user", content: "json" }],
  498. tools: {},
  499. })
  500. expect(value).toBe("stop")
  501. expect(yield* llm.calls).toBe(1)
  502. expect(handle.message.error?.name).toBe("APIError")
  503. }),
  504. { config: (url) => providerCfg(url) },
  505. ),
  506. )
  507. it.live("session.processor effect tests retry recognized structured json errors", () =>
  508. provideTmpdirServer(
  509. ({ dir, llm }) =>
  510. Effect.gen(function* () {
  511. const { processors, session, provider } = yield* boot()
  512. yield* llm.error(429, { type: "error", error: { type: "too_many_requests" } })
  513. yield* llm.text("after")
  514. const chat = yield* session.create({})
  515. const parent = yield* user(chat.id, "retry json")
  516. const msg = yield* assistant(chat.id, parent.id, path.resolve(dir))
  517. const mdl = yield* provider.getModel(ref.providerID, ref.modelID)
  518. const handle = yield* processors.create({
  519. assistantMessage: msg,
  520. sessionID: chat.id,
  521. model: mdl,
  522. })
  523. const value = yield* handle.process({
  524. user: {
  525. id: parent.id,
  526. sessionID: chat.id,
  527. role: "user",
  528. time: parent.time,
  529. agent: parent.agent,
  530. model: { providerID: ref.providerID, modelID: ref.modelID },
  531. } satisfies SessionV1.User,
  532. sessionID: chat.id,
  533. model: mdl,
  534. agent: agent(),
  535. system: [],
  536. messages: [{ role: "user", content: "retry json" }],
  537. tools: {},
  538. })
  539. const parts = yield* MessageV2.parts(msg.id)
  540. expect(value).toBe("continue")
  541. expect(yield* llm.calls).toBe(2)
  542. expect(parts.some((part) => part.type === "text" && part.text === "after")).toBe(true)
  543. expect(handle.message.error).toBeUndefined()
  544. }),
  545. { config: (url) => providerCfg(url) },
  546. ),
  547. )
  548. it.live("session.processor effect tests publish retry status updates", () =>
  549. provideTmpdirServer(
  550. ({ dir, llm }) =>
  551. Effect.gen(function* () {
  552. const { processors, session, provider } = yield* boot()
  553. const events = yield* EventV2Bridge.Service
  554. yield* llm.error(503, { error: "boom" })
  555. yield* llm.text("")
  556. const chat = yield* session.create({})
  557. const parent = yield* user(chat.id, "retry")
  558. const msg = yield* assistant(chat.id, parent.id, path.resolve(dir))
  559. const mdl = yield* provider.getModel(ref.providerID, ref.modelID)
  560. const states: number[] = []
  561. const off = yield* events.listen((evt) => {
  562. if (evt.type !== SessionStatus.Event.Status.type) return Effect.void
  563. const data = evt.data as typeof SessionStatus.Event.Status.data.Type
  564. if (data.sessionID === chat.id && data.status.type === "retry") states.push(data.status.attempt)
  565. return Effect.void
  566. })
  567. const handle = yield* processors.create({
  568. assistantMessage: msg,
  569. sessionID: chat.id,
  570. model: mdl,
  571. })
  572. const value = yield* handle.process({
  573. user: {
  574. id: parent.id,
  575. sessionID: chat.id,
  576. role: "user",
  577. time: parent.time,
  578. agent: parent.agent,
  579. model: { providerID: ref.providerID, modelID: ref.modelID },
  580. } satisfies SessionV1.User,
  581. sessionID: chat.id,
  582. model: mdl,
  583. agent: agent(),
  584. system: [],
  585. messages: [{ role: "user", content: "retry" }],
  586. tools: {},
  587. })
  588. yield* off
  589. expect(value).toBe("continue")
  590. expect(yield* llm.calls).toBe(2)
  591. expect(states).toStrictEqual([1])
  592. }),
  593. { config: (url) => providerCfg(url) },
  594. ),
  595. )
  596. it.live("session.processor effect tests compact on structured context overflow", () =>
  597. provideTmpdirServer(
  598. ({ dir, llm }) =>
  599. Effect.gen(function* () {
  600. const { processors, session, provider } = yield* boot()
  601. yield* llm.error(400, { type: "error", error: { code: "context_length_exceeded" } })
  602. const chat = yield* session.create({})
  603. const parent = yield* user(chat.id, "compact json")
  604. const msg = yield* assistant(chat.id, parent.id, path.resolve(dir))
  605. const mdl = yield* provider.getModel(ref.providerID, ref.modelID)
  606. const handle = yield* processors.create({
  607. assistantMessage: msg,
  608. sessionID: chat.id,
  609. model: mdl,
  610. })
  611. const value = yield* handle.process({
  612. user: {
  613. id: parent.id,
  614. sessionID: chat.id,
  615. role: "user",
  616. time: parent.time,
  617. agent: parent.agent,
  618. model: { providerID: ref.providerID, modelID: ref.modelID },
  619. } satisfies SessionV1.User,
  620. sessionID: chat.id,
  621. model: mdl,
  622. agent: agent(),
  623. system: [],
  624. messages: [{ role: "user", content: "compact json" }],
  625. tools: {},
  626. })
  627. expect(value).toBe("compact")
  628. expect(yield* llm.calls).toBe(1)
  629. expect(handle.message.error).toBeUndefined()
  630. }),
  631. { config: (url) => providerCfg(url) },
  632. ),
  633. )
  634. it.live("session.processor effect tests complete AI SDK tool calls when native flag is off", () =>
  635. provideTmpdirServer(
  636. ({ dir, llm }) =>
  637. Effect.gen(function* () {
  638. const { processors, session, provider } = yield* boot()
  639. yield* llm.tool("lookup", { query: "weather" })
  640. const chat = yield* session.create({})
  641. const parent = yield* user(chat.id, "tool")
  642. const msg = yield* assistant(chat.id, parent.id, path.resolve(dir))
  643. const mdl = yield* provider.getModel(ref.providerID, ref.modelID)
  644. const handle = yield* processors.create({
  645. assistantMessage: msg,
  646. sessionID: chat.id,
  647. model: mdl,
  648. })
  649. const value = yield* handle.process({
  650. user: {
  651. id: parent.id,
  652. sessionID: chat.id,
  653. role: "user",
  654. time: parent.time,
  655. agent: parent.agent,
  656. model: { providerID: ref.providerID, modelID: ref.modelID },
  657. } satisfies SessionV1.User,
  658. sessionID: chat.id,
  659. model: mdl,
  660. agent: agent(),
  661. system: [],
  662. messages: [{ role: "user", content: "tool" }],
  663. tools: {
  664. lookup: tool({
  665. description: "Look up information",
  666. inputSchema: z.object({ query: z.string() }),
  667. execute: async (input) => ({
  668. title: "Weather lookup",
  669. output: `result:${input.query}`,
  670. metadata: { source: "test" },
  671. }),
  672. }),
  673. },
  674. })
  675. const parts = yield* MessageV2.parts(msg.id)
  676. const call = parts.find((part): part is SessionV1.ToolPart => part.type === "tool")
  677. expect(value).toBe("continue")
  678. expect(yield* llm.calls).toBe(1)
  679. expect(call?.callID).toBe("call_1")
  680. expect(call?.tool).toBe("lookup")
  681. expect(call?.state.status).toBe("completed")
  682. if (call?.state.status !== "completed") return
  683. expect(call.state.input).toEqual({ query: "weather" })
  684. expect(call.state.output).toBe("result:weather")
  685. expect(call.state.title).toBe("Weather lookup")
  686. expect(call.state.metadata).toEqual({ source: "test" })
  687. expect(call.state.time.start).toBeDefined()
  688. expect(call.state.time.end).toBeDefined()
  689. }),
  690. { config: (url) => providerCfg(url) },
  691. ),
  692. )
  693. it.live("session.processor effect tests mark pending tools as aborted on cleanup", () =>
  694. provideTmpdirServer(
  695. ({ dir, llm }) =>
  696. Effect.gen(function* () {
  697. const database = yield* Database.Service
  698. const { processors, session, provider } = yield* boot()
  699. yield* llm.toolHang("bash", { cmd: "pwd" })
  700. const chat = yield* session.create({})
  701. const parent = yield* user(chat.id, "tool abort")
  702. const msg = yield* assistant(chat.id, parent.id, path.resolve(dir))
  703. const mdl = yield* provider.getModel(ref.providerID, ref.modelID)
  704. const handle = yield* processors.create({
  705. assistantMessage: msg,
  706. sessionID: chat.id,
  707. model: mdl,
  708. })
  709. const run = yield* handle
  710. .process({
  711. user: {
  712. id: parent.id,
  713. sessionID: chat.id,
  714. role: "user",
  715. time: parent.time,
  716. agent: parent.agent,
  717. model: { providerID: ref.providerID, modelID: ref.modelID },
  718. } satisfies SessionV1.User,
  719. sessionID: chat.id,
  720. model: mdl,
  721. agent: agent(),
  722. system: [],
  723. messages: [{ role: "user", content: "tool abort" }],
  724. tools: {},
  725. })
  726. .pipe(Effect.forkChild)
  727. yield* llm.wait(1)
  728. yield* waitFor(
  729. MessageV2.parts(msg.id).pipe(
  730. Effect.map((parts) => parts.find((part): part is SessionV1.ToolPart => part.type === "tool")),
  731. Effect.provideService(Database.Service, database),
  732. ),
  733. "timed out waiting for tool part",
  734. )
  735. yield* Fiber.interrupt(run)
  736. const exit = yield* Fiber.await(run)
  737. const parts = yield* MessageV2.parts(msg.id)
  738. const call = parts.find((part): part is SessionV1.ToolPart => part.type === "tool")
  739. expect(Exit.isFailure(exit)).toBe(true)
  740. if (Exit.isFailure(exit)) {
  741. expect(Cause.hasInterruptsOnly(exit.cause)).toBe(true)
  742. }
  743. expect(yield* llm.calls).toBe(1)
  744. expect(call?.state.status).toBe("error")
  745. if (call?.state.status === "error") {
  746. expect(call.state.error).toBe("Tool execution aborted")
  747. expect(call.state.metadata?.interrupted).toBe(true)
  748. expect(call.state.time.end).toBeDefined()
  749. }
  750. }),
  751. { config: (url) => providerCfg(url) },
  752. ),
  753. )
  754. it.live("session.processor effect tests record aborted errors and idle state", () =>
  755. provideTmpdirServer(
  756. ({ dir, llm }) =>
  757. Effect.gen(function* () {
  758. const seen = defer<void>()
  759. const { processors, session, provider } = yield* boot()
  760. const events = yield* EventV2Bridge.Service
  761. const sts = yield* SessionStatus.Service
  762. yield* llm.hang
  763. const chat = yield* session.create({})
  764. const parent = yield* user(chat.id, "abort")
  765. const msg = yield* assistant(chat.id, parent.id, path.resolve(dir))
  766. const mdl = yield* provider.getModel(ref.providerID, ref.modelID)
  767. const errs: string[] = []
  768. const off = yield* events.listen((evt) => {
  769. if (evt.type !== Session.Event.Error.type) return Effect.void
  770. const data = evt.data as typeof Session.Event.Error.data.Type
  771. if (data.sessionID !== chat.id || !data.error) return Effect.void
  772. errs.push(data.error.name)
  773. seen.resolve()
  774. return Effect.void
  775. })
  776. const handle = yield* processors.create({
  777. assistantMessage: msg,
  778. sessionID: chat.id,
  779. model: mdl,
  780. })
  781. const run = yield* handle
  782. .process({
  783. user: {
  784. id: parent.id,
  785. sessionID: chat.id,
  786. role: "user",
  787. time: parent.time,
  788. agent: parent.agent,
  789. model: { providerID: ref.providerID, modelID: ref.modelID },
  790. } satisfies SessionV1.User,
  791. sessionID: chat.id,
  792. model: mdl,
  793. agent: agent(),
  794. system: [],
  795. messages: [{ role: "user", content: "abort" }],
  796. tools: {},
  797. })
  798. .pipe(Effect.forkChild)
  799. yield* llm.wait(1)
  800. yield* Fiber.interrupt(run)
  801. const exit = yield* Fiber.await(run)
  802. yield* Effect.promise(() => seen.promise)
  803. const stored = yield* MessageV2.get({ sessionID: chat.id, messageID: msg.id })
  804. const state = yield* sts.get(chat.id)
  805. yield* off
  806. expect(Exit.isFailure(exit)).toBe(true)
  807. if (Exit.isFailure(exit)) {
  808. expect(Cause.hasInterruptsOnly(exit.cause)).toBe(true)
  809. }
  810. expect(handle.message.error?.name).toBe("MessageAbortedError")
  811. expect(stored.info.role).toBe("assistant")
  812. if (stored.info.role === "assistant") {
  813. expect(stored.info.error?.name).toBe("MessageAbortedError")
  814. }
  815. expect(state).toMatchObject({ type: "idle" })
  816. expect(errs).toContain("MessageAbortedError")
  817. }),
  818. { config: (url) => providerCfg(url) },
  819. ),
  820. )
  821. it.live("session.processor effect tests mark interruptions aborted without manual abort", () =>
  822. provideTmpdirServer(
  823. ({ dir, llm }) =>
  824. Effect.gen(function* () {
  825. const { processors, session, provider } = yield* boot()
  826. const sts = yield* SessionStatus.Service
  827. yield* llm.hang
  828. const chat = yield* session.create({})
  829. const parent = yield* user(chat.id, "interrupt")
  830. const msg = yield* assistant(chat.id, parent.id, path.resolve(dir))
  831. const mdl = yield* provider.getModel(ref.providerID, ref.modelID)
  832. const handle = yield* processors.create({
  833. assistantMessage: msg,
  834. sessionID: chat.id,
  835. model: mdl,
  836. })
  837. const run = yield* handle
  838. .process({
  839. user: {
  840. id: parent.id,
  841. sessionID: chat.id,
  842. role: "user",
  843. time: parent.time,
  844. agent: parent.agent,
  845. model: { providerID: ref.providerID, modelID: ref.modelID },
  846. } satisfies SessionV1.User,
  847. sessionID: chat.id,
  848. model: mdl,
  849. agent: agent(),
  850. system: [],
  851. messages: [{ role: "user", content: "interrupt" }],
  852. tools: {},
  853. })
  854. .pipe(Effect.forkChild)
  855. yield* llm.wait(1)
  856. yield* Fiber.interrupt(run)
  857. const exit = yield* Fiber.await(run)
  858. const stored = yield* MessageV2.get({ sessionID: chat.id, messageID: msg.id })
  859. const state = yield* sts.get(chat.id)
  860. expect(Exit.isFailure(exit)).toBe(true)
  861. expect(handle.message.error?.name).toBe("MessageAbortedError")
  862. expect(stored.info.role).toBe("assistant")
  863. if (stored.info.role === "assistant") {
  864. expect(stored.info.error?.name).toBe("MessageAbortedError")
  865. }
  866. expect(state).toMatchObject({ type: "idle" })
  867. }),
  868. { config: (url) => providerCfg(url) },
  869. ),
  870. )
  871. itProviderError.live("session.processor effect tests fail provider-executed error results", () =>
  872. provideTmpdirInstance(
  873. (dir) =>
  874. Effect.gen(function* () {
  875. const { processors, session, provider } = yield* boot()
  876. const events = yield* EventV2Bridge.Service
  877. const chat = yield* session.create({})
  878. const parent = yield* user(chat.id, "provider tool error")
  879. const msg = yield* assistant(chat.id, parent.id, path.resolve(dir))
  880. const mdl = yield* provider.getModel(ref.providerID, ref.modelID)
  881. const seen: string[] = []
  882. const off = yield* events.listen((event) => {
  883. seen.push(event.type)
  884. return Effect.void
  885. })
  886. const handle = yield* processors.create({ assistantMessage: msg, sessionID: chat.id, model: mdl })
  887. yield* handle.process({
  888. user: {
  889. id: parent.id,
  890. sessionID: chat.id,
  891. role: "user",
  892. time: parent.time,
  893. agent: parent.agent,
  894. model: { providerID: ref.providerID, modelID: ref.modelID },
  895. } satisfies SessionV1.User,
  896. sessionID: chat.id,
  897. model: mdl,
  898. agent: agent(),
  899. system: [],
  900. messages: [{ role: "user", content: "provider tool error" }],
  901. tools: {},
  902. })
  903. yield* off
  904. const parts = yield* MessageV2.parts(msg.id)
  905. const call = parts.find((part): part is SessionV1.ToolPart => part.type === "tool")
  906. expect(call?.state.status).toBe("error")
  907. if (call?.state.status === "error") expect(call.state.error).toBe("provider boom")
  908. expect(seen).toContain(MessageV2.Event.PartUpdated.type)
  909. expect(seen).toContain(MessageV2.Event.Updated.type)
  910. expect(seen.filter((type) => type.startsWith("session.next."))).toEqual([])
  911. }),
  912. { config: cfg },
  913. ),
  914. )
  915. itFragmentFailure.live("session.processor effect tests retain partial legacy parts without v2 events", () =>
  916. provideTmpdirInstance(
  917. (dir) =>
  918. Effect.gen(function* () {
  919. const { processors, session, provider } = yield* boot()
  920. const events = yield* EventV2Bridge.Service
  921. const chat = yield* session.create({})
  922. const parent = yield* user(chat.id, "provider failure")
  923. const msg = yield* assistant(chat.id, parent.id, path.resolve(dir))
  924. const mdl = yield* provider.getModel(ref.providerID, ref.modelID)
  925. const seen: string[] = []
  926. const off = yield* events.listen((event) => {
  927. seen.push(event.type)
  928. return Effect.void
  929. })
  930. const handle = yield* processors.create({ assistantMessage: msg, sessionID: chat.id, model: mdl })
  931. expect(
  932. yield* handle.process({
  933. user: {
  934. id: parent.id,
  935. sessionID: chat.id,
  936. role: "user",
  937. time: parent.time,
  938. agent: parent.agent,
  939. model: { providerID: ref.providerID, modelID: ref.modelID },
  940. } satisfies SessionV1.User,
  941. sessionID: chat.id,
  942. model: mdl,
  943. agent: agent(),
  944. system: [],
  945. messages: [{ role: "user", content: "provider failure" }],
  946. tools: {},
  947. }),
  948. ).toBe("stop")
  949. yield* off
  950. const parts = yield* MessageV2.parts(msg.id)
  951. expect(parts).toEqual(
  952. expect.arrayContaining([
  953. expect.objectContaining({ type: "text", text: "partial" }),
  954. expect.objectContaining({ type: "reasoning", text: "thinking" }),
  955. ]),
  956. )
  957. expect(seen).toContain(MessageV2.Event.PartUpdated.type)
  958. expect(seen).toContain(Session.Event.Error.type)
  959. expect(seen.filter((type) => type.startsWith("session.next."))).toEqual([])
  960. }),
  961. { config: cfg },
  962. ),
  963. )