processor-effect.test.ts 36 KB

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