| 12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073107410751076107710781079 |
- import { SessionV1 } from "@opencode-ai/core/v1/session"
- import { Database } from "@opencode-ai/core/database/database"
- import { LayerNode } from "@opencode-ai/core/effect/layer-node"
- import { EventV2Bridge } from "@/event-v2-bridge"
- import { expect } from "bun:test"
- import { tool } from "ai"
- import { Cause, Effect, Exit, Fiber, Layer, Stream } from "effect"
- import path from "path"
- import z from "zod"
- import type { Agent } from "../../src/agent/agent"
- import { Provider } from "@/provider/provider"
- import { Session } from "@/session/session"
- import { LLM } from "../../src/session/llm"
- import { MessageV2 } from "../../src/session/message-v2"
- import { SessionProcessor } from "../../src/session/processor"
- import { MessageID, PartID, SessionID } from "../../src/session/schema"
- import { SessionStatus } from "../../src/session/status"
- import { SessionSummary } from "../../src/session/summary"
- import { CrossSpawnSpawner } from "@opencode-ai/core/cross-spawn-spawner"
- import { provideTmpdirInstance, provideTmpdirServer } from "../fixture/fixture"
- import { testEffect } from "../lib/effect"
- import { raw, reply, TestLLMServer } from "../lib/llm-server"
- import { RuntimeFlags } from "@/effect/runtime-flags"
- import { ProviderV2 } from "@opencode-ai/core/provider"
- import { ModelV2 } from "@opencode-ai/core/model"
- import { SessionEvent } from "@opencode-ai/core/session/event"
- import { SessionProjector } from "@opencode-ai/core/session/projector"
- import { LLMEvent } from "@opencode-ai/llm"
- const summary = Layer.succeed(
- SessionSummary.Service,
- SessionSummary.Service.of({
- summarize: () => Effect.void,
- diff: () => Effect.succeed([]),
- computeDiff: () => Effect.succeed([]),
- }),
- )
- const ref = {
- providerID: ProviderV2.ID.make("test"),
- modelID: ModelV2.ID.make("test-model"),
- }
- const cfg = {
- provider: {
- test: {
- name: "Test",
- id: "test",
- env: [],
- npm: "@ai-sdk/openai-compatible",
- models: {
- "test-model": {
- id: "test-model",
- name: "Test Model",
- attachment: false,
- reasoning: false,
- temperature: false,
- tool_call: true,
- release_date: "2025-01-01",
- limit: { context: 100000, output: 10000 },
- cost: { input: 0, output: 0 },
- options: {},
- },
- },
- options: {
- apiKey: "test-key",
- baseURL: "http://localhost:1/v1",
- },
- },
- },
- }
- function providerCfg(url: string) {
- return {
- ...cfg,
- provider: {
- ...cfg.provider,
- test: {
- ...cfg.provider.test,
- options: {
- ...cfg.provider.test.options,
- baseURL: url,
- },
- },
- },
- }
- }
- function agent(): Agent.Info {
- return {
- name: "build",
- mode: "primary",
- options: {},
- permission: [{ permission: "*", pattern: "*", action: "allow" }],
- }
- }
- function defer<T>() {
- let resolve!: (value: T | PromiseLike<T>) => void
- const promise = new Promise<T>((done) => {
- resolve = done
- })
- return { promise, resolve }
- }
- const waitFor = <A>(check: Effect.Effect<A | undefined>, message: string) =>
- Effect.gen(function* () {
- const stop = Date.now() + 500
- while (Date.now() < stop) {
- const value = yield* check
- if (value !== undefined) return value
- yield* Effect.sleep("10 millis")
- }
- return yield* Effect.fail(new Error(message))
- })
- const user = Effect.fn("TestSession.user")(function* (sessionID: SessionID, text: string) {
- const session = yield* Session.Service
- const msg = yield* session.updateMessage({
- id: MessageID.ascending(),
- role: "user",
- sessionID,
- agent: "build",
- model: ref,
- time: { created: Date.now() },
- })
- yield* session.updatePart({
- id: PartID.ascending(),
- messageID: msg.id,
- sessionID,
- type: "text",
- text,
- })
- return msg
- })
- const assistant = Effect.fn("TestSession.assistant")(function* (
- sessionID: SessionID,
- parentID: MessageID,
- root: string,
- ) {
- const session = yield* Session.Service
- const msg: SessionV1.Assistant = {
- id: MessageID.ascending(),
- role: "assistant",
- sessionID,
- mode: "build",
- agent: "build",
- path: { cwd: root, root },
- cost: 0,
- tokens: {
- total: 0,
- input: 0,
- output: 0,
- reasoning: 0,
- cache: { read: 0, write: 0 },
- },
- modelID: ref.modelID,
- providerID: ref.providerID,
- parentID,
- time: { created: Date.now() },
- finish: "end_turn",
- }
- yield* session.updateMessage(msg)
- return msg
- })
- const root = LayerNode.group([
- SessionProcessor.node,
- Session.node,
- SessionProjector.node,
- Provider.node,
- Database.node,
- EventV2Bridge.node,
- SessionStatus.node,
- CrossSpawnSpawner.node,
- ])
- const replacements = [
- LayerNode.replace(SessionSummary.layer, summary),
- LayerNode.replace(RuntimeFlags.defaultLayer, RuntimeFlags.layer({ experimentalEventSystem: true })),
- ]
- const env = LayerNode.buildLayer(
- LayerNode.group([root, LayerNode.make({ service: TestLLMServer, layer: TestLLMServer.layer, deps: [] })]),
- { replacements },
- )
- const it = testEffect(env)
- const providerErrorLLM = Layer.succeed(
- LLM.Service,
- LLM.Service.of({
- stream: () =>
- Stream.make(
- LLMEvent.stepStart({ index: 0 }),
- LLMEvent.toolInputStart({ id: "call-1", name: "lookup" }),
- LLMEvent.toolInputEnd({ id: "call-1", name: "lookup" }),
- LLMEvent.toolCall({ id: "call-1", name: "lookup", input: {}, providerExecuted: true }),
- LLMEvent.toolResult({
- id: "call-1",
- name: "lookup",
- result: { type: "error", value: "provider boom" },
- providerExecuted: true,
- }),
- LLMEvent.stepFinish({ index: 0, reason: "stop" }),
- LLMEvent.finish({ reason: "stop" }),
- ),
- }),
- )
- const providerErrorEnv = LayerNode.buildLayer(root, {
- replacements: [...replacements, LayerNode.replace(LLM.layer, providerErrorLLM)],
- })
- const itProviderError = testEffect(providerErrorEnv)
- const fragmentFailureLLM = Layer.succeed(
- LLM.Service,
- LLM.Service.of({
- stream: () =>
- Stream.make(
- LLMEvent.stepStart({ index: 0 }),
- LLMEvent.reasoningStart({ id: "reasoning-1" }),
- LLMEvent.reasoningDelta({ id: "reasoning-1", text: "thinking" }),
- LLMEvent.textStart({ id: "text-1" }),
- LLMEvent.textDelta({ id: "text-1", text: "partial" }),
- LLMEvent.providerError({ message: "provider boom" }),
- ),
- }),
- )
- const fragmentFailureEnv = LayerNode.buildLayer(root, {
- replacements: [...replacements, LayerNode.replace(LLM.layer, fragmentFailureLLM)],
- })
- const itFragmentFailure = testEffect(fragmentFailureEnv)
- const boot = Effect.fn("test.boot")(function* () {
- const processors = yield* SessionProcessor.Service
- const session = yield* Session.Service
- const provider = yield* Provider.Service
- return { processors, session, provider }
- })
- // ---------------------------------------------------------------------------
- // Tests
- // ---------------------------------------------------------------------------
- it.live("session.processor effect tests capture llm input cleanly", () =>
- provideTmpdirServer(
- ({ dir, llm }) =>
- Effect.gen(function* () {
- const database = yield* Database.Service
- const { processors, session, provider } = yield* boot()
- yield* llm.text("hello")
- const chat = yield* session.create({})
- const parent = yield* user(chat.id, "hi")
- const msg = yield* assistant(chat.id, parent.id, path.resolve(dir))
- const mdl = yield* provider.getModel(ref.providerID, ref.modelID)
- const handle = yield* processors.create({
- assistantMessage: msg,
- sessionID: chat.id,
- model: mdl,
- })
- const input = {
- user: {
- id: parent.id,
- sessionID: chat.id,
- role: "user",
- time: parent.time,
- agent: parent.agent,
- model: { providerID: ref.providerID, modelID: ref.modelID },
- } satisfies SessionV1.User,
- sessionID: chat.id,
- model: mdl,
- agent: agent(),
- system: [],
- messages: [{ role: "user", content: "hi" }],
- tools: {},
- } satisfies LLM.StreamInput
- const value = yield* handle.process(input)
- const parts = yield* MessageV2.parts(msg.id)
- const calls = yield* llm.calls
- expect(value).toBe("continue")
- expect(calls).toBe(1)
- expect(parts.some((part) => part.type === "text" && part.text === "hello")).toBe(true)
- }),
- { config: (url) => providerCfg(url) },
- ),
- )
- it.live("session.processor effect tests preserve text start time", () =>
- provideTmpdirServer(
- ({ dir, llm }) =>
- Effect.gen(function* () {
- const database = yield* Database.Service
- const gate = defer<void>()
- const { processors, session, provider } = yield* boot()
- yield* llm.push(
- raw({
- head: [
- {
- id: "chatcmpl-test",
- object: "chat.completion.chunk",
- choices: [{ delta: { role: "assistant" } }],
- },
- {
- id: "chatcmpl-test",
- object: "chat.completion.chunk",
- choices: [{ delta: { content: "hello" } }],
- },
- ],
- wait: gate.promise,
- tail: [
- {
- id: "chatcmpl-test",
- object: "chat.completion.chunk",
- choices: [{ delta: {}, finish_reason: "stop" }],
- },
- ],
- }),
- )
- const chat = yield* session.create({})
- const parent = yield* user(chat.id, "hi")
- const msg = yield* assistant(chat.id, parent.id, path.resolve(dir))
- const mdl = yield* provider.getModel(ref.providerID, ref.modelID)
- const handle = yield* processors.create({
- assistantMessage: msg,
- sessionID: chat.id,
- model: mdl,
- })
- const run = yield* handle
- .process({
- user: {
- id: parent.id,
- sessionID: chat.id,
- role: "user",
- time: parent.time,
- agent: parent.agent,
- model: { providerID: ref.providerID, modelID: ref.modelID },
- } satisfies SessionV1.User,
- sessionID: chat.id,
- model: mdl,
- agent: agent(),
- system: [],
- messages: [{ role: "user", content: "hi" }],
- tools: {},
- })
- .pipe(Effect.forkChild)
- yield* waitFor(
- MessageV2.parts(msg.id).pipe(
- Effect.map((parts) => parts.find((part): part is SessionV1.TextPart => part.type === "text")),
- Effect.provideService(Database.Service, database),
- ),
- "timed out waiting for text part",
- )
- yield* Effect.sleep("20 millis")
- gate.resolve()
- const exit = yield* Fiber.await(run)
- const text = (yield* MessageV2.parts(msg.id)).find((part): part is SessionV1.TextPart => part.type === "text")
- expect(Exit.isSuccess(exit)).toBe(true)
- expect(text?.text).toBe("hello")
- expect(text?.time?.start).toBeDefined()
- expect(text?.time?.end).toBeDefined()
- if (!text?.time?.start || !text.time.end) return
- expect(text.time.start).toBeLessThan(text.time.end)
- }),
- { config: (url) => providerCfg(url) },
- ),
- )
- it.live("session.processor effect tests stop after token overflow requests compaction", () =>
- provideTmpdirServer(
- ({ dir, llm }) =>
- Effect.gen(function* () {
- const database = yield* Database.Service
- const { processors, session, provider } = yield* boot()
- yield* llm.text("after", { usage: { input: 100, output: 0 } })
- const chat = yield* session.create({})
- const parent = yield* user(chat.id, "compact")
- const msg = yield* assistant(chat.id, parent.id, path.resolve(dir))
- const base = yield* provider.getModel(ref.providerID, ref.modelID)
- const mdl = { ...base, limit: { context: 20, output: 10 } }
- const handle = yield* processors.create({
- assistantMessage: msg,
- sessionID: chat.id,
- model: mdl,
- })
- const value = yield* handle.process({
- user: {
- id: parent.id,
- sessionID: chat.id,
- role: "user",
- time: parent.time,
- agent: parent.agent,
- model: { providerID: ref.providerID, modelID: ref.modelID },
- } satisfies SessionV1.User,
- sessionID: chat.id,
- model: mdl,
- agent: agent(),
- system: [],
- messages: [{ role: "user", content: "compact" }],
- tools: {},
- })
- const parts = yield* MessageV2.parts(msg.id)
- expect(value).toBe("compact")
- expect(parts.some((part) => part.type === "text" && part.text === "after")).toBe(true)
- expect(parts.some((part) => part.type === "step-finish")).toBe(true)
- }),
- { config: (url) => providerCfg(url) },
- ),
- )
- it.live("session.processor effect tests capture reasoning from http mock", () =>
- provideTmpdirServer(
- ({ dir, llm }) =>
- Effect.gen(function* () {
- const database = yield* Database.Service
- const { processors, session, provider } = yield* boot()
- yield* llm.push(reply().reason("think").text("done").stop())
- const chat = yield* session.create({})
- const parent = yield* user(chat.id, "reason")
- const msg = yield* assistant(chat.id, parent.id, path.resolve(dir))
- const mdl = yield* provider.getModel(ref.providerID, ref.modelID)
- const handle = yield* processors.create({
- assistantMessage: msg,
- sessionID: chat.id,
- model: mdl,
- })
- const value = yield* handle.process({
- user: {
- id: parent.id,
- sessionID: chat.id,
- role: "user",
- time: parent.time,
- agent: parent.agent,
- model: { providerID: ref.providerID, modelID: ref.modelID },
- } satisfies SessionV1.User,
- sessionID: chat.id,
- model: mdl,
- agent: agent(),
- system: [],
- messages: [{ role: "user", content: "reason" }],
- tools: {},
- })
- const parts = yield* MessageV2.parts(msg.id)
- const reasoning = parts.find((part): part is SessionV1.ReasoningPart => part.type === "reasoning")
- const text = parts.find((part): part is SessionV1.TextPart => part.type === "text")
- expect(value).toBe("continue")
- expect(yield* llm.calls).toBe(1)
- expect(reasoning?.text).toBe("think")
- expect(text?.text).toBe("done")
- }),
- { config: (url) => providerCfg(url) },
- ),
- )
- it.live("session.processor effect tests reset reasoning state across retries", () =>
- provideTmpdirServer(
- ({ dir, llm }) =>
- Effect.gen(function* () {
- const { processors, session, provider } = yield* boot()
- yield* llm.push(reply().reason("one").reset(), reply().reason("two").stop())
- const chat = yield* session.create({})
- const parent = yield* user(chat.id, "reason")
- const msg = yield* assistant(chat.id, parent.id, path.resolve(dir))
- const mdl = yield* provider.getModel(ref.providerID, ref.modelID)
- const handle = yield* processors.create({
- assistantMessage: msg,
- sessionID: chat.id,
- model: mdl,
- })
- const value = yield* handle.process({
- user: {
- id: parent.id,
- sessionID: chat.id,
- role: "user",
- time: parent.time,
- agent: parent.agent,
- model: { providerID: ref.providerID, modelID: ref.modelID },
- } satisfies SessionV1.User,
- sessionID: chat.id,
- model: mdl,
- agent: agent(),
- system: [],
- messages: [{ role: "user", content: "reason" }],
- tools: {},
- })
- const parts = yield* MessageV2.parts(msg.id)
- const reasoning = parts.filter((part): part is SessionV1.ReasoningPart => part.type === "reasoning")
- expect(value).toBe("continue")
- expect(yield* llm.calls).toBe(2)
- expect(reasoning.some((part) => part.text === "two")).toBe(true)
- expect(reasoning.some((part) => part.text === "onetwo")).toBe(false)
- }),
- { config: (url) => providerCfg(url) },
- ),
- )
- it.live("session.processor effect tests do not retry unknown json errors", () =>
- provideTmpdirServer(
- ({ dir, llm }) =>
- Effect.gen(function* () {
- const { processors, session, provider } = yield* boot()
- yield* llm.error(400, { error: { message: "no_kv_space" } })
- const chat = yield* session.create({})
- const parent = yield* user(chat.id, "json")
- const msg = yield* assistant(chat.id, parent.id, path.resolve(dir))
- const mdl = yield* provider.getModel(ref.providerID, ref.modelID)
- const handle = yield* processors.create({
- assistantMessage: msg,
- sessionID: chat.id,
- model: mdl,
- })
- const value = yield* handle.process({
- user: {
- id: parent.id,
- sessionID: chat.id,
- role: "user",
- time: parent.time,
- agent: parent.agent,
- model: { providerID: ref.providerID, modelID: ref.modelID },
- } satisfies SessionV1.User,
- sessionID: chat.id,
- model: mdl,
- agent: agent(),
- system: [],
- messages: [{ role: "user", content: "json" }],
- tools: {},
- })
- expect(value).toBe("stop")
- expect(yield* llm.calls).toBe(1)
- expect(handle.message.error?.name).toBe("APIError")
- }),
- { config: (url) => providerCfg(url) },
- ),
- )
- it.live("session.processor effect tests retry recognized structured json errors", () =>
- provideTmpdirServer(
- ({ dir, llm }) =>
- Effect.gen(function* () {
- const { processors, session, provider } = yield* boot()
- yield* llm.error(429, { type: "error", error: { type: "too_many_requests" } })
- yield* llm.text("after")
- const chat = yield* session.create({})
- const parent = yield* user(chat.id, "retry json")
- const msg = yield* assistant(chat.id, parent.id, path.resolve(dir))
- const mdl = yield* provider.getModel(ref.providerID, ref.modelID)
- const handle = yield* processors.create({
- assistantMessage: msg,
- sessionID: chat.id,
- model: mdl,
- })
- const value = yield* handle.process({
- user: {
- id: parent.id,
- sessionID: chat.id,
- role: "user",
- time: parent.time,
- agent: parent.agent,
- model: { providerID: ref.providerID, modelID: ref.modelID },
- } satisfies SessionV1.User,
- sessionID: chat.id,
- model: mdl,
- agent: agent(),
- system: [],
- messages: [{ role: "user", content: "retry json" }],
- tools: {},
- })
- const parts = yield* MessageV2.parts(msg.id)
- expect(value).toBe("continue")
- expect(yield* llm.calls).toBe(2)
- expect(parts.some((part) => part.type === "text" && part.text === "after")).toBe(true)
- expect(handle.message.error).toBeUndefined()
- }),
- { config: (url) => providerCfg(url) },
- ),
- )
- it.live("session.processor effect tests publish retry status updates", () =>
- provideTmpdirServer(
- ({ dir, llm }) =>
- Effect.gen(function* () {
- const { processors, session, provider } = yield* boot()
- const events = yield* EventV2Bridge.Service
- yield* llm.error(503, { error: "boom" })
- yield* llm.text("")
- const chat = yield* session.create({})
- const parent = yield* user(chat.id, "retry")
- const msg = yield* assistant(chat.id, parent.id, path.resolve(dir))
- const mdl = yield* provider.getModel(ref.providerID, ref.modelID)
- const states: number[] = []
- const off = yield* events.listen((evt) => {
- if (evt.type !== SessionStatus.Event.Status.type) return Effect.void
- const data = evt.data as typeof SessionStatus.Event.Status.data.Type
- if (data.sessionID === chat.id && data.status.type === "retry") states.push(data.status.attempt)
- return Effect.void
- })
- const handle = yield* processors.create({
- assistantMessage: msg,
- sessionID: chat.id,
- model: mdl,
- })
- const value = yield* handle.process({
- user: {
- id: parent.id,
- sessionID: chat.id,
- role: "user",
- time: parent.time,
- agent: parent.agent,
- model: { providerID: ref.providerID, modelID: ref.modelID },
- } satisfies SessionV1.User,
- sessionID: chat.id,
- model: mdl,
- agent: agent(),
- system: [],
- messages: [{ role: "user", content: "retry" }],
- tools: {},
- })
- yield* off
- expect(value).toBe("continue")
- expect(yield* llm.calls).toBe(2)
- expect(states).toStrictEqual([1])
- }),
- { config: (url) => providerCfg(url) },
- ),
- )
- it.live("session.processor effect tests compact on structured context overflow", () =>
- provideTmpdirServer(
- ({ dir, llm }) =>
- Effect.gen(function* () {
- const { processors, session, provider } = yield* boot()
- yield* llm.error(400, { type: "error", error: { code: "context_length_exceeded" } })
- const chat = yield* session.create({})
- const parent = yield* user(chat.id, "compact json")
- const msg = yield* assistant(chat.id, parent.id, path.resolve(dir))
- const mdl = yield* provider.getModel(ref.providerID, ref.modelID)
- const handle = yield* processors.create({
- assistantMessage: msg,
- sessionID: chat.id,
- model: mdl,
- })
- const value = yield* handle.process({
- user: {
- id: parent.id,
- sessionID: chat.id,
- role: "user",
- time: parent.time,
- agent: parent.agent,
- model: { providerID: ref.providerID, modelID: ref.modelID },
- } satisfies SessionV1.User,
- sessionID: chat.id,
- model: mdl,
- agent: agent(),
- system: [],
- messages: [{ role: "user", content: "compact json" }],
- tools: {},
- })
- expect(value).toBe("compact")
- expect(yield* llm.calls).toBe(1)
- expect(handle.message.error).toBeUndefined()
- }),
- { config: (url) => providerCfg(url) },
- ),
- )
- it.live("session.processor effect tests complete AI SDK tool calls when native flag is off", () =>
- provideTmpdirServer(
- ({ dir, llm }) =>
- Effect.gen(function* () {
- const { processors, session, provider } = yield* boot()
- yield* llm.tool("lookup", { query: "weather" })
- const chat = yield* session.create({})
- const parent = yield* user(chat.id, "tool")
- const msg = yield* assistant(chat.id, parent.id, path.resolve(dir))
- const mdl = yield* provider.getModel(ref.providerID, ref.modelID)
- const handle = yield* processors.create({
- assistantMessage: msg,
- sessionID: chat.id,
- model: mdl,
- })
- const value = yield* handle.process({
- user: {
- id: parent.id,
- sessionID: chat.id,
- role: "user",
- time: parent.time,
- agent: parent.agent,
- model: { providerID: ref.providerID, modelID: ref.modelID },
- } satisfies SessionV1.User,
- sessionID: chat.id,
- model: mdl,
- agent: agent(),
- system: [],
- messages: [{ role: "user", content: "tool" }],
- tools: {
- lookup: tool({
- description: "Look up information",
- inputSchema: z.object({ query: z.string() }),
- execute: async (input) => ({
- title: "Weather lookup",
- output: `result:${input.query}`,
- metadata: { source: "test" },
- }),
- }),
- },
- })
- const parts = yield* MessageV2.parts(msg.id)
- const call = parts.find((part): part is SessionV1.ToolPart => part.type === "tool")
- expect(value).toBe("continue")
- expect(yield* llm.calls).toBe(1)
- expect(call?.callID).toBe("call_1")
- expect(call?.tool).toBe("lookup")
- expect(call?.state.status).toBe("completed")
- if (call?.state.status !== "completed") return
- expect(call.state.input).toEqual({ query: "weather" })
- expect(call.state.output).toBe("result:weather")
- expect(call.state.title).toBe("Weather lookup")
- expect(call.state.metadata).toEqual({ source: "test" })
- expect(call.state.time.start).toBeDefined()
- expect(call.state.time.end).toBeDefined()
- }),
- { config: (url) => providerCfg(url) },
- ),
- )
- it.live("session.processor effect tests mark pending tools as aborted on cleanup", () =>
- provideTmpdirServer(
- ({ dir, llm }) =>
- Effect.gen(function* () {
- const database = yield* Database.Service
- const { processors, session, provider } = yield* boot()
- yield* llm.toolHang("bash", { cmd: "pwd" })
- const chat = yield* session.create({})
- const parent = yield* user(chat.id, "tool abort")
- const msg = yield* assistant(chat.id, parent.id, path.resolve(dir))
- const mdl = yield* provider.getModel(ref.providerID, ref.modelID)
- const handle = yield* processors.create({
- assistantMessage: msg,
- sessionID: chat.id,
- model: mdl,
- })
- const run = yield* handle
- .process({
- user: {
- id: parent.id,
- sessionID: chat.id,
- role: "user",
- time: parent.time,
- agent: parent.agent,
- model: { providerID: ref.providerID, modelID: ref.modelID },
- } satisfies SessionV1.User,
- sessionID: chat.id,
- model: mdl,
- agent: agent(),
- system: [],
- messages: [{ role: "user", content: "tool abort" }],
- tools: {},
- })
- .pipe(Effect.forkChild)
- yield* llm.wait(1)
- yield* waitFor(
- MessageV2.parts(msg.id).pipe(
- Effect.map((parts) => parts.find((part): part is SessionV1.ToolPart => part.type === "tool")),
- Effect.provideService(Database.Service, database),
- ),
- "timed out waiting for tool part",
- )
- yield* Fiber.interrupt(run)
- const exit = yield* Fiber.await(run)
- const parts = yield* MessageV2.parts(msg.id)
- const call = parts.find((part): part is SessionV1.ToolPart => part.type === "tool")
- expect(Exit.isFailure(exit)).toBe(true)
- if (Exit.isFailure(exit)) {
- expect(Cause.hasInterruptsOnly(exit.cause)).toBe(true)
- }
- expect(yield* llm.calls).toBe(1)
- expect(call?.state.status).toBe("error")
- if (call?.state.status === "error") {
- expect(call.state.error).toBe("Tool execution aborted")
- expect(call.state.metadata?.interrupted).toBe(true)
- expect(call.state.time.end).toBeDefined()
- }
- }),
- { config: (url) => providerCfg(url) },
- ),
- )
- it.live("session.processor effect tests record aborted errors and idle state", () =>
- provideTmpdirServer(
- ({ dir, llm }) =>
- Effect.gen(function* () {
- const seen = defer<void>()
- const { processors, session, provider } = yield* boot()
- const events = yield* EventV2Bridge.Service
- const sts = yield* SessionStatus.Service
- yield* llm.hang
- const chat = yield* session.create({})
- const parent = yield* user(chat.id, "abort")
- const msg = yield* assistant(chat.id, parent.id, path.resolve(dir))
- const mdl = yield* provider.getModel(ref.providerID, ref.modelID)
- const errs: string[] = []
- const off = yield* events.listen((evt) => {
- if (evt.type !== Session.Event.Error.type) return Effect.void
- const data = evt.data as typeof Session.Event.Error.data.Type
- if (data.sessionID !== chat.id || !data.error) return Effect.void
- errs.push(data.error.name)
- seen.resolve()
- return Effect.void
- })
- const handle = yield* processors.create({
- assistantMessage: msg,
- sessionID: chat.id,
- model: mdl,
- })
- const run = yield* handle
- .process({
- user: {
- id: parent.id,
- sessionID: chat.id,
- role: "user",
- time: parent.time,
- agent: parent.agent,
- model: { providerID: ref.providerID, modelID: ref.modelID },
- } satisfies SessionV1.User,
- sessionID: chat.id,
- model: mdl,
- agent: agent(),
- system: [],
- messages: [{ role: "user", content: "abort" }],
- tools: {},
- })
- .pipe(Effect.forkChild)
- yield* llm.wait(1)
- yield* Fiber.interrupt(run)
- const exit = yield* Fiber.await(run)
- yield* Effect.promise(() => seen.promise)
- const stored = yield* MessageV2.get({ sessionID: chat.id, messageID: msg.id })
- const state = yield* sts.get(chat.id)
- yield* off
- expect(Exit.isFailure(exit)).toBe(true)
- if (Exit.isFailure(exit)) {
- expect(Cause.hasInterruptsOnly(exit.cause)).toBe(true)
- }
- expect(handle.message.error?.name).toBe("MessageAbortedError")
- expect(stored.info.role).toBe("assistant")
- if (stored.info.role === "assistant") {
- expect(stored.info.error?.name).toBe("MessageAbortedError")
- }
- expect(state).toMatchObject({ type: "idle" })
- expect(errs).toContain("MessageAbortedError")
- }),
- { config: (url) => providerCfg(url) },
- ),
- )
- it.live("session.processor effect tests mark interruptions aborted without manual abort", () =>
- provideTmpdirServer(
- ({ dir, llm }) =>
- Effect.gen(function* () {
- const { processors, session, provider } = yield* boot()
- const sts = yield* SessionStatus.Service
- yield* llm.hang
- const chat = yield* session.create({})
- const parent = yield* user(chat.id, "interrupt")
- const msg = yield* assistant(chat.id, parent.id, path.resolve(dir))
- const mdl = yield* provider.getModel(ref.providerID, ref.modelID)
- const handle = yield* processors.create({
- assistantMessage: msg,
- sessionID: chat.id,
- model: mdl,
- })
- const run = yield* handle
- .process({
- user: {
- id: parent.id,
- sessionID: chat.id,
- role: "user",
- time: parent.time,
- agent: parent.agent,
- model: { providerID: ref.providerID, modelID: ref.modelID },
- } satisfies SessionV1.User,
- sessionID: chat.id,
- model: mdl,
- agent: agent(),
- system: [],
- messages: [{ role: "user", content: "interrupt" }],
- tools: {},
- })
- .pipe(Effect.forkChild)
- yield* llm.wait(1)
- yield* Fiber.interrupt(run)
- const exit = yield* Fiber.await(run)
- const stored = yield* MessageV2.get({ sessionID: chat.id, messageID: msg.id })
- const state = yield* sts.get(chat.id)
- expect(Exit.isFailure(exit)).toBe(true)
- expect(handle.message.error?.name).toBe("MessageAbortedError")
- expect(stored.info.role).toBe("assistant")
- if (stored.info.role === "assistant") {
- expect(stored.info.error?.name).toBe("MessageAbortedError")
- }
- expect(state).toMatchObject({ type: "idle" })
- }),
- { config: (url) => providerCfg(url) },
- ),
- )
- itProviderError.live("session.processor effect tests fail provider-executed error results", () =>
- provideTmpdirInstance(
- (dir) =>
- Effect.gen(function* () {
- const { processors, session, provider } = yield* boot()
- const events = yield* EventV2Bridge.Service
- const chat = yield* session.create({})
- const parent = yield* user(chat.id, "provider tool error")
- const msg = yield* assistant(chat.id, parent.id, path.resolve(dir))
- const mdl = yield* provider.getModel(ref.providerID, ref.modelID)
- const settlements: Array<typeof SessionEvent.Tool.Failed.Type> = []
- const off = yield* events.listen((event) => {
- if (event.type === SessionEvent.Tool.Failed.type)
- settlements.push(event as typeof SessionEvent.Tool.Failed.Type)
- return Effect.void
- })
- const handle = yield* processors.create({ assistantMessage: msg, sessionID: chat.id, model: mdl })
- yield* handle.process({
- user: {
- id: parent.id,
- sessionID: chat.id,
- role: "user",
- time: parent.time,
- agent: parent.agent,
- model: { providerID: ref.providerID, modelID: ref.modelID },
- } satisfies SessionV1.User,
- sessionID: chat.id,
- model: mdl,
- agent: agent(),
- system: [],
- messages: [{ role: "user", content: "provider tool error" }],
- tools: {},
- })
- yield* off
- const parts = yield* MessageV2.parts(msg.id)
- const call = parts.find((part): part is SessionV1.ToolPart => part.type === "tool")
- expect(call?.state.status).toBe("error")
- if (call?.state.status === "error") expect(call.state.error).toBe("provider boom")
- expect(settlements).toHaveLength(1)
- expect(settlements[0]?.data).toMatchObject({
- callID: "call-1",
- error: { type: "unknown", message: "provider boom" },
- result: { type: "error", value: "provider boom" },
- provider: { executed: true },
- })
- }),
- { config: cfg },
- ),
- )
- itFragmentFailure.live("session.processor effect tests flush partial v2 fragments before step failure", () =>
- provideTmpdirInstance(
- (dir) =>
- Effect.gen(function* () {
- const { processors, session, provider } = yield* boot()
- const events = yield* EventV2Bridge.Service
- const chat = yield* session.create({})
- const parent = yield* user(chat.id, "provider failure")
- const msg = yield* assistant(chat.id, parent.id, path.resolve(dir))
- const mdl = yield* provider.getModel(ref.providerID, ref.modelID)
- const seen: string[] = []
- let text: string | undefined
- let reasoning: string | undefined
- const off = yield* events.listen((event) => {
- seen.push(event.type)
- if (event.type === SessionEvent.Text.Ended.type)
- text = (event.data as typeof SessionEvent.Text.Ended.data.Type).text
- if (event.type === SessionEvent.Reasoning.Ended.type)
- reasoning = (event.data as typeof SessionEvent.Reasoning.Ended.data.Type).text
- return Effect.void
- })
- const handle = yield* processors.create({ assistantMessage: msg, sessionID: chat.id, model: mdl })
- expect(
- yield* handle.process({
- user: {
- id: parent.id,
- sessionID: chat.id,
- role: "user",
- time: parent.time,
- agent: parent.agent,
- model: { providerID: ref.providerID, modelID: ref.modelID },
- } satisfies SessionV1.User,
- sessionID: chat.id,
- model: mdl,
- agent: agent(),
- system: [],
- messages: [{ role: "user", content: "provider failure" }],
- tools: {},
- }),
- ).toBe("stop")
- yield* off
- const failed = seen.indexOf(SessionEvent.Step.Failed.type)
- expect(failed).toBeGreaterThan(-1)
- expect(seen.indexOf(SessionEvent.Text.Ended.type)).toBeLessThan(failed)
- expect(seen.indexOf(SessionEvent.Reasoning.Ended.type)).toBeLessThan(failed)
- expect(text).toBe("partial")
- expect(reasoning).toBe("thinking")
- }),
- { config: cfg },
- ),
- )
|