processor-effect.test.ts 37 KB

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