processor-effect.test.ts 27 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845
  1. import { NodeFileSystem } from "@effect/platform-node"
  2. import { expect } from "bun:test"
  3. import { Cause, Effect, Exit, Fiber, Layer } from "effect"
  4. import path from "path"
  5. import type { Agent } from "../../src/agent/agent"
  6. import { Agent as AgentSvc } from "../../src/agent/agent"
  7. import { Bus } from "../../src/bus"
  8. import { Config } from "@/config/config"
  9. import { Image } from "@/image/image"
  10. import { Permission } from "../../src/permission"
  11. import { Plugin } from "../../src/plugin"
  12. import { Provider } from "@/provider/provider"
  13. import { ModelID, ProviderID } from "../../src/provider/schema"
  14. import { Session } from "@/session/session"
  15. import { LLM } from "../../src/session/llm"
  16. import { MessageV2 } from "../../src/session/message-v2"
  17. import { SessionProcessor } from "../../src/session/processor"
  18. import { MessageID, PartID, SessionID } from "../../src/session/schema"
  19. import { SessionStatus } from "../../src/session/status"
  20. import { SessionSummary } from "../../src/session/summary"
  21. import { Snapshot } from "../../src/snapshot"
  22. import * as Log from "@opencode-ai/core/util/log"
  23. import { CrossSpawnSpawner } from "@opencode-ai/core/cross-spawn-spawner"
  24. import { provideTmpdirServer } from "../fixture/fixture"
  25. import { testEffect } from "../lib/effect"
  26. import { raw, reply, TestLLMServer } from "../lib/llm-server"
  27. import { SyncEvent } from "@/sync"
  28. void Log.init({ print: false })
  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: ProviderID.make("test"),
  39. modelID: ModelID.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 user = Effect.fn("TestSession.user")(function* (sessionID: SessionID, text: string) {
  100. const session = yield* Session.Service
  101. const msg = yield* session.updateMessage({
  102. id: MessageID.ascending(),
  103. role: "user",
  104. sessionID,
  105. agent: "build",
  106. model: ref,
  107. time: { created: Date.now() },
  108. })
  109. yield* session.updatePart({
  110. id: PartID.ascending(),
  111. messageID: msg.id,
  112. sessionID,
  113. type: "text",
  114. text,
  115. })
  116. return msg
  117. })
  118. const assistant = Effect.fn("TestSession.assistant")(function* (
  119. sessionID: SessionID,
  120. parentID: MessageID,
  121. root: string,
  122. ) {
  123. const session = yield* Session.Service
  124. const msg: MessageV2.Assistant = {
  125. id: MessageID.ascending(),
  126. role: "assistant",
  127. sessionID,
  128. mode: "build",
  129. agent: "build",
  130. path: { cwd: root, root },
  131. cost: 0,
  132. tokens: {
  133. total: 0,
  134. input: 0,
  135. output: 0,
  136. reasoning: 0,
  137. cache: { read: 0, write: 0 },
  138. },
  139. modelID: ref.modelID,
  140. providerID: ref.providerID,
  141. parentID,
  142. time: { created: Date.now() },
  143. finish: "end_turn",
  144. }
  145. yield* session.updateMessage(msg)
  146. return msg
  147. })
  148. const status = SessionStatus.layer.pipe(Layer.provideMerge(Bus.layer))
  149. const infra = Layer.mergeAll(NodeFileSystem.layer, CrossSpawnSpawner.defaultLayer)
  150. const deps = Layer.mergeAll(
  151. Session.defaultLayer,
  152. Snapshot.defaultLayer,
  153. AgentSvc.defaultLayer,
  154. Permission.defaultLayer,
  155. Plugin.defaultLayer,
  156. Config.defaultLayer,
  157. LLM.defaultLayer,
  158. Provider.defaultLayer,
  159. status,
  160. SyncEvent.defaultLayer,
  161. ).pipe(Layer.provideMerge(infra))
  162. const env = Layer.mergeAll(
  163. TestLLMServer.layer,
  164. SessionProcessor.layer.pipe(Layer.provide(summary), Layer.provide(Image.defaultLayer), Layer.provideMerge(deps)),
  165. )
  166. const it = testEffect(env)
  167. const boot = Effect.fn("test.boot")(function* () {
  168. const processors = yield* SessionProcessor.Service
  169. const session = yield* Session.Service
  170. const provider = yield* Provider.Service
  171. return { processors, session, provider }
  172. })
  173. // ---------------------------------------------------------------------------
  174. // Tests
  175. // ---------------------------------------------------------------------------
  176. it.live("session.processor effect tests capture llm input cleanly", () =>
  177. provideTmpdirServer(
  178. ({ dir, llm }) =>
  179. Effect.gen(function* () {
  180. const { processors, session, provider } = yield* boot()
  181. yield* llm.text("hello")
  182. const chat = yield* session.create({})
  183. const parent = yield* user(chat.id, "hi")
  184. const msg = yield* assistant(chat.id, parent.id, path.resolve(dir))
  185. const mdl = yield* provider.getModel(ref.providerID, ref.modelID)
  186. const handle = yield* processors.create({
  187. assistantMessage: msg,
  188. sessionID: chat.id,
  189. model: mdl,
  190. })
  191. const input = {
  192. user: {
  193. id: parent.id,
  194. sessionID: chat.id,
  195. role: "user",
  196. time: parent.time,
  197. agent: parent.agent,
  198. model: { providerID: ref.providerID, modelID: ref.modelID },
  199. } satisfies MessageV2.User,
  200. sessionID: chat.id,
  201. model: mdl,
  202. agent: agent(),
  203. system: [],
  204. messages: [{ role: "user", content: "hi" }],
  205. tools: {},
  206. } satisfies LLM.StreamInput
  207. const value = yield* handle.process(input)
  208. const parts = MessageV2.parts(msg.id)
  209. const calls = yield* llm.calls
  210. expect(value).toBe("continue")
  211. expect(calls).toBe(1)
  212. expect(parts.some((part) => part.type === "text" && part.text === "hello")).toBe(true)
  213. }),
  214. { git: true, config: (url) => providerCfg(url) },
  215. ),
  216. )
  217. it.live("session.processor effect tests preserve text start time", () =>
  218. provideTmpdirServer(
  219. ({ dir, llm }) =>
  220. Effect.gen(function* () {
  221. const gate = defer<void>()
  222. const { processors, session, provider } = yield* boot()
  223. yield* llm.push(
  224. raw({
  225. head: [
  226. {
  227. id: "chatcmpl-test",
  228. object: "chat.completion.chunk",
  229. choices: [{ delta: { role: "assistant" } }],
  230. },
  231. {
  232. id: "chatcmpl-test",
  233. object: "chat.completion.chunk",
  234. choices: [{ delta: { content: "hello" } }],
  235. },
  236. ],
  237. wait: gate.promise,
  238. tail: [
  239. {
  240. id: "chatcmpl-test",
  241. object: "chat.completion.chunk",
  242. choices: [{ delta: {}, finish_reason: "stop" }],
  243. },
  244. ],
  245. }),
  246. )
  247. const chat = yield* session.create({})
  248. const parent = yield* user(chat.id, "hi")
  249. const msg = yield* assistant(chat.id, parent.id, path.resolve(dir))
  250. const mdl = yield* provider.getModel(ref.providerID, ref.modelID)
  251. const handle = yield* processors.create({
  252. assistantMessage: msg,
  253. sessionID: chat.id,
  254. model: mdl,
  255. })
  256. const run = yield* handle
  257. .process({
  258. user: {
  259. id: parent.id,
  260. sessionID: chat.id,
  261. role: "user",
  262. time: parent.time,
  263. agent: parent.agent,
  264. model: { providerID: ref.providerID, modelID: ref.modelID },
  265. } satisfies MessageV2.User,
  266. sessionID: chat.id,
  267. model: mdl,
  268. agent: agent(),
  269. system: [],
  270. messages: [{ role: "user", content: "hi" }],
  271. tools: {},
  272. })
  273. .pipe(Effect.forkChild)
  274. yield* Effect.promise(async () => {
  275. const stop = Date.now() + 500
  276. while (Date.now() < stop) {
  277. const text = MessageV2.parts(msg.id).find((part): part is MessageV2.TextPart => part.type === "text")
  278. if (text?.time?.start) return
  279. await Bun.sleep(10)
  280. }
  281. throw new Error("timed out waiting for text part")
  282. })
  283. yield* Effect.sleep("20 millis")
  284. gate.resolve()
  285. const exit = yield* Fiber.await(run)
  286. const text = MessageV2.parts(msg.id).find((part): part is MessageV2.TextPart => part.type === "text")
  287. expect(Exit.isSuccess(exit)).toBe(true)
  288. expect(text?.text).toBe("hello")
  289. expect(text?.time?.start).toBeDefined()
  290. expect(text?.time?.end).toBeDefined()
  291. if (!text?.time?.start || !text.time.end) return
  292. expect(text.time.start).toBeLessThan(text.time.end)
  293. }),
  294. { git: true, config: (url) => providerCfg(url) },
  295. ),
  296. )
  297. it.live("session.processor effect tests stop after token overflow requests compaction", () =>
  298. provideTmpdirServer(
  299. ({ dir, llm }) =>
  300. Effect.gen(function* () {
  301. const { processors, session, provider } = yield* boot()
  302. yield* llm.text("after", { usage: { input: 100, output: 0 } })
  303. const chat = yield* session.create({})
  304. const parent = yield* user(chat.id, "compact")
  305. const msg = yield* assistant(chat.id, parent.id, path.resolve(dir))
  306. const base = yield* provider.getModel(ref.providerID, ref.modelID)
  307. const mdl = { ...base, limit: { context: 20, output: 10 } }
  308. const handle = yield* processors.create({
  309. assistantMessage: msg,
  310. sessionID: chat.id,
  311. model: mdl,
  312. })
  313. const value = yield* handle.process({
  314. user: {
  315. id: parent.id,
  316. sessionID: chat.id,
  317. role: "user",
  318. time: parent.time,
  319. agent: parent.agent,
  320. model: { providerID: ref.providerID, modelID: ref.modelID },
  321. } satisfies MessageV2.User,
  322. sessionID: chat.id,
  323. model: mdl,
  324. agent: agent(),
  325. system: [],
  326. messages: [{ role: "user", content: "compact" }],
  327. tools: {},
  328. })
  329. const parts = MessageV2.parts(msg.id)
  330. expect(value).toBe("compact")
  331. expect(parts.some((part) => part.type === "text" && part.text === "after")).toBe(true)
  332. expect(parts.some((part) => part.type === "step-finish")).toBe(true)
  333. }),
  334. { git: true, config: (url) => providerCfg(url) },
  335. ),
  336. )
  337. it.live("session.processor effect tests capture reasoning from http mock", () =>
  338. provideTmpdirServer(
  339. ({ dir, llm }) =>
  340. Effect.gen(function* () {
  341. const { processors, session, provider } = yield* boot()
  342. yield* llm.push(reply().reason("think").text("done").stop())
  343. const chat = yield* session.create({})
  344. const parent = yield* user(chat.id, "reason")
  345. const msg = yield* assistant(chat.id, parent.id, path.resolve(dir))
  346. const mdl = yield* provider.getModel(ref.providerID, ref.modelID)
  347. const handle = yield* processors.create({
  348. assistantMessage: msg,
  349. sessionID: chat.id,
  350. model: mdl,
  351. })
  352. const value = yield* handle.process({
  353. user: {
  354. id: parent.id,
  355. sessionID: chat.id,
  356. role: "user",
  357. time: parent.time,
  358. agent: parent.agent,
  359. model: { providerID: ref.providerID, modelID: ref.modelID },
  360. } satisfies MessageV2.User,
  361. sessionID: chat.id,
  362. model: mdl,
  363. agent: agent(),
  364. system: [],
  365. messages: [{ role: "user", content: "reason" }],
  366. tools: {},
  367. })
  368. const parts = MessageV2.parts(msg.id)
  369. const reasoning = parts.find((part): part is MessageV2.ReasoningPart => part.type === "reasoning")
  370. const text = parts.find((part): part is MessageV2.TextPart => part.type === "text")
  371. expect(value).toBe("continue")
  372. expect(yield* llm.calls).toBe(1)
  373. expect(reasoning?.text).toBe("think")
  374. expect(text?.text).toBe("done")
  375. }),
  376. { git: true, config: (url) => providerCfg(url) },
  377. ),
  378. )
  379. it.live("session.processor effect tests reset reasoning state across retries", () =>
  380. provideTmpdirServer(
  381. ({ dir, llm }) =>
  382. Effect.gen(function* () {
  383. const { processors, session, provider } = yield* boot()
  384. yield* llm.push(reply().reason("one").reset(), reply().reason("two").stop())
  385. const chat = yield* session.create({})
  386. const parent = yield* user(chat.id, "reason")
  387. const msg = yield* assistant(chat.id, parent.id, path.resolve(dir))
  388. const mdl = yield* provider.getModel(ref.providerID, ref.modelID)
  389. const handle = yield* processors.create({
  390. assistantMessage: msg,
  391. sessionID: chat.id,
  392. model: mdl,
  393. })
  394. const value = yield* handle.process({
  395. user: {
  396. id: parent.id,
  397. sessionID: chat.id,
  398. role: "user",
  399. time: parent.time,
  400. agent: parent.agent,
  401. model: { providerID: ref.providerID, modelID: ref.modelID },
  402. } satisfies MessageV2.User,
  403. sessionID: chat.id,
  404. model: mdl,
  405. agent: agent(),
  406. system: [],
  407. messages: [{ role: "user", content: "reason" }],
  408. tools: {},
  409. })
  410. const parts = MessageV2.parts(msg.id)
  411. const reasoning = parts.filter((part): part is MessageV2.ReasoningPart => part.type === "reasoning")
  412. expect(value).toBe("continue")
  413. expect(yield* llm.calls).toBe(2)
  414. expect(reasoning.some((part) => part.text === "two")).toBe(true)
  415. expect(reasoning.some((part) => part.text === "onetwo")).toBe(false)
  416. }),
  417. { git: true, config: (url) => providerCfg(url) },
  418. ),
  419. )
  420. it.live("session.processor effect tests do not retry unknown json errors", () =>
  421. provideTmpdirServer(
  422. ({ dir, llm }) =>
  423. Effect.gen(function* () {
  424. const { processors, session, provider } = yield* boot()
  425. yield* llm.error(400, { error: { message: "no_kv_space" } })
  426. const chat = yield* session.create({})
  427. const parent = yield* user(chat.id, "json")
  428. const msg = yield* assistant(chat.id, parent.id, path.resolve(dir))
  429. const mdl = yield* provider.getModel(ref.providerID, ref.modelID)
  430. const handle = yield* processors.create({
  431. assistantMessage: msg,
  432. sessionID: chat.id,
  433. model: mdl,
  434. })
  435. const value = yield* handle.process({
  436. user: {
  437. id: parent.id,
  438. sessionID: chat.id,
  439. role: "user",
  440. time: parent.time,
  441. agent: parent.agent,
  442. model: { providerID: ref.providerID, modelID: ref.modelID },
  443. } satisfies MessageV2.User,
  444. sessionID: chat.id,
  445. model: mdl,
  446. agent: agent(),
  447. system: [],
  448. messages: [{ role: "user", content: "json" }],
  449. tools: {},
  450. })
  451. expect(value).toBe("stop")
  452. expect(yield* llm.calls).toBe(1)
  453. expect(handle.message.error?.name).toBe("APIError")
  454. }),
  455. { git: true, config: (url) => providerCfg(url) },
  456. ),
  457. )
  458. it.live("session.processor effect tests retry recognized structured json errors", () =>
  459. provideTmpdirServer(
  460. ({ dir, llm }) =>
  461. Effect.gen(function* () {
  462. const { processors, session, provider } = yield* boot()
  463. yield* llm.error(429, { type: "error", error: { type: "too_many_requests" } })
  464. yield* llm.text("after")
  465. const chat = yield* session.create({})
  466. const parent = yield* user(chat.id, "retry json")
  467. const msg = yield* assistant(chat.id, parent.id, path.resolve(dir))
  468. const mdl = yield* provider.getModel(ref.providerID, ref.modelID)
  469. const handle = yield* processors.create({
  470. assistantMessage: msg,
  471. sessionID: chat.id,
  472. model: mdl,
  473. })
  474. const value = yield* handle.process({
  475. user: {
  476. id: parent.id,
  477. sessionID: chat.id,
  478. role: "user",
  479. time: parent.time,
  480. agent: parent.agent,
  481. model: { providerID: ref.providerID, modelID: ref.modelID },
  482. } satisfies MessageV2.User,
  483. sessionID: chat.id,
  484. model: mdl,
  485. agent: agent(),
  486. system: [],
  487. messages: [{ role: "user", content: "retry json" }],
  488. tools: {},
  489. })
  490. const parts = MessageV2.parts(msg.id)
  491. expect(value).toBe("continue")
  492. expect(yield* llm.calls).toBe(2)
  493. expect(parts.some((part) => part.type === "text" && part.text === "after")).toBe(true)
  494. expect(handle.message.error).toBeUndefined()
  495. }),
  496. { git: true, config: (url) => providerCfg(url) },
  497. ),
  498. )
  499. it.live("session.processor effect tests publish retry status updates", () =>
  500. provideTmpdirServer(
  501. ({ dir, llm }) =>
  502. Effect.gen(function* () {
  503. const { processors, session, provider } = yield* boot()
  504. const bus = yield* Bus.Service
  505. yield* llm.error(503, { error: "boom" })
  506. yield* llm.text("")
  507. const chat = yield* session.create({})
  508. const parent = yield* user(chat.id, "retry")
  509. const msg = yield* assistant(chat.id, parent.id, path.resolve(dir))
  510. const mdl = yield* provider.getModel(ref.providerID, ref.modelID)
  511. const states: number[] = []
  512. const off = yield* bus.subscribeCallback(SessionStatus.Event.Status, (evt) => {
  513. if (evt.properties.sessionID !== chat.id) return
  514. if (evt.properties.status.type === "retry") states.push(evt.properties.status.attempt)
  515. })
  516. const handle = yield* processors.create({
  517. assistantMessage: msg,
  518. sessionID: chat.id,
  519. model: mdl,
  520. })
  521. const value = yield* handle.process({
  522. user: {
  523. id: parent.id,
  524. sessionID: chat.id,
  525. role: "user",
  526. time: parent.time,
  527. agent: parent.agent,
  528. model: { providerID: ref.providerID, modelID: ref.modelID },
  529. } satisfies MessageV2.User,
  530. sessionID: chat.id,
  531. model: mdl,
  532. agent: agent(),
  533. system: [],
  534. messages: [{ role: "user", content: "retry" }],
  535. tools: {},
  536. })
  537. off()
  538. expect(value).toBe("continue")
  539. expect(yield* llm.calls).toBe(2)
  540. expect(states).toStrictEqual([1])
  541. }),
  542. { git: true, config: (url) => providerCfg(url) },
  543. ),
  544. )
  545. it.live("session.processor effect tests compact on structured context overflow", () =>
  546. provideTmpdirServer(
  547. ({ dir, llm }) =>
  548. Effect.gen(function* () {
  549. const { processors, session, provider } = yield* boot()
  550. yield* llm.error(400, { type: "error", error: { code: "context_length_exceeded" } })
  551. const chat = yield* session.create({})
  552. const parent = yield* user(chat.id, "compact json")
  553. const msg = yield* assistant(chat.id, parent.id, path.resolve(dir))
  554. const mdl = yield* provider.getModel(ref.providerID, ref.modelID)
  555. const handle = yield* processors.create({
  556. assistantMessage: msg,
  557. sessionID: chat.id,
  558. model: mdl,
  559. })
  560. const value = yield* handle.process({
  561. user: {
  562. id: parent.id,
  563. sessionID: chat.id,
  564. role: "user",
  565. time: parent.time,
  566. agent: parent.agent,
  567. model: { providerID: ref.providerID, modelID: ref.modelID },
  568. } satisfies MessageV2.User,
  569. sessionID: chat.id,
  570. model: mdl,
  571. agent: agent(),
  572. system: [],
  573. messages: [{ role: "user", content: "compact json" }],
  574. tools: {},
  575. })
  576. expect(value).toBe("compact")
  577. expect(yield* llm.calls).toBe(1)
  578. expect(handle.message.error).toBeUndefined()
  579. }),
  580. { git: true, config: (url) => providerCfg(url) },
  581. ),
  582. )
  583. it.live("session.processor effect tests mark pending tools as aborted on cleanup", () =>
  584. provideTmpdirServer(
  585. ({ dir, llm }) =>
  586. Effect.gen(function* () {
  587. const { processors, session, provider } = yield* boot()
  588. yield* llm.toolHang("bash", { cmd: "pwd" })
  589. const chat = yield* session.create({})
  590. const parent = yield* user(chat.id, "tool abort")
  591. const msg = yield* assistant(chat.id, parent.id, path.resolve(dir))
  592. const mdl = yield* provider.getModel(ref.providerID, ref.modelID)
  593. const handle = yield* processors.create({
  594. assistantMessage: msg,
  595. sessionID: chat.id,
  596. model: mdl,
  597. })
  598. const run = yield* handle
  599. .process({
  600. user: {
  601. id: parent.id,
  602. sessionID: chat.id,
  603. role: "user",
  604. time: parent.time,
  605. agent: parent.agent,
  606. model: { providerID: ref.providerID, modelID: ref.modelID },
  607. } satisfies MessageV2.User,
  608. sessionID: chat.id,
  609. model: mdl,
  610. agent: agent(),
  611. system: [],
  612. messages: [{ role: "user", content: "tool abort" }],
  613. tools: {},
  614. })
  615. .pipe(Effect.forkChild)
  616. yield* llm.wait(1)
  617. yield* Effect.promise(async () => {
  618. const end = Date.now() + 500
  619. while (Date.now() < end) {
  620. const parts = await MessageV2.parts(msg.id)
  621. if (parts.some((part) => part.type === "tool")) return
  622. await Bun.sleep(10)
  623. }
  624. })
  625. yield* Fiber.interrupt(run)
  626. const exit = yield* Fiber.await(run)
  627. const parts = MessageV2.parts(msg.id)
  628. const call = parts.find((part): part is MessageV2.ToolPart => part.type === "tool")
  629. expect(Exit.isFailure(exit)).toBe(true)
  630. if (Exit.isFailure(exit)) {
  631. expect(Cause.hasInterruptsOnly(exit.cause)).toBe(true)
  632. }
  633. expect(yield* llm.calls).toBe(1)
  634. expect(call?.state.status).toBe("error")
  635. if (call?.state.status === "error") {
  636. expect(call.state.error).toBe("Tool execution aborted")
  637. expect(call.state.metadata?.interrupted).toBe(true)
  638. expect(call.state.time.end).toBeDefined()
  639. }
  640. }),
  641. { git: true, config: (url) => providerCfg(url) },
  642. ),
  643. )
  644. it.live("session.processor effect tests record aborted errors and idle state", () =>
  645. provideTmpdirServer(
  646. ({ dir, llm }) =>
  647. Effect.gen(function* () {
  648. const seen = defer<void>()
  649. const { processors, session, provider } = yield* boot()
  650. const bus = yield* Bus.Service
  651. const sts = yield* SessionStatus.Service
  652. yield* llm.hang
  653. const chat = yield* session.create({})
  654. const parent = yield* user(chat.id, "abort")
  655. const msg = yield* assistant(chat.id, parent.id, path.resolve(dir))
  656. const mdl = yield* provider.getModel(ref.providerID, ref.modelID)
  657. const errs: string[] = []
  658. const off = yield* bus.subscribeCallback(Session.Event.Error, (evt) => {
  659. if (evt.properties.sessionID !== chat.id) return
  660. if (!evt.properties.error) return
  661. errs.push(evt.properties.error.name)
  662. seen.resolve()
  663. })
  664. const handle = yield* processors.create({
  665. assistantMessage: msg,
  666. sessionID: chat.id,
  667. model: mdl,
  668. })
  669. const run = yield* handle
  670. .process({
  671. user: {
  672. id: parent.id,
  673. sessionID: chat.id,
  674. role: "user",
  675. time: parent.time,
  676. agent: parent.agent,
  677. model: { providerID: ref.providerID, modelID: ref.modelID },
  678. } satisfies MessageV2.User,
  679. sessionID: chat.id,
  680. model: mdl,
  681. agent: agent(),
  682. system: [],
  683. messages: [{ role: "user", content: "abort" }],
  684. tools: {},
  685. })
  686. .pipe(Effect.forkChild)
  687. yield* llm.wait(1)
  688. yield* Fiber.interrupt(run)
  689. const exit = yield* Fiber.await(run)
  690. yield* Effect.promise(() => seen.promise)
  691. const stored = yield* MessageV2.get({ sessionID: chat.id, messageID: msg.id })
  692. const state = yield* sts.get(chat.id)
  693. off()
  694. expect(Exit.isFailure(exit)).toBe(true)
  695. if (Exit.isFailure(exit)) {
  696. expect(Cause.hasInterruptsOnly(exit.cause)).toBe(true)
  697. }
  698. expect(handle.message.error?.name).toBe("MessageAbortedError")
  699. expect(stored.info.role).toBe("assistant")
  700. if (stored.info.role === "assistant") {
  701. expect(stored.info.error?.name).toBe("MessageAbortedError")
  702. }
  703. expect(state).toMatchObject({ type: "idle" })
  704. expect(errs).toContain("MessageAbortedError")
  705. }),
  706. { git: true, config: (url) => providerCfg(url) },
  707. ),
  708. )
  709. it.live("session.processor effect tests mark interruptions aborted without manual abort", () =>
  710. provideTmpdirServer(
  711. ({ dir, llm }) =>
  712. Effect.gen(function* () {
  713. const { processors, session, provider } = yield* boot()
  714. const sts = yield* SessionStatus.Service
  715. yield* llm.hang
  716. const chat = yield* session.create({})
  717. const parent = yield* user(chat.id, "interrupt")
  718. const msg = yield* assistant(chat.id, parent.id, path.resolve(dir))
  719. const mdl = yield* provider.getModel(ref.providerID, ref.modelID)
  720. const handle = yield* processors.create({
  721. assistantMessage: msg,
  722. sessionID: chat.id,
  723. model: mdl,
  724. })
  725. const run = yield* handle
  726. .process({
  727. user: {
  728. id: parent.id,
  729. sessionID: chat.id,
  730. role: "user",
  731. time: parent.time,
  732. agent: parent.agent,
  733. model: { providerID: ref.providerID, modelID: ref.modelID },
  734. } satisfies MessageV2.User,
  735. sessionID: chat.id,
  736. model: mdl,
  737. agent: agent(),
  738. system: [],
  739. messages: [{ role: "user", content: "interrupt" }],
  740. tools: {},
  741. })
  742. .pipe(Effect.forkChild)
  743. yield* llm.wait(1)
  744. yield* Fiber.interrupt(run)
  745. const exit = yield* Fiber.await(run)
  746. const stored = yield* MessageV2.get({ sessionID: chat.id, messageID: msg.id })
  747. const state = yield* sts.get(chat.id)
  748. expect(Exit.isFailure(exit)).toBe(true)
  749. expect(handle.message.error?.name).toBe("MessageAbortedError")
  750. expect(stored.info.role).toBe("assistant")
  751. if (stored.info.role === "assistant") {
  752. expect(stored.info.error?.name).toBe("MessageAbortedError")
  753. }
  754. expect(state).toMatchObject({ type: "idle" })
  755. }),
  756. { git: true, config: (url) => providerCfg(url) },
  757. ),
  758. )