processor-effect.test.ts 36 KB

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