processor-effect.test.ts 26 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838
  1. import { NodeFileSystem } from "@effect/platform-node"
  2. import { expect } from "bun:test"
  3. import { APICallError } from "ai"
  4. import { Effect, Layer, ServiceMap } from "effect"
  5. import * as Stream from "effect/Stream"
  6. import path from "path"
  7. import type { Agent } from "../../src/agent/agent"
  8. import { Agent as AgentSvc } from "../../src/agent/agent"
  9. import { Bus } from "../../src/bus"
  10. import { Config } from "../../src/config/config"
  11. import { Permission } from "../../src/permission"
  12. import { Plugin } from "../../src/plugin"
  13. import { Instance } from "../../src/project/instance"
  14. import type { Provider } from "../../src/provider/provider"
  15. import { ModelID, ProviderID } from "../../src/provider/schema"
  16. import { Session } from "../../src/session"
  17. import { LLM } from "../../src/session/llm"
  18. import { MessageV2 } from "../../src/session/message-v2"
  19. import { SessionProcessor } from "../../src/session/processor"
  20. import { MessageID, PartID, SessionID } from "../../src/session/schema"
  21. import { SessionStatus } from "../../src/session/status"
  22. import { Snapshot } from "../../src/snapshot"
  23. import { Log } from "../../src/util/log"
  24. import * as CrossSpawnSpawner from "../../src/effect/cross-spawn-spawner"
  25. import { provideTmpdirInstance } from "../fixture/fixture"
  26. import { testEffect } from "../lib/effect"
  27. Log.init({ print: false })
  28. const ref = {
  29. providerID: ProviderID.make("test"),
  30. modelID: ModelID.make("test-model"),
  31. }
  32. type Script = Stream.Stream<LLM.Event, unknown> | ((input: LLM.StreamInput) => Stream.Stream<LLM.Event, unknown>)
  33. class TestLLM extends ServiceMap.Service<
  34. TestLLM,
  35. {
  36. readonly push: (stream: Script) => Effect.Effect<void>
  37. readonly reply: (...items: LLM.Event[]) => Effect.Effect<void>
  38. readonly calls: Effect.Effect<number>
  39. readonly inputs: Effect.Effect<LLM.StreamInput[]>
  40. }
  41. >()("@test/SessionProcessorLLM") {}
  42. function stream(...items: LLM.Event[]) {
  43. return Stream.make(...items)
  44. }
  45. function usage(input = 1, output = 1, total = input + output) {
  46. return {
  47. inputTokens: input,
  48. outputTokens: output,
  49. totalTokens: total,
  50. inputTokenDetails: {
  51. noCacheTokens: undefined,
  52. cacheReadTokens: undefined,
  53. cacheWriteTokens: undefined,
  54. },
  55. outputTokenDetails: {
  56. textTokens: undefined,
  57. reasoningTokens: undefined,
  58. },
  59. }
  60. }
  61. function start(): LLM.Event {
  62. return { type: "start" }
  63. }
  64. function textStart(id = "t"): LLM.Event {
  65. return { type: "text-start", id }
  66. }
  67. function textDelta(id: string, text: string): LLM.Event {
  68. return { type: "text-delta", id, text }
  69. }
  70. function textEnd(id = "t"): LLM.Event {
  71. return { type: "text-end", id }
  72. }
  73. function reasoningStart(id: string): LLM.Event {
  74. return { type: "reasoning-start", id }
  75. }
  76. function reasoningDelta(id: string, text: string): LLM.Event {
  77. return { type: "reasoning-delta", id, text }
  78. }
  79. function reasoningEnd(id: string): LLM.Event {
  80. return { type: "reasoning-end", id }
  81. }
  82. function finishStep(): LLM.Event {
  83. return {
  84. type: "finish-step",
  85. finishReason: "stop",
  86. rawFinishReason: "stop",
  87. response: { id: "res", modelId: "test-model", timestamp: new Date() },
  88. providerMetadata: undefined,
  89. usage: usage(),
  90. }
  91. }
  92. function finish(): LLM.Event {
  93. return { type: "finish", finishReason: "stop", rawFinishReason: "stop", totalUsage: usage() }
  94. }
  95. function toolInputStart(id: string, toolName: string): LLM.Event {
  96. return { type: "tool-input-start", id, toolName }
  97. }
  98. function toolCall(toolCallId: string, toolName: string, input: unknown): LLM.Event {
  99. return { type: "tool-call", toolCallId, toolName, input }
  100. }
  101. function fail<E>(err: E, ...items: LLM.Event[]) {
  102. return stream(...items).pipe(Stream.concat(Stream.fail(err)))
  103. }
  104. function wait(abort: AbortSignal) {
  105. return Effect.promise(
  106. () =>
  107. new Promise<void>((done) => {
  108. abort.addEventListener("abort", () => done(), { once: true })
  109. }),
  110. )
  111. }
  112. function hang(input: LLM.StreamInput, ...items: LLM.Event[]) {
  113. return stream(...items).pipe(
  114. Stream.concat(
  115. Stream.unwrap(wait(input.abort).pipe(Effect.as(Stream.fail(new DOMException("Aborted", "AbortError"))))),
  116. ),
  117. )
  118. }
  119. function model(context: number): Provider.Model {
  120. return {
  121. id: "test-model",
  122. providerID: "test",
  123. name: "Test",
  124. limit: { context, output: 10 },
  125. cost: { input: 0, output: 0, cache: { read: 0, write: 0 } },
  126. capabilities: {
  127. toolcall: true,
  128. attachment: false,
  129. reasoning: false,
  130. temperature: true,
  131. input: { text: true, image: false, audio: false, video: false },
  132. output: { text: true, image: false, audio: false, video: false },
  133. },
  134. api: { npm: "@ai-sdk/anthropic" },
  135. options: {},
  136. } as Provider.Model
  137. }
  138. function agent(): Agent.Info {
  139. return {
  140. name: "build",
  141. mode: "primary",
  142. options: {},
  143. permission: [{ permission: "*", pattern: "*", action: "allow" }],
  144. }
  145. }
  146. function defer<T>() {
  147. let resolve!: (value: T | PromiseLike<T>) => void
  148. const promise = new Promise<T>((done) => {
  149. resolve = done
  150. })
  151. return { promise, resolve }
  152. }
  153. const user = Effect.fn("TestSession.user")(function* (sessionID: SessionID, text: string) {
  154. const session = yield* Session.Service
  155. const msg = yield* session.updateMessage({
  156. id: MessageID.ascending(),
  157. role: "user",
  158. sessionID,
  159. agent: "build",
  160. model: ref,
  161. time: { created: Date.now() },
  162. })
  163. yield* session.updatePart({
  164. id: PartID.ascending(),
  165. messageID: msg.id,
  166. sessionID,
  167. type: "text",
  168. text,
  169. })
  170. return msg
  171. })
  172. const assistant = Effect.fn("TestSession.assistant")(function* (
  173. sessionID: SessionID,
  174. parentID: MessageID,
  175. root: string,
  176. ) {
  177. const session = yield* Session.Service
  178. const msg: MessageV2.Assistant = {
  179. id: MessageID.ascending(),
  180. role: "assistant",
  181. sessionID,
  182. mode: "build",
  183. agent: "build",
  184. path: { cwd: root, root },
  185. cost: 0,
  186. tokens: {
  187. total: 0,
  188. input: 0,
  189. output: 0,
  190. reasoning: 0,
  191. cache: { read: 0, write: 0 },
  192. },
  193. modelID: ref.modelID,
  194. providerID: ref.providerID,
  195. parentID,
  196. time: { created: Date.now() },
  197. finish: "end_turn",
  198. }
  199. yield* session.updateMessage(msg)
  200. return msg
  201. })
  202. const llm = Layer.unwrap(
  203. Effect.gen(function* () {
  204. const queue: Script[] = []
  205. const inputs: LLM.StreamInput[] = []
  206. let calls = 0
  207. const push = Effect.fn("TestLLM.push")((item: Script) => {
  208. queue.push(item)
  209. return Effect.void
  210. })
  211. const reply = Effect.fn("TestLLM.reply")((...items: LLM.Event[]) => push(stream(...items)))
  212. return Layer.mergeAll(
  213. Layer.succeed(
  214. LLM.Service,
  215. LLM.Service.of({
  216. stream: (input) => {
  217. calls += 1
  218. inputs.push(input)
  219. const item = queue.shift() ?? Stream.empty
  220. return typeof item === "function" ? item(input) : item
  221. },
  222. }),
  223. ),
  224. Layer.succeed(
  225. TestLLM,
  226. TestLLM.of({
  227. push,
  228. reply,
  229. calls: Effect.sync(() => calls),
  230. inputs: Effect.sync(() => [...inputs]),
  231. }),
  232. ),
  233. )
  234. }),
  235. )
  236. const status = SessionStatus.layer.pipe(Layer.provideMerge(Bus.layer))
  237. const infra = Layer.mergeAll(NodeFileSystem.layer, CrossSpawnSpawner.defaultLayer)
  238. const deps = Layer.mergeAll(
  239. Session.defaultLayer,
  240. Snapshot.defaultLayer,
  241. AgentSvc.defaultLayer,
  242. Permission.layer,
  243. Plugin.defaultLayer,
  244. Config.defaultLayer,
  245. status,
  246. llm,
  247. ).pipe(Layer.provideMerge(infra))
  248. const env = SessionProcessor.layer.pipe(Layer.provideMerge(deps))
  249. const it = testEffect(env)
  250. it.effect("session.processor effect tests capture llm input cleanly", () => {
  251. return provideTmpdirInstance(
  252. (dir) =>
  253. Effect.gen(function* () {
  254. const test = yield* TestLLM
  255. const processors = yield* SessionProcessor.Service
  256. const session = yield* Session.Service
  257. yield* test.reply(start(), textStart(), textDelta("t", "hello"), textEnd(), finishStep(), finish())
  258. const chat = yield* session.create({})
  259. const parent = yield* user(chat.id, "hi")
  260. const msg = yield* assistant(chat.id, parent.id, path.resolve(dir))
  261. const abort = new AbortController()
  262. const mdl = model(100)
  263. const handle = yield* processors.create({
  264. assistantMessage: msg,
  265. sessionID: chat.id,
  266. model: mdl,
  267. abort: abort.signal,
  268. })
  269. const input = {
  270. user: {
  271. id: parent.id,
  272. sessionID: chat.id,
  273. role: "user",
  274. time: parent.time,
  275. agent: parent.agent,
  276. model: { providerID: ref.providerID, modelID: ref.modelID },
  277. } satisfies MessageV2.User,
  278. sessionID: chat.id,
  279. model: mdl,
  280. agent: agent(),
  281. system: [],
  282. abort: abort.signal,
  283. messages: [{ role: "user", content: "hi" }],
  284. tools: {},
  285. } satisfies LLM.StreamInput
  286. const value = yield* handle.process(input)
  287. const parts = yield* Effect.promise(() => MessageV2.parts(msg.id))
  288. const calls = yield* test.calls
  289. const inputs = yield* test.inputs
  290. expect(value).toBe("continue")
  291. expect(calls).toBe(1)
  292. expect(inputs).toHaveLength(1)
  293. expect(inputs[0].messages).toStrictEqual([{ role: "user", content: "hi" }])
  294. expect(parts.some((part) => part.type === "text" && part.text === "hello")).toBe(true)
  295. }),
  296. { git: true },
  297. )
  298. })
  299. it.effect("session.processor effect tests stop after token overflow requests compaction", () => {
  300. return provideTmpdirInstance(
  301. (dir) =>
  302. Effect.gen(function* () {
  303. const test = yield* TestLLM
  304. const processors = yield* SessionProcessor.Service
  305. const session = yield* Session.Service
  306. yield* test.reply(
  307. start(),
  308. {
  309. type: "finish-step",
  310. finishReason: "stop",
  311. rawFinishReason: "stop",
  312. response: { id: "res", modelId: "test-model", timestamp: new Date() },
  313. providerMetadata: undefined,
  314. usage: usage(100, 0, 100),
  315. },
  316. textStart(),
  317. textDelta("t", "after"),
  318. textEnd(),
  319. )
  320. const chat = yield* session.create({})
  321. const parent = yield* user(chat.id, "compact")
  322. const msg = yield* assistant(chat.id, parent.id, path.resolve(dir))
  323. const abort = new AbortController()
  324. const mdl = model(20)
  325. const handle = yield* processors.create({
  326. assistantMessage: msg,
  327. sessionID: chat.id,
  328. model: mdl,
  329. abort: abort.signal,
  330. })
  331. const value = yield* handle.process({
  332. user: {
  333. id: parent.id,
  334. sessionID: chat.id,
  335. role: "user",
  336. time: parent.time,
  337. agent: parent.agent,
  338. model: { providerID: ref.providerID, modelID: ref.modelID },
  339. } satisfies MessageV2.User,
  340. sessionID: chat.id,
  341. model: mdl,
  342. agent: agent(),
  343. system: [],
  344. abort: abort.signal,
  345. messages: [{ role: "user", content: "compact" }],
  346. tools: {},
  347. })
  348. const parts = yield* Effect.promise(() => MessageV2.parts(msg.id))
  349. expect(value).toBe("compact")
  350. expect(parts.some((part) => part.type === "text")).toBe(false)
  351. expect(parts.some((part) => part.type === "step-finish")).toBe(true)
  352. }),
  353. { git: true },
  354. )
  355. })
  356. it.effect("session.processor effect tests reset reasoning state across retries", () => {
  357. return provideTmpdirInstance(
  358. (dir) =>
  359. Effect.gen(function* () {
  360. const test = yield* TestLLM
  361. const processors = yield* SessionProcessor.Service
  362. const session = yield* Session.Service
  363. yield* test.push(
  364. fail(
  365. new APICallError({
  366. message: "boom",
  367. url: "https://example.com/v1/chat/completions",
  368. requestBodyValues: {},
  369. statusCode: 503,
  370. responseHeaders: { "retry-after-ms": "0" },
  371. responseBody: '{"error":"boom"}',
  372. isRetryable: true,
  373. }),
  374. start(),
  375. reasoningStart("r"),
  376. reasoningDelta("r", "one"),
  377. ),
  378. )
  379. yield* test.reply(
  380. start(),
  381. reasoningStart("r"),
  382. reasoningDelta("r", "two"),
  383. reasoningEnd("r"),
  384. finishStep(),
  385. finish(),
  386. )
  387. const chat = yield* session.create({})
  388. const parent = yield* user(chat.id, "reason")
  389. const msg = yield* assistant(chat.id, parent.id, path.resolve(dir))
  390. const abort = new AbortController()
  391. const mdl = model(100)
  392. const handle = yield* processors.create({
  393. assistantMessage: msg,
  394. sessionID: chat.id,
  395. model: mdl,
  396. abort: abort.signal,
  397. })
  398. const value = yield* handle.process({
  399. user: {
  400. id: parent.id,
  401. sessionID: chat.id,
  402. role: "user",
  403. time: parent.time,
  404. agent: parent.agent,
  405. model: { providerID: ref.providerID, modelID: ref.modelID },
  406. } satisfies MessageV2.User,
  407. sessionID: chat.id,
  408. model: mdl,
  409. agent: agent(),
  410. system: [],
  411. abort: abort.signal,
  412. messages: [{ role: "user", content: "reason" }],
  413. tools: {},
  414. })
  415. const parts = yield* Effect.promise(() => MessageV2.parts(msg.id))
  416. const reasoning = parts.filter((part): part is MessageV2.ReasoningPart => part.type === "reasoning")
  417. expect(value).toBe("continue")
  418. expect(yield* test.calls).toBe(2)
  419. expect(reasoning.some((part) => part.text === "two")).toBe(true)
  420. expect(reasoning.some((part) => part.text === "onetwo")).toBe(false)
  421. }),
  422. { git: true },
  423. )
  424. })
  425. it.effect("session.processor effect tests do not retry unknown json errors", () => {
  426. return provideTmpdirInstance(
  427. (dir) =>
  428. Effect.gen(function* () {
  429. const test = yield* TestLLM
  430. const processors = yield* SessionProcessor.Service
  431. const session = yield* Session.Service
  432. yield* test.push(fail({ error: { message: "no_kv_space" } }, start()))
  433. const chat = yield* session.create({})
  434. const parent = yield* user(chat.id, "json")
  435. const msg = yield* assistant(chat.id, parent.id, path.resolve(dir))
  436. const abort = new AbortController()
  437. const mdl = model(100)
  438. const handle = yield* processors.create({
  439. assistantMessage: msg,
  440. sessionID: chat.id,
  441. model: mdl,
  442. abort: abort.signal,
  443. })
  444. const value = yield* handle.process({
  445. user: {
  446. id: parent.id,
  447. sessionID: chat.id,
  448. role: "user",
  449. time: parent.time,
  450. agent: parent.agent,
  451. model: { providerID: ref.providerID, modelID: ref.modelID },
  452. } satisfies MessageV2.User,
  453. sessionID: chat.id,
  454. model: mdl,
  455. agent: agent(),
  456. system: [],
  457. abort: abort.signal,
  458. messages: [{ role: "user", content: "json" }],
  459. tools: {},
  460. })
  461. expect(value).toBe("stop")
  462. expect(yield* test.calls).toBe(1)
  463. expect(yield* test.inputs).toHaveLength(1)
  464. expect(handle.message.error?.name).toBe("UnknownError")
  465. }),
  466. { git: true },
  467. )
  468. })
  469. it.effect("session.processor effect tests retry recognized structured json errors", () => {
  470. return provideTmpdirInstance(
  471. (dir) =>
  472. Effect.gen(function* () {
  473. const test = yield* TestLLM
  474. const processors = yield* SessionProcessor.Service
  475. const session = yield* Session.Service
  476. yield* test.push(fail({ type: "error", error: { type: "too_many_requests" } }, start()))
  477. yield* test.reply(start(), textStart(), textDelta("t", "after"), textEnd(), finishStep(), finish())
  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 abort = new AbortController()
  482. const mdl = model(100)
  483. const handle = yield* processors.create({
  484. assistantMessage: msg,
  485. sessionID: chat.id,
  486. model: mdl,
  487. abort: abort.signal,
  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 MessageV2.User,
  498. sessionID: chat.id,
  499. model: mdl,
  500. agent: agent(),
  501. system: [],
  502. abort: abort.signal,
  503. messages: [{ role: "user", content: "retry json" }],
  504. tools: {},
  505. })
  506. const parts = yield* Effect.promise(() => MessageV2.parts(msg.id))
  507. expect(value).toBe("continue")
  508. expect(yield* test.calls).toBe(2)
  509. expect(parts.some((part) => part.type === "text" && part.text === "after")).toBe(true)
  510. expect(handle.message.error).toBeUndefined()
  511. }),
  512. { git: true },
  513. )
  514. })
  515. it.effect("session.processor effect tests publish retry status updates", () => {
  516. return provideTmpdirInstance(
  517. (dir) =>
  518. Effect.gen(function* () {
  519. const test = yield* TestLLM
  520. const processors = yield* SessionProcessor.Service
  521. const session = yield* Session.Service
  522. const bus = yield* Bus.Service
  523. yield* test.push(
  524. fail(
  525. new APICallError({
  526. message: "boom",
  527. url: "https://example.com/v1/chat/completions",
  528. requestBodyValues: {},
  529. statusCode: 503,
  530. responseHeaders: { "retry-after-ms": "0" },
  531. responseBody: '{"error":"boom"}',
  532. isRetryable: true,
  533. }),
  534. start(),
  535. ),
  536. )
  537. yield* test.reply(start(), finishStep(), finish())
  538. const chat = yield* session.create({})
  539. const parent = yield* user(chat.id, "retry")
  540. const msg = yield* assistant(chat.id, parent.id, path.resolve(dir))
  541. const abort = new AbortController()
  542. const mdl = model(100)
  543. const states: number[] = []
  544. const off = yield* bus.subscribeCallback(SessionStatus.Event.Status, (evt) => {
  545. if (evt.properties.sessionID !== chat.id) return
  546. if (evt.properties.status.type === "retry") states.push(evt.properties.status.attempt)
  547. })
  548. const handle = yield* processors.create({
  549. assistantMessage: msg,
  550. sessionID: chat.id,
  551. model: mdl,
  552. abort: abort.signal,
  553. })
  554. const value = yield* handle.process({
  555. user: {
  556. id: parent.id,
  557. sessionID: chat.id,
  558. role: "user",
  559. time: parent.time,
  560. agent: parent.agent,
  561. model: { providerID: ref.providerID, modelID: ref.modelID },
  562. } satisfies MessageV2.User,
  563. sessionID: chat.id,
  564. model: mdl,
  565. agent: agent(),
  566. system: [],
  567. abort: abort.signal,
  568. messages: [{ role: "user", content: "retry" }],
  569. tools: {},
  570. })
  571. off()
  572. expect(value).toBe("continue")
  573. expect(yield* test.calls).toBe(2)
  574. expect(states).toStrictEqual([1])
  575. }),
  576. { git: true },
  577. )
  578. })
  579. it.effect("session.processor effect tests compact on structured context overflow", () => {
  580. return provideTmpdirInstance(
  581. (dir) =>
  582. Effect.gen(function* () {
  583. const test = yield* TestLLM
  584. const processors = yield* SessionProcessor.Service
  585. const session = yield* Session.Service
  586. yield* test.push(fail({ type: "error", error: { code: "context_length_exceeded" } }, start()))
  587. const chat = yield* session.create({})
  588. const parent = yield* user(chat.id, "compact json")
  589. const msg = yield* assistant(chat.id, parent.id, path.resolve(dir))
  590. const abort = new AbortController()
  591. const mdl = model(100)
  592. const handle = yield* processors.create({
  593. assistantMessage: msg,
  594. sessionID: chat.id,
  595. model: mdl,
  596. abort: abort.signal,
  597. })
  598. const value = yield* handle.process({
  599. user: {
  600. id: parent.id,
  601. sessionID: chat.id,
  602. role: "user",
  603. time: parent.time,
  604. agent: parent.agent,
  605. model: { providerID: ref.providerID, modelID: ref.modelID },
  606. } satisfies MessageV2.User,
  607. sessionID: chat.id,
  608. model: mdl,
  609. agent: agent(),
  610. system: [],
  611. abort: abort.signal,
  612. messages: [{ role: "user", content: "compact json" }],
  613. tools: {},
  614. })
  615. expect(value).toBe("compact")
  616. expect(yield* test.calls).toBe(1)
  617. expect(handle.message.error).toBeUndefined()
  618. }),
  619. { git: true },
  620. )
  621. })
  622. it.effect("session.processor effect tests mark pending tools as aborted on cleanup", () => {
  623. return provideTmpdirInstance(
  624. (dir) =>
  625. Effect.gen(function* () {
  626. const ready = defer<void>()
  627. const seen = defer<void>()
  628. const test = yield* TestLLM
  629. const processors = yield* SessionProcessor.Service
  630. const session = yield* Session.Service
  631. yield* test.push((input) =>
  632. hang(input, start(), toolInputStart("tool-1", "bash"), toolCall("tool-1", "bash", { cmd: "pwd" })).pipe(
  633. Stream.tap((event) => (event.type === "tool-call" ? Effect.sync(() => ready.resolve()) : Effect.void)),
  634. ),
  635. )
  636. const chat = yield* session.create({})
  637. const parent = yield* user(chat.id, "tool abort")
  638. const msg = yield* assistant(chat.id, parent.id, path.resolve(dir))
  639. const abort = new AbortController()
  640. const mdl = model(100)
  641. const handle = yield* processors.create({
  642. assistantMessage: msg,
  643. sessionID: chat.id,
  644. model: mdl,
  645. abort: abort.signal,
  646. })
  647. const run = Effect.runPromise(
  648. handle.process({
  649. user: {
  650. id: parent.id,
  651. sessionID: chat.id,
  652. role: "user",
  653. time: parent.time,
  654. agent: parent.agent,
  655. model: { providerID: ref.providerID, modelID: ref.modelID },
  656. } satisfies MessageV2.User,
  657. sessionID: chat.id,
  658. model: mdl,
  659. agent: agent(),
  660. system: [],
  661. abort: abort.signal,
  662. messages: [{ role: "user", content: "tool abort" }],
  663. tools: {},
  664. }),
  665. )
  666. yield* Effect.promise(() => ready.promise)
  667. abort.abort()
  668. const value = yield* Effect.promise(() => run)
  669. const parts = yield* Effect.promise(() => MessageV2.parts(msg.id))
  670. const tool = parts.find((part): part is MessageV2.ToolPart => part.type === "tool")
  671. expect(value).toBe("stop")
  672. expect(yield* test.calls).toBe(1)
  673. expect(tool?.state.status).toBe("error")
  674. if (tool?.state.status === "error") {
  675. expect(tool.state.error).toBe("Tool execution aborted")
  676. expect(tool.state.time.end).toBeDefined()
  677. }
  678. }),
  679. { git: true },
  680. )
  681. })
  682. it.effect("session.processor effect tests record aborted errors and idle state", () => {
  683. return provideTmpdirInstance(
  684. (dir) =>
  685. Effect.gen(function* () {
  686. const ready = defer<void>()
  687. const seen = defer<void>()
  688. const test = yield* TestLLM
  689. const processors = yield* SessionProcessor.Service
  690. const session = yield* Session.Service
  691. const bus = yield* Bus.Service
  692. const status = yield* SessionStatus.Service
  693. yield* test.push((input) =>
  694. hang(input, start()).pipe(
  695. Stream.tap((event) => (event.type === "start" ? Effect.sync(() => ready.resolve()) : Effect.void)),
  696. ),
  697. )
  698. const chat = yield* session.create({})
  699. const parent = yield* user(chat.id, "abort")
  700. const msg = yield* assistant(chat.id, parent.id, path.resolve(dir))
  701. const abort = new AbortController()
  702. const mdl = model(100)
  703. const errs: string[] = []
  704. const off = yield* bus.subscribeCallback(Session.Event.Error, (evt) => {
  705. if (evt.properties.sessionID !== chat.id) return
  706. if (!evt.properties.error) return
  707. errs.push(evt.properties.error.name)
  708. seen.resolve()
  709. })
  710. const handle = yield* processors.create({
  711. assistantMessage: msg,
  712. sessionID: chat.id,
  713. model: mdl,
  714. abort: abort.signal,
  715. })
  716. const run = Effect.runPromise(
  717. handle.process({
  718. user: {
  719. id: parent.id,
  720. sessionID: chat.id,
  721. role: "user",
  722. time: parent.time,
  723. agent: parent.agent,
  724. model: { providerID: ref.providerID, modelID: ref.modelID },
  725. } satisfies MessageV2.User,
  726. sessionID: chat.id,
  727. model: mdl,
  728. agent: agent(),
  729. system: [],
  730. abort: abort.signal,
  731. messages: [{ role: "user", content: "abort" }],
  732. tools: {},
  733. }),
  734. )
  735. yield* Effect.promise(() => ready.promise)
  736. abort.abort()
  737. const value = yield* Effect.promise(() => run)
  738. yield* Effect.promise(() => seen.promise)
  739. const stored = yield* Effect.promise(() => MessageV2.get({ sessionID: chat.id, messageID: msg.id }))
  740. const state = yield* status.get(chat.id)
  741. off()
  742. expect(value).toBe("stop")
  743. expect(handle.message.error?.name).toBe("MessageAbortedError")
  744. expect(stored.info.role).toBe("assistant")
  745. if (stored.info.role === "assistant") {
  746. expect(stored.info.error?.name).toBe("MessageAbortedError")
  747. }
  748. expect(state).toMatchObject({ type: "idle" })
  749. expect(errs).toContain("MessageAbortedError")
  750. }),
  751. { git: true },
  752. )
  753. })