processor-effect.test.ts 27 KB

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