diff --git a/packages/ai/package.json b/packages/ai/package.json index f83d89430f..c8a16bbf7c 100644 --- a/packages/ai/package.json +++ b/packages/ai/package.json @@ -15,6 +15,7 @@ ], "exports": { ".": "./src/index.ts", + "./testing": "./src/testing.ts", "./*": "./src/*.ts" }, "devDependencies": { diff --git a/packages/ai/src/testing.ts b/packages/ai/src/testing.ts new file mode 100644 index 0000000000..69619b96e5 --- /dev/null +++ b/packages/ai/src/testing.ts @@ -0,0 +1,157 @@ +export * as TestLLM from "./testing" + +import { LLMClient, type Interface as LLMClientShape } from "./route/client" +import { + LLMEvent, + LLMResponse, + type FinishReasonDetails, + type LLMError, + type LLMRequest, + type UsageInput, +} from "./schema" +import { Context, Deferred, Effect, Latch, Layer, Queue, Scope, Stream } from "effect" + +export type Response = readonly LLMEvent[] | Stream.Stream + +export type Gate = Readonly<{ started: Effect.Effect; release: Effect.Effect }> + +export interface Interface { + readonly requests: LLMRequest[] + readonly push: (...responses: readonly Response[]) => Effect.Effect + readonly always: (response: Response) => Effect.Effect + readonly wait: (count: number) => Effect.Effect + readonly gate: Effect.Effect + readonly client: LLMClientShape +} + +export interface LayerOptions { + readonly transformRequest?: (request: LLMRequest) => LLMRequest + /** Used after the one-shot response queue is exhausted. Omit to defect on unexpected requests. */ + readonly fallback?: Response +} + +export class Service extends Context.Service()("@opencode/ai/TestLLM") {} + +export const complete = ( + options: { readonly reason: FinishReasonDetails; readonly usage?: UsageInput }, + ...events: readonly LLMEvent[] +) => [ + LLMEvent.stepStart({ index: 0 }), + ...events, + LLMEvent.stepFinish({ index: 0, reason: options.reason, usage: options.usage }), + LLMEvent.finish({ reason: options.reason }), +] + +export const stop = (...events: readonly LLMEvent[]) => complete({ reason: { normalized: "stop" } }, ...events) + +export const toolCalls = (...events: readonly LLMEvent[]) => + complete({ reason: { normalized: "tool-calls" } }, ...events) + +const textEvents = (value: string, id: string) => [ + LLMEvent.textStart({ id }), + LLMEvent.textDelta({ id, text: value }), + LLMEvent.textEnd({ id }), +] + +export const text = (value: string, id: string) => stop(...textEvents(value, id)) + +export const textWithUsage = (value: string, id: string, inputTokens: number) => + complete( + { reason: { normalized: "stop" }, usage: { inputTokens, nonCachedInputTokens: inputTokens } }, + ...textEvents(value, id), + ) + +export const tool = (id: string, name: string, input: unknown) => toolCalls(LLMEvent.toolCall({ id, name, input })) + +export const failAfter = (error: LLMError, ...events: readonly LLMEvent[]) => + Stream.fromIterable(events).pipe(Stream.concat(Stream.fail(error))) + +export const hangAfter = (...events: readonly LLMEvent[]) => Stream.concat(Stream.fromIterable(events), Stream.never) + +const toStream = (response: Response) => (Stream.isStream(response) ? response : Stream.fromIterable(response)) + +export const layer = (options: LayerOptions = {}) => + Layer.effect( + Service, + Effect.gen(function* () { + const requests: LLMRequest[] = [] + const responses: Response[] = [] + let started = Deferred.makeUnsafe() + let fallback = options.fallback + let activeGate: { readonly started: Queue.Queue; readonly release: Latch.Latch } | undefined + const wait = (count: number): Effect.Effect => + Effect.suspend(() => + requests.length >= count ? Effect.void : Deferred.await(started).pipe(Effect.andThen(wait(count))), + ) + + const stream = ((request: LLMRequest) => { + requests.push(options.transformRequest?.(request) ?? request) + const waiting = started + started = Deferred.makeUnsafe() + Deferred.doneUnsafe(waiting, Effect.void) + const response = responses.shift() ?? fallback + if (!response) return Stream.die(new Error(`TestLLM has no response for request ${requests.length}`)) + const streamed = toStream(response) + const gate = activeGate + if (!gate) return streamed + return Stream.unwrap( + Queue.offer(gate.started, undefined).pipe(Effect.andThen(gate.release.await), Effect.as(streamed)), + ) + }) as LLMClientShape["stream"] + const client = LLMClient.Service.of({ + prepare: () => Effect.die("TestLLM does not prepare provider-native requests"), + stream, + generate: (request) => + stream(request).pipe( + Stream.runFold(LLMResponse.empty, LLMResponse.reduce), + Effect.flatMap((state) => { + const response = LLMResponse.complete(state) + if (response) return Effect.succeed(response) + return Effect.die("TestLLM response ended without a terminal finish event") + }), + ), + }) + + return Service.of({ + requests, + push: (...input) => + Effect.sync(() => { + responses.push(...input) + }), + always: (response) => + Effect.sync(() => { + fallback = response + }), + wait, + gate: Effect.gen(function* () { + const gate = { + started: yield* Effect.acquireRelease(Queue.unbounded(), Queue.shutdown), + release: yield* Latch.make(), + } + activeGate = gate + const release = Effect.sync(() => { + if (activeGate === gate) activeGate = undefined + }).pipe(Effect.andThen(gate.release.open), Effect.asVoid) + yield* Effect.addFinalizer(() => release) + return { + started: Queue.take(gate.started), + release, + } + }), + client, + }) + }), + ) + +export const clientLayer = Layer.effect( + LLMClient.Service, + Effect.map(Service, (service) => service.client), +) + +export const push = (...responses: readonly Response[]) => Service.use((service) => service.push(...responses)) + +export const always = (response: Response) => Service.use((service) => service.always(response)) + +export const wait = (count: number) => Service.use((service) => service.wait(count)) + +export const gate = Service.use((service) => service.gate) diff --git a/packages/ai/test/exports.test.ts b/packages/ai/test/exports.test.ts index bea36c4c3b..cb32528634 100644 --- a/packages/ai/test/exports.test.ts +++ b/packages/ai/test/exports.test.ts @@ -19,6 +19,7 @@ import { OpenResponses, } from "@opencode-ai/ai/protocols" import * as AnthropicMessages from "@opencode-ai/ai/protocols/anthropic-messages" +import { TestLLM } from "@opencode-ai/ai/testing" describe("public exports", () => { test("root exposes app-facing runtime APIs", () => { @@ -28,6 +29,7 @@ describe("public exports", () => { expect(ImageInput.bytes).toBeFunction() expect(Provider.make).toBeFunction() expect(ProviderSubpath.make).toBe(Provider.make) + expect(TestLLM.layer).toBeFunction() }) test("route barrel exposes route-authoring APIs", () => { diff --git a/packages/core/test/generate.test.ts b/packages/core/test/generate.test.ts index b669e60f48..f90d83ed1b 100644 --- a/packages/core/test/generate.test.ts +++ b/packages/core/test/generate.test.ts @@ -1,6 +1,7 @@ import { expect } from "bun:test" -import { LLMClient, LLMEvent, LLMResponse, Model } from "@opencode-ai/ai" +import { Model } from "@opencode-ai/ai" import { OpenAIChat } from "@opencode-ai/ai/protocols" +import { TestLLM } from "@opencode-ai/ai/testing" import { AISDK } from "@opencode-ai/core/aisdk" import { Catalog } from "@opencode-ai/core/catalog" import { Generate } from "@opencode-ai/core/generate" @@ -9,7 +10,7 @@ import { ModelResolver } from "@opencode-ai/core/model-resolver" import { ID, Info, Ref } from "@opencode-ai/core/model" import { Provider } from "@opencode-ai/core/provider" import { Npm } from "@opencode-ai/util/npm" -import { Effect, Layer, Stream } from "effect" +import { Effect, Layer } from "effect" import { testEffect } from "./lib/effect" const selected = Info.make({ @@ -64,21 +65,7 @@ const aisdk = Layer.mock(AISDK.Service, { }, model: () => Effect.succeed(runtime), }) -const client = Layer.mock(LLMClient.Service)({ - prepare: () => Effect.die("unused"), - stream: () => Stream.die("unused"), - generate: () => - Effect.sync(() => { - const response = LLMResponse.fromEvents([ - LLMEvent.textStart({ id: "generate" }), - LLMEvent.textDelta({ id: "generate", text: "OK" }), - LLMEvent.textEnd({ id: "generate" }), - LLMEvent.finish({ reason: { normalized: "stop" } }), - ]) - if (!response) throw new Error("Incomplete generate response") - return response - }), -}) +const client = TestLLM.clientLayer.pipe(Layer.provide(TestLLM.layer({ fallback: TestLLM.text("OK", "generate") }))) const resolver = ModelResolver.layer.pipe(Layer.provide(Layer.mergeAll(catalog, integrations, npm, aisdk))) const it = testEffect(Generate.layer.pipe(Layer.provide(Layer.merge(resolver, client)))) diff --git a/packages/core/test/session-runner.test.ts b/packages/core/test/session-runner.test.ts index 13f1ed5c3b..3fc569c2f9 100644 --- a/packages/core/test/session-runner.test.ts +++ b/packages/core/test/session-runner.test.ts @@ -1,8 +1,8 @@ import { describe, expect, test } from "bun:test" import { - LLMClient, LLMError, LLMEvent, + LLMRequest, Message, Model, SystemPart, @@ -11,10 +11,9 @@ import { InvalidProviderOutputReason, InvalidRequestReason, RateLimitReason, - type LLMClientShape, - type LLMRequest, } from "@opencode-ai/ai" import * as OpenAIChat from "@opencode-ai/ai/protocols/openai-chat" +import { TestLLM } from "@opencode-ai/ai/testing" import { Catalog } from "@opencode-ai/core/catalog" import { Database } from "@opencode-ai/core/database/database" import { makeLocationNode } from "@opencode-ai/util/effect/app-node" @@ -78,77 +77,59 @@ import { agentHost, catalogHost, host } from "./plugin/host" import PROMPT_DEFAULT from "../src/session/runner/prompt/base.txt" import { CodeModeInstructions } from "@opencode-ai/core/codemode/instructions" -const requests: LLMRequest[] = [] +let requests: LLMRequest[] = [] const emptyCodeMode = `\n\n${CodeModeInstructions.render({ total: 0, shown: 0, namespaces: [] })}` -let response: LLMEvent[] = [] -let responses: LLMEvent[][] | undefined -let responseStream: Stream.Stream | undefined -let responseStreams: Stream.Stream[] | undefined -let streamGate: Deferred.Deferred | undefined -let streamStarted: Deferred.Deferred | undefined -let streamFailure: LLMError | undefined -let toolExecutionGate: Deferred.Deferred | undefined -let toolExecutionsStarted: Deferred.Deferred | undefined -let toolExecutionsReady = 5 -let activeToolExecutions = 0 -let maxActiveToolExecutions = 0 -const client = Layer.succeed( - LLMClient.Service, - LLMClient.Service.of({ - prepare: () => Effect.die("unused"), - stream: ((request: LLMRequest) => { - requests.push({ - ...request, - system: request.system.map((part) => ({ - ...part, - text: part.text.replace(emptyCodeMode, ""), - })), - tools: request.tools.filter((tool) => tool.name !== "execute"), - }) - if (responseStreams) return responseStreams.shift() ?? Stream.empty - if (responseStream) { - const stream = responseStream - responseStream = undefined - return stream - } - const bus = streamFailure - ? Stream.fail(streamFailure) - : Stream.fromIterable(responses === undefined ? response : (responses.shift() ?? [])) - if (!streamGate) return bus - return Stream.unwrap( - (streamStarted ? Deferred.succeed(streamStarted, undefined) : Effect.void).pipe( - Effect.andThen(Deferred.await(streamGate)), - Effect.as(bus), - ), - ) - }) as unknown as LLMClientShape["stream"], - generate: () => Effect.die("unused"), - }), -) -const reply = { - stop: () => [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.stepFinish({ index: 0, reason: { normalized: "stop" } }), - LLMEvent.finish({ reason: { normalized: "stop" } }), - ], - text: (text: string, id: string) => fragmentFixture("text", id, [text]).completeEvents, - textWithUsage: (text: string, id: string, inputTokens: number) => - fragmentFixture("text", id, [text]).completeEvents.map((event) => - LLMEvent.is.stepFinish(event) - ? LLMEvent.stepFinish({ - index: event.index, - reason: event.reason, - usage: { inputTokens, nonCachedInputTokens: inputTokens }, - }) - : event, - ), - tool: (id: string, name: string, input: unknown) => [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.toolCall({ id, name, input }), - LLMEvent.stepFinish({ index: 0, reason: { normalized: "tool-calls" } }), - LLMEvent.finish({ reason: { normalized: "tool-calls" } }), - ], +type ToolBarrier = { + readonly count: number + readonly started: Deferred.Deferred + readonly release: Deferred.Deferred + active: number + maxActive: number } +let toolBarrier: ToolBarrier | undefined +const releaseTools = (barrier: ToolBarrier) => + Effect.sync(() => { + if (toolBarrier === barrier) toolBarrier = undefined + }).pipe(Effect.andThen(Deferred.succeed(barrier.release, undefined)), Effect.asVoid) +const blockTools = (count = 1) => + Effect.acquireRelease( + Effect.all({ started: Deferred.make(), release: Deferred.make() }).pipe( + Effect.map((deferreds) => { + const barrier = { count, ...deferreds, active: 0, maxActive: 0 } + toolBarrier = barrier + return barrier + }), + ), + releaseTools, + ).pipe( + Effect.map((barrier) => ({ + started: Deferred.await(barrier.started), + release: releaseTools(barrier), + maxActive: Effect.sync(() => barrier.maxActive), + })), + ) +const awaitToolBarrier = Effect.suspend(() => { + const barrier = toolBarrier + if (!barrier) return Effect.void + barrier.active++ + barrier.maxActive = Math.max(barrier.maxActive, barrier.active) + return (barrier.active === barrier.count ? Deferred.succeed(barrier.started, undefined) : Effect.void).pipe( + Effect.andThen(Deferred.await(barrier.release)), + Effect.ensuring(Effect.sync(() => barrier.active--)), + ) +}) +const testLLM = TestLLM.layer({ + fallback: [], + transformRequest: (request) => + LLMRequest.update(request, { + system: request.system.map((part) => ({ + ...part, + text: part.text.replace(emptyCodeMode, ""), + })), + tools: request.tools.filter((tool) => tool.name !== "execute"), + }), +}) +const client = TestLLM.clientLayer const model = Model.make({ id: "fake-model", provider: "fake", route: OpenAIChat.route }) const defaultSystem = PROMPT_DEFAULT const replacementModel = Model.make({ id: "replacement", provider: "fake", route: OpenAIChat.route }) @@ -221,7 +202,7 @@ test("does not apply an ineligible tier without base pricing", () => { const authorizations: Tool.Context[] = [] const executions: string[] = [] -const permissionFail = ({ +const permissionFail = { name: "permission_fail", description: "Reject a permission", input: Schema.Struct({}), @@ -235,7 +216,7 @@ const permissionFail = ({ resources: ["src/index.ts"], }), }), -}) +} const permission = Layer.succeed( Permission.Service, Permission.Service.of({ @@ -247,11 +228,7 @@ const permission = Layer.succeed( list: () => Effect.die("unused"), }), ) -const transformTools = ( - registry: Tool.Interface, - tools: Readonly>, - options?: Tool.Options, -) => +const transformTools = (registry: Tool.Interface, tools: Readonly>, options?: Tool.Options) => registry.transform((draft) => Object.entries(tools).forEach(([name, tool]) => draft.add({ ...tool, name, options: { ...tool.options, ...options } }), @@ -259,9 +236,10 @@ const transformTools = ( ) const echo = Layer.effectDiscard( Tool.Service.use((registry) => - transformTools(registry, + transformTools( + registry, { - echo: ({ + echo: { name: "echo", description: "Echo text", input: Schema.Struct({ text: Schema.String }), @@ -270,32 +248,24 @@ const echo = Layer.effectDiscard( Effect.gen(function* () { authorizations.push(context) executions.push(text) - activeToolExecutions++ - maxActiveToolExecutions = Math.max(maxActiveToolExecutions, activeToolExecutions) - if (activeToolExecutions === toolExecutionsReady && toolExecutionsStarted) { - yield* Deferred.succeed(toolExecutionsStarted, undefined) - } - if (toolExecutionGate) yield* Deferred.await(toolExecutionGate) + yield* awaitToolBarrier return { output: { text }, content: text } - }).pipe(Effect.ensuring(Effect.sync(() => activeToolExecutions--))), - }), - defect: ({ + }), + }, + defect: { name: "defect", description: "Fail unexpectedly", input: Schema.Struct({}), output: Schema.Struct({}), - execute: () => - (toolExecutionGate ? Deferred.await(toolExecutionGate) : Effect.void).pipe( - Effect.andThen(Effect.die("unexpected tool defect")), - ), - }), - storefail: ({ + execute: () => awaitToolBarrier.pipe(Effect.andThen(Effect.die("unexpected tool defect"))), + }, + storefail: { name: "storefail", description: "Produce output that cannot be persisted", input: Schema.Struct({}), output: Schema.Struct({}), execute: () => Effect.succeed({ output: {} }), - }), + }, }, { codemode: false }, ), @@ -469,11 +439,16 @@ const it = testEffect( [Config.node, config], [PluginSupervisor.node, pluginSupervisor], ], - ), + ).pipe(Layer.provideMerge(testLLM)), ) const sessionID = Session.ID.make("ses_runner_test") const otherSessionID = Session.ID.make("ses_runner_other") const admit = (session: Session.Interface, text: string) => session.prompt({ sessionID, text, resume: false }) +const runPrompt = Effect.fnUntraced(function* (session: Session.Interface, text: string) { + const message = yield* admit(session, text) + yield* session.resume(sessionID) + return message +}) const insertSession = (id: Session.ID) => Effect.gen(function* () { @@ -506,10 +481,9 @@ const setup = Effect.gen(function* () { yield* Effect.forEach(SystemPromptPlugin.Plugins, (plugin) => plugin.effect(pluginHost), { discard: true, }) - requests.length = 0 + requests = (yield* TestLLM.Service).requests authorizations.length = 0 executions.length = 0 - response = [] systemBaseline = "Initial context" systemRemoved = false systemUnavailable = false @@ -518,17 +492,7 @@ const setup = Effect.gen(function* () { pluginFlushHook = Effect.void currentModel = model skillBaselines.clear() - responses = undefined - streamFailure = undefined - responseStream = undefined - responseStreams = undefined - streamGate = undefined - streamStarted = undefined - toolExecutionGate = undefined - toolExecutionsStarted = undefined - toolExecutionsReady = 5 - activeToolExecutions = 0 - maxActiveToolExecutions = 0 + toolBarrier = undefined yield* agents.transform((draft) => draft.update(Agent.ID.make("build"), (agent) => { agent.mode = "primary" @@ -567,9 +531,8 @@ const rateLimited = (retryAfterMs?: number) => const setupOverflowRecovery = Effect.gen(function* () { const session = yield* setup - response = reply.text("Earlier answer", "text-earlier") - yield* admit(session, "Earlier question ".repeat(700)) - yield* session.resume(sessionID) + yield* TestLLM.push(TestLLM.text("Earlier answer", "text-earlier")) + yield* runPrompt(session, "Earlier question ".repeat(700)) currentModel = recoveryModel requests.length = 0 return session @@ -581,6 +544,7 @@ const messageTexts = (request: LLMRequest, role: "user" | "system") => ) const userTexts = (request: LLMRequest) => messageTexts(request, "user") const systemTexts = (request: LLMRequest) => messageTexts(request, "system") +const messageRoles = (request: LLMRequest | undefined) => request?.messages.map((message) => message.role) const recordedEventTypes = (id: Session.ID) => Effect.gen(function* () { @@ -619,6 +583,9 @@ const recordedStepSettlementEvents = (id: Session.ID, assistantMessageID: Sessio ) }) +const recordedStepSettlementTypes = (id: Session.ID, assistantMessageID: SessionMessage.ID) => + recordedStepSettlementEvents(id, assistantMessageID).pipe(Effect.map((events) => events.map((event) => event.type))) + const hostedCall = (id: string, query: string) => LLMEvent.toolCall({ id, name: "web_search", input: { query }, providerExecuted: true }) @@ -742,7 +709,7 @@ const verifyEphemeralDeltas = (kind: FragmentKind) => const bus = yield* Bus.Service const live = yield* bus.subscribe(fixture.delta).pipe(Stream.take(32), Stream.runCollect, Effect.forkScoped) yield* Effect.yieldNow - response = fixture.completeEvents + yield* TestLLM.push(fixture.completeEvents) yield* session.resume(sessionID) @@ -769,7 +736,7 @@ const verifyPartialFlushOnFailure = (kind: FragmentKind) => const fixture = fragmentFixture(kind, fragmentID(kind, "partial"), ["Partial"]) const failure = providerUnavailable() yield* admit(session, prompt) - responseStream = Stream.concat(Stream.fromIterable(fixture.partialEvents), Stream.fail(failure)) + yield* TestLLM.push(TestLLM.failAfter(failure, ...fixture.partialEvents)) expect(yield* session.resume(sessionID).pipe(Effect.flip)).toBe(failure) expect(yield* session.context(sessionID)).toMatchObject([ @@ -802,9 +769,11 @@ const verifyPartialFlushOnInterruption = (kind: FragmentKind) => const fixture = fragmentFixture(kind, fragmentID(kind, "interrupted"), ["Partial"]) const streamed = yield* Deferred.make() yield* admit(session, prompt) - responseStream = Stream.concat( - Stream.fromIterable(fixture.partialEvents), - Stream.fromEffect(Deferred.succeed(streamed, undefined)).pipe(Stream.flatMap(() => Stream.never)), + yield* TestLLM.push( + Stream.concat( + Stream.fromIterable(fixture.partialEvents), + Stream.fromEffect(Deferred.succeed(streamed, undefined)).pipe(Stream.flatMap(() => Stream.never)), + ), ) const runner = yield* SessionRunner.Service @@ -840,7 +809,7 @@ describe("SessionRunnerLLM", () => { }), ) yield* admit(session, "Original message") - responses = [reply.tool("call-removed", "echo", { text: "blocked" })] + yield* TestLLM.push(TestLLM.tool("call-removed", "echo", { text: "blocked" })) yield* session.resume(sessionID) @@ -872,9 +841,10 @@ describe("SessionRunnerLLM", () => { const session = yield* setup const registry = yield* Tool.Service const contexts: Tool.Context[] = [] - yield* transformTools(registry, + yield* transformTools( + registry, { - location_context: ({ + location_context: { name: "location_context", description: "Read application context", input: Schema.Struct({ query: Schema.String }), @@ -885,12 +855,12 @@ describe("SessionRunnerLLM", () => { yield* context.progress({ phase: "reading" }) return { output: { answer: query.toUpperCase() } } }), - }), + }, }, { codemode: false }, ) yield* admit(session, "Use application context") - responses = [reply.tool("call-location", "location_context", { query: "hello" }), []] + yield* TestLLM.push(TestLLM.tool("call-location", "location_context", { query: "hello" }), []) const bus = yield* Bus.Service const progressFiber = yield* bus.subscribe(SessionEvent.Tool.Progress).pipe( Stream.filter((event) => event.data.sessionID === sessionID && event.data.callID === "call-location"), @@ -934,50 +904,42 @@ describe("SessionRunnerLLM", () => { const registry = yield* Tool.Service const scope = yield* Scope.make() const executions: string[] = [] - yield* transformTools(registry, - { - reloaded: ({ - name: "reloaded", - description: "Record the advertised tool", - input: Schema.Struct({}), - output: Schema.Struct({ value: Schema.String }), - execute: () => - Effect.sync(() => executions.push("advertised")).pipe(Effect.as({ output: { value: "advertised" } })), - }), + yield* transformTools( + registry, + { + reloaded: { + name: "reloaded", + description: "Record the advertised tool", + input: Schema.Struct({}), + output: Schema.Struct({ value: Schema.String }), + execute: () => + Effect.sync(() => executions.push("advertised")).pipe(Effect.as({ output: { value: "advertised" } })), }, - { codemode: false }, - ) - .pipe(Scope.provide(scope)) + }, + { codemode: false }, + ).pipe(Scope.provide(scope)) yield* admit(session, "Use the reloaded tool") - responses = [ - [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.toolCall({ id: "call-reloaded", name: "reloaded", input: {} }), - LLMEvent.stepFinish({ index: 0, reason: { normalized: "tool-calls" } }), - LLMEvent.finish({ reason: { normalized: "tool-calls" } }), - ], - [], - ] - streamGate = yield* Deferred.make() - streamStarted = yield* Deferred.make() + yield* TestLLM.push(TestLLM.tool("call-reloaded", "reloaded", {}), []) + const stream = yield* TestLLM.gate const run = yield* session.resume(sessionID).pipe(Effect.forkChild) - yield* Deferred.await(streamStarted) + yield* stream.started yield* Scope.close(scope, Exit.void) - yield* transformTools(registry, + yield* transformTools( + registry, { - reloaded: ({ + reloaded: { name: "reloaded", description: "Record the replacement tool", input: Schema.Struct({}), output: Schema.Struct({ value: Schema.String }), execute: () => Effect.sync(() => executions.push("replacement")).pipe(Effect.as({ output: { value: "replacement" } })), - }), + }, }, { codemode: false }, ) - yield* Deferred.succeed(streamGate, undefined) + yield* stream.release yield* Fiber.join(run) expect(executions).toEqual(["advertised"]) @@ -1019,16 +981,16 @@ describe("SessionRunnerLLM", () => { const session = yield* setup const secondStarted = yield* Deferred.make() const releaseSecond = yield* Deferred.make() - responseStreams = [ - Stream.fromIterable(reply.tool("call-echo", "echo", { text: "background started" })), + yield* TestLLM.push( + Stream.fromIterable(TestLLM.tool("call-echo", "echo", { text: "background started" })), Stream.unwrap( Deferred.succeed(secondStarted, undefined).pipe( Effect.andThen(Deferred.await(releaseSecond)), - Effect.as(Stream.fromIterable(reply.stop())), + Effect.as(Stream.fromIterable(TestLLM.stop())), ), ), - Stream.fromIterable(reply.text("Handled completion", "text-completion")), - ] + Stream.fromIterable(TestLLM.text("Handled completion", "text-completion")), + ) yield* admit(session, "Start background work") const running = yield* session.resume(sessionID).pipe(Effect.forkChild({ startImmediately: true })) yield* Deferred.await(secondStarted) @@ -1038,7 +1000,7 @@ describe("SessionRunnerLLM", () => { yield* Fiber.join(running) expect(requests).toHaveLength(3) - expect(userTexts(requests[2]!)).toContain("Background work completed") + expect(userTexts(requests[2])).toContain("Background work completed") }), ) @@ -1046,9 +1008,7 @@ describe("SessionRunnerLLM", () => { Effect.gen(function* () { const session = yield* setup yield* admit(session, "First") - yield* admit(session, "Second") - - yield* session.resume(sessionID) + yield* runPrompt(session, "Second") expect(requests).toHaveLength(1) expect(requests[0]?.model).toBe(model) @@ -1071,12 +1031,9 @@ describe("SessionRunnerLLM", () => { if (event.type === "session.instructions.updated") instructionEvents.push(event) }), ) - yield* admit(session, "First") - - yield* session.resume(sessionID) + yield* runPrompt(session, "First") systemBaseline = "Changed context" - yield* admit(session, "Second") - yield* session.resume(sessionID) + yield* runPrompt(session, "Second") yield* unsubscribe expect(instructionEvents).toHaveLength(2) @@ -1113,7 +1070,7 @@ describe("SessionRunnerLLM", () => { yield* session.wait(sessionID) expect(requests).toHaveLength(1) - expect(requests[0]?.messages.map((message) => message.role)).toEqual(["user"]) + expect(messageRoles(requests[0])).toEqual(["user"]) }), ) @@ -1122,8 +1079,7 @@ describe("SessionRunnerLLM", () => { const session = yield* setup const bus = yield* Bus.Service const { db } = yield* Database.Service - yield* admit(session, "First") - yield* session.resume(sessionID) + yield* runPrompt(session, "First") yield* bus.publish(SessionEvent.Moved, { sessionID, @@ -1145,14 +1101,11 @@ describe("SessionRunnerLLM", () => { it.effect("forks instruction values at the selected message instead of the parent's latest state", () => Effect.gen(function* () { const session = yield* setup - const first = yield* admit(session, "First") - yield* session.resume(sessionID) + yield* runPrompt(session, "First") systemBaseline = "Changed context" - const second = yield* admit(session, "Second") - yield* session.resume(sessionID) + const second = yield* runPrompt(session, "Second") systemBaseline = "Latest context" - yield* admit(session, "Third") - yield* session.resume(sessionID) + yield* runPrompt(session, "Third") const forked = yield* session.fork({ sessionID, messageID: second.id }) expect( @@ -1201,11 +1154,9 @@ describe("SessionRunnerLLM", () => { it.effect("caps nested fork instruction ancestry at the selected message", () => Effect.gen(function* () { const session = yield* setup - yield* admit(session, "First") - yield* session.resume(sessionID) + yield* runPrompt(session, "First") systemBaseline = "Changed context" - const second = yield* admit(session, "Second") - yield* session.resume(sessionID) + const second = yield* runPrompt(session, "Second") const child = yield* session.fork({ sessionID, messageID: second.id }) const inheritedFirst = (yield* session.messages({ sessionID: child.id })).find( @@ -1224,6 +1175,7 @@ describe("SessionRunnerLLM", () => { initial_values: { "test/context": Instructions.hash("Initial context") }, current_values: { "test/context": Instructions.hash("Initial context") }, }) + return undefined }), ) @@ -1231,8 +1183,7 @@ describe("SessionRunnerLLM", () => { Effect.gen(function* () { const session = yield* setup const { db } = yield* Database.Service - yield* admit(session, "First") - yield* session.resume(sessionID) + yield* runPrompt(session, "First") yield* db.delete(InstructionStateTable).where(eq(InstructionStateTable.session_id, sessionID)).run() yield* admit(session, "Second") requests.length = 0 @@ -1241,7 +1192,7 @@ describe("SessionRunnerLLM", () => { expect(requests).toHaveLength(1) expect(requests[0]?.system.map((part) => part.text)).toEqual([defaultSystem, "Initial context"]) - expect(requests[0]?.messages.map((message) => message.role)).toEqual(["user", "user"]) + expect(messageRoles(requests[0])).toEqual(["user", "user"]) expect( yield* db .select({ id: EventTable.id }) @@ -1259,24 +1210,21 @@ describe("SessionRunnerLLM", () => { it.effect("keeps the initial instructions stable and derives a chronological update from values", () => Effect.gen(function* () { const session = yield* setup - yield* admit(session, "First") - - yield* session.resume(sessionID) + yield* runPrompt(session, "First") systemBaseline = "Changed context" - yield* admit(session, "Second") - yield* session.resume(sessionID) + yield* runPrompt(session, "Second") expect( PromptCacheDiagnostics.compare( - PromptCacheDiagnostics.snapshot(requests[0]!), - PromptCacheDiagnostics.snapshot(requests[1]!), + PromptCacheDiagnostics.snapshot(requests[0]), + PromptCacheDiagnostics.snapshot(requests[1]), ), ).toEqual({ status: "append-only", previousMessages: 1, currentMessages: 3 }) expect(requests.map((request) => request.system.map((part) => part.text))).toEqual([ [defaultSystem, "Initial context"], [defaultSystem, "Initial context"], ]) - expect(requests[1]?.messages.map((message) => message.role)).toEqual(["user", "system", "user"]) + expect(messageRoles(requests[1])).toEqual(["user", "system", "user"]) expect(requests[1]?.messages.at(1)?.content).toEqual([{ type: "text", text: "Changed context" }]) expect(yield* session.messages({ sessionID })).toHaveLength(2) const { db } = yield* Database.Service @@ -1307,7 +1255,7 @@ describe("SessionRunnerLLM", () => { currentModel = Model.make({ id: "gpt-5", provider: "openai", route: OpenAIChat.route }) yield* admit(session, "First") - response = reply.text("Done", "text-provider-prompt") + yield* TestLLM.push(TestLLM.text("Done", "text-provider-prompt")) yield* session.resume(sessionID) expect(requests.at(-1)?.system.map((part) => part.text)).toEqual([ @@ -1330,7 +1278,7 @@ describe("SessionRunnerLLM", () => { ) yield* admit(session, "First") - response = reply.text("Done", "text-empty-agent-system") + yield* TestLLM.push(TestLLM.text("Done", "text-empty-agent-system")) yield* session.resume(sessionID) expect(requests.at(-1)?.system.map((part) => part.text)).toEqual([ @@ -1352,7 +1300,7 @@ describe("SessionRunnerLLM", () => { ) yield* admit(session, "First") - response = reply.text("Done", "text-build") + yield* TestLLM.push(TestLLM.text("Done", "text-build")) yield* session.resume(sessionID) expect(requests.at(-1)?.system.map((part) => part.text)).toEqual(["Build agent instructions", "Initial context"]) @@ -1376,7 +1324,7 @@ describe("SessionRunnerLLM", () => { }) yield* admit(session, "First") - response = reply.text("Done", "text-reviewer") + yield* TestLLM.push(TestLLM.text("Done", "text-reviewer")) yield* session.resume(sessionID) expect(requests.at(-1)?.system.map((part) => part.text)).toEqual(["Reviewer instructions", "Initial context"]) @@ -1396,7 +1344,7 @@ describe("SessionRunnerLLM", () => { ) yield* admit(session, "First") - response = reply.text("Done", "text-no-system") + yield* TestLLM.push(TestLLM.text("Done", "text-no-system")) yield* session.resume(sessionID) expect(requests.at(-1)?.system.map((part) => part.text)).toEqual(["Build agent instructions", "Initial context"]) @@ -1422,7 +1370,7 @@ describe("SessionRunnerLLM", () => { .pipe(Effect.orDie) yield* admit(session, "First") - response = reply.text("Done", "text-selected") + yield* TestLLM.push(TestLLM.text("Done", "text-selected")) yield* session.resume(sessionID) expect(requests.at(-1)?.system.map((part) => part.text)).toEqual(["Reviewer instructions", "Initial context"]) @@ -1444,7 +1392,7 @@ describe("SessionRunnerLLM", () => { yield* session.prompt({ sessionID, text: "Inspect files", resume: false }) requests.length = 0 - response = [] + yield* TestLLM.push([]) const failure = yield* session.resume(sessionID).pipe(Effect.flip) expect(failure).toMatchObject({ @@ -1465,7 +1413,7 @@ describe("SessionRunnerLLM", () => { yield* session.prompt({ sessionID, text: "Wait for plugins", resume: false }) requests.length = 0 - response = [] + yield* TestLLM.push([]) const running = yield* session.resume(sessionID).pipe(Effect.forkChild({ startImmediately: true })) yield* Effect.yieldNow @@ -1489,22 +1437,19 @@ describe("SessionRunnerLLM", () => { }), ) skillBaselines.set(Agent.ID.make("build"), "Build skills") - yield* admit(session, "First") - - yield* session.resume(sessionID) + yield* runPrompt(session, "First") skillBaselines.set(Agent.ID.make("reviewer"), "Reviewer skills") yield* bus.publish(SessionEvent.AgentSelected, { sessionID, agent: Agent.ID.make("reviewer"), }) - yield* admit(session, "Second") - yield* session.resume(sessionID) + yield* runPrompt(session, "Second") expect(requests.map((request) => request.system.map((part) => part.text))).toEqual([ [defaultSystem, "Initial context\n\nBuild skills"], [defaultSystem, "Initial context\n\nBuild skills"], ]) - expect(systemTexts(requests[1]!)).toContainEqual(expect.stringContaining("Reviewer skills")) + expect(systemTexts(requests[1])).toContainEqual(expect.stringContaining("Reviewer skills")) }), ) @@ -1525,9 +1470,7 @@ describe("SessionRunnerLLM", () => { }) .pipe(Effect.asVoid) }) - yield* admit(session, "First") - - yield* session.resume(sessionID) + yield* runPrompt(session, "First") expect(requests.map((request) => request.system.map((part) => part.text))).toEqual([ [defaultSystem, "Initial context\n\nBuild skills"], @@ -1550,9 +1493,7 @@ describe("SessionRunnerLLM", () => { }) .pipe(Effect.asVoid) }) - yield* admit(session, "First") - - yield* session.resume(sessionID) + yield* runPrompt(session, "First") expect(requests.map((request) => request.model)).toEqual([model]) expect(requests.map((request) => request.system.map((part) => part.text))).toEqual([ [defaultSystem, "Initial context"], @@ -1563,14 +1504,11 @@ describe("SessionRunnerLLM", () => { it.effect("admits removed context as a chronological System message", () => Effect.gen(function* () { const session = yield* setup - yield* admit(session, "First") - - yield* session.resume(sessionID) + yield* runPrompt(session, "First") systemRemoved = true - yield* admit(session, "Second") - yield* session.resume(sessionID) + yield* runPrompt(session, "Second") - expect(requests[1]?.messages.map((message) => message.role)).toEqual(["user", "system", "user"]) + expect(messageRoles(requests[1])).toEqual(["user", "system", "user"]) expect(requests[1]?.messages.at(1)?.content).toEqual([ { type: "text", text: "System context source removed: test/context" }, ]) @@ -1583,9 +1521,7 @@ describe("SessionRunnerLLM", () => { const session = yield* setup const contextEntries = yield* InstructionEntry.Service yield* contextEntries.put({ sessionID, key: "deploy-target", value: "production" }) - yield* admit(session, "First") - - yield* session.resume(sessionID) + yield* runPrompt(session, "First") // String values render verbatim inside the initial tagged block. expect(requests[0]?.system.map((part) => part.text)).toEqual([ @@ -1595,10 +1531,9 @@ describe("SessionRunnerLLM", () => { // Non-string JSON pretty-prints; the change narrates as a System update. yield* contextEntries.put({ sessionID, key: "deploy-target", value: { region: "us-east-1" } }) - yield* admit(session, "Second") - yield* session.resume(sessionID) + yield* runPrompt(session, "Second") - expect(requests[1]?.messages.map((message) => message.role)).toEqual(["user", "system", "user"]) + expect(messageRoles(requests[1])).toEqual(["user", "system", "user"]) expect(requests[1]?.messages.at(1)?.content).toEqual([ { type: "text", @@ -1616,10 +1551,9 @@ describe("SessionRunnerLLM", () => { // Deleting the row announces removal through the stored removal text. yield* contextEntries.remove({ sessionID, key: "deploy-target" }) - yield* admit(session, "Third") - yield* session.resume(sessionID) + yield* runPrompt(session, "Third") - expect(requests[2]?.messages.map((message) => message.role)).toEqual(["user", "system", "user", "system", "user"]) + expect(messageRoles(requests[2])).toEqual(["user", "system", "user", "system", "user"]) expect(requests[2]?.messages.at(-2)?.content).toEqual([ { type: "text", text: 'The context under "deploy-target" no longer applies. Disregard it.' }, ]) @@ -1632,12 +1566,10 @@ describe("SessionRunnerLLM", () => { const session = yield* setup const entries = yield* InstructionEntry.Service yield* entries.put({ sessionID, key: "nullable", value: "present" }) - yield* admit(session, "First") - yield* session.resume(sessionID) + yield* runPrompt(session, "First") yield* entries.put({ sessionID, key: "nullable", value: null }) - yield* admit(session, "Second") - yield* session.resume(sessionID) + yield* runPrompt(session, "Second") expect(requests[1]?.messages.at(1)?.content).toEqual([ { @@ -1673,26 +1605,22 @@ describe("SessionRunnerLLM", () => { Effect.gen(function* () { const session = yield* setup const bus = yield* Bus.Service - yield* admit(session, "First") - - yield* session.resume(sessionID) + yield* runPrompt(session, "First") systemBaseline = "Changed context" - yield* admit(session, "Second") - yield* session.resume(sessionID) + yield* runPrompt(session, "Second") yield* bus.publish(SessionEvent.ModelSelected, { sessionID, model: { id: ID.make("replacement"), providerID: Provider.ID.make("fake") }, }) systemBaseline = "Replacement context" - yield* admit(session, "Third") - yield* session.resume(sessionID) + yield* runPrompt(session, "Third") expect(requests.map((request) => request.system.map((part) => part.text))).toEqual([ [defaultSystem, "Initial context"], [defaultSystem, "Initial context"], [defaultSystem, "Initial context"], ]) - expect(requests[1]?.messages.map((message) => message.role)).toEqual(["user", "system", "user"]) + expect(messageRoles(requests[1])).toEqual(["user", "system", "user"]) expect(requests[2]?.messages.filter((message) => message.role === "system")).toHaveLength(2) expect((yield* session.context(sessionID)).map((message) => message.type)).toEqual([ "user", @@ -1702,8 +1630,7 @@ describe("SessionRunnerLLM", () => { ]) yield* replaySessionProjection(sessionID) expect(yield* session.messages({ sessionID })).toHaveLength(4) - yield* admit(session, "Fourth") - yield* session.resume(sessionID) + yield* runPrompt(session, "Fourth") }), ) @@ -1711,20 +1638,16 @@ describe("SessionRunnerLLM", () => { Effect.gen(function* () { const session = yield* setup const bus = yield* Bus.Service - yield* admit(session, "First") - - yield* session.resume(sessionID) + yield* runPrompt(session, "First") yield* bus.publish(SessionEvent.ModelSelected, { sessionID, model: { id: ID.make("replacement"), providerID: Provider.ID.make("fake") }, }) systemUnavailable = true - yield* admit(session, "Second") - yield* session.resume(sessionID) + yield* runPrompt(session, "Second") systemUnavailable = false systemBaseline = "Replacement context" - yield* admit(session, "Third") - yield* session.resume(sessionID) + yield* runPrompt(session, "Third") expect(requests.map((request) => request.system.map((part) => part.text))).toEqual([ [defaultSystem, "Initial context"], @@ -1738,9 +1661,7 @@ describe("SessionRunnerLLM", () => { Effect.gen(function* () { const session = yield* setup const bus = yield* Bus.Service - yield* admit(session, "First") - - yield* session.resume(sessionID) + yield* runPrompt(session, "First") yield* bus.publish(SessionEvent.Compaction.Started, { sessionID, reason: "manual", @@ -1753,18 +1674,16 @@ describe("SessionRunnerLLM", () => { recent: "", }) systemBaseline = "Replacement context" - yield* admit(session, "Second") - yield* session.resume(sessionID) + yield* runPrompt(session, "Second") expect(requests.map((request) => request.system.map((part) => part.text))).toEqual([ [defaultSystem, "Initial context"], [defaultSystem, "Initial context"], ]) - expect(requests[1]?.messages.map((message) => message.role)).toEqual(["user", "system", "user"]) + expect(messageRoles(requests[1])).toEqual(["user", "system", "user"]) expect(requests[1]?.messages.at(1)?.content).toEqual([{ type: "text", text: "Replacement context" }]) yield* replaySessionProjection(sessionID) - yield* admit(session, "Third") - yield* session.resume(sessionID) + yield* runPrompt(session, "Third") }), ) @@ -1772,17 +1691,16 @@ describe("SessionRunnerLLM", () => { Effect.gen(function* () { const session = yield* setup currentModel = recoveryModel - streamGate = yield* Deferred.make() - streamStarted = yield* Deferred.make() - responses = [ - reply.tool("call-active", "echo", { text: "active" }), + const stream = yield* TestLLM.gate + yield* TestLLM.push( + TestLLM.tool("call-active", "echo", { text: "active" }), [LLMEvent.textDelta({ id: "summary", text: "durable summary" })], - reply.text("Steer complete", "text-steer"), - reply.text("Queue complete", "text-queue"), - ] + TestLLM.text("Steer complete", "text-steer"), + TestLLM.text("Queue complete", "text-queue"), + ) yield* admit(session, "Active work") const active = yield* session.resume(sessionID).pipe(Effect.forkChild) - yield* Deferred.await(streamStarted) + yield* stream.started const first = yield* session.compact({ sessionID }) const second = yield* session.compact({ sessionID }) @@ -1802,7 +1720,7 @@ describe("SessionRunnerLLM", () => { }) expect(yield* SessionPending.has((yield* Database.Service).db, sessionID, "steer")).toBe(false) - yield* Deferred.succeed(streamGate, undefined) + yield* stream.release yield* Fiber.join(active) expect(requests).toHaveLength(4) @@ -1823,16 +1741,15 @@ describe("SessionRunnerLLM", () => { Effect.gen(function* () { const session = yield* setup currentModel = recoveryModel - streamGate = yield* Deferred.make() - streamStarted = yield* Deferred.make() - responses = [ - reply.text("Active complete", "text-active-failure"), + const stream = yield* TestLLM.gate + yield* TestLLM.push( + TestLLM.text("Active complete", "text-active-failure"), [], - reply.text("Continued", "text-after-failure"), - ] + TestLLM.text("Continued", "text-after-failure"), + ) yield* admit(session, "Active work") const active = yield* session.resume(sessionID).pipe(Effect.forkChild) - yield* Deferred.await(streamStarted) + yield* stream.started const compaction = yield* session.compact({ sessionID }) yield* session.prompt({ @@ -1841,7 +1758,7 @@ describe("SessionRunnerLLM", () => { delivery: "queue", resume: false, }) - yield* Deferred.succeed(streamGate, undefined) + yield* stream.release yield* Fiber.join(active) expect(requests).toHaveLength(3) @@ -1886,12 +1803,11 @@ describe("SessionRunnerLLM", () => { it.effect("manually compacts when the model has no context limit", () => Effect.gen(function* () { const session = yield* setup - response = reply.text("Earlier answer", "text-manual-unknown-history") - yield* admit(session, "Earlier question") - yield* session.resume(sessionID) + yield* TestLLM.push(TestLLM.text("Earlier answer", "text-manual-unknown-history")) + yield* runPrompt(session, "Earlier question") requests.length = 0 - response = reply.text("Manual summary", "text-manual-unknown-summary") + yield* TestLLM.push(TestLLM.text("Manual summary", "text-manual-unknown-summary")) const compaction = yield* session.compact({ sessionID }) yield* session.resume(sessionID) @@ -1908,11 +1824,10 @@ describe("SessionRunnerLLM", () => { it.effect("preserves provider errors from manual compaction", () => Effect.gen(function* () { const session = yield* setup - response = reply.text("Earlier answer", "text-manual-provider-history") - yield* admit(session, "Earlier question") - yield* session.resume(sessionID) + yield* TestLLM.push(TestLLM.text("Earlier answer", "text-manual-provider-history")) + yield* runPrompt(session, "Earlier question") - response = [LLMEvent.providerError({ message: "summary unavailable" })] + yield* TestLLM.push([LLMEvent.providerError({ message: "summary unavailable" })]) const compaction = yield* session.compact({ sessionID }) yield* session.resume(sessionID) @@ -1927,11 +1842,10 @@ describe("SessionRunnerLLM", () => { it.effect("preserves typed provider failures from manual compaction", () => Effect.gen(function* () { const session = yield* setup - response = reply.text("Earlier answer", "text-manual-failure-history") - yield* admit(session, "Earlier question") - yield* session.resume(sessionID) + yield* TestLLM.push(TestLLM.text("Earlier answer", "text-manual-failure-history")) + yield* runPrompt(session, "Earlier question") - responseStream = Stream.fail(providerUnavailable()) + yield* TestLLM.push(Stream.fail(providerUnavailable())) const compaction = yield* session.compact({ sessionID }) yield* session.resume(sessionID) @@ -1946,15 +1860,16 @@ describe("SessionRunnerLLM", () => { it.effect("records cancelled manual compaction without surfacing an internal failure", () => Effect.gen(function* () { const session = yield* setup - response = reply.text("Earlier answer", "text-manual-interrupt-history") - yield* admit(session, "Earlier question") - yield* session.resume(sessionID) + yield* TestLLM.push(TestLLM.text("Earlier answer", "text-manual-interrupt-history")) + yield* runPrompt(session, "Earlier question") const streamed = yield* Deferred.make() const partial = fragmentFixture("text", "text-manual-interrupt-summary", ["Partial summary"]) - responseStream = Stream.concat( - Stream.fromIterable(partial.partialEvents), - Stream.fromEffect(Deferred.succeed(streamed, undefined)).pipe(Stream.flatMap(() => Stream.never)), + yield* TestLLM.push( + Stream.concat( + Stream.fromIterable(partial.partialEvents), + Stream.fromEffect(Deferred.succeed(streamed, undefined)).pipe(Stream.flatMap(() => Stream.never)), + ), ) const compaction = yield* session.compact({ sessionID }) const run = yield* session.resume(sessionID).pipe(Effect.forkChild) @@ -1975,9 +1890,8 @@ describe("SessionRunnerLLM", () => { it.effect("settles an admitted manual compaction when pre-start resolution throws", () => Effect.gen(function* () { const session = yield* setup - response = reply.text("Earlier answer", "text-manual-resolution-history") - yield* admit(session, "Earlier question") - yield* session.resume(sessionID) + yield* TestLLM.push(TestLLM.text("Earlier answer", "text-manual-resolution-history")) + yield* runPrompt(session, "Earlier question") const compaction = yield* session.compact({ sessionID }) modelResolveHook = Effect.die("model resolution failed") @@ -2001,18 +1915,16 @@ describe("SessionRunnerLLM", () => { it.effect("automatically compacts into a completed summary and retained recent turn", () => Effect.gen(function* () { const session = yield* setup - response = reply.textWithUsage("Earlier answer", "text-first", 3_950) - yield* admit(session, "Earlier question ".repeat(180)) - yield* session.resume(sessionID) + yield* TestLLM.push(TestLLM.textWithUsage("Earlier answer", "text-first", 3_950)) + yield* runPrompt(session, "Earlier question ".repeat(180)) currentModel = compactModel requests.length = 0 - responses = [ - reply.text("## Objective\n- Preserve the task", "text-summary"), - reply.textWithUsage("Continued", "text-final", 3_950), - ] - yield* admit(session, "Recent exact request ".repeat(180)) - yield* session.resume(sessionID) + yield* TestLLM.push( + TestLLM.text("## Objective\n- Preserve the task", "text-summary"), + TestLLM.textWithUsage("Continued", "text-final", 3_950), + ) + yield* runPrompt(session, "Recent exact request ".repeat(180)) expect(requests).toHaveLength(2) expect(userTexts(requests[0])[0]).toContain("## Objective") @@ -2029,12 +1941,11 @@ describe("SessionRunnerLLM", () => { requests.length = 0 executions.length = 0 - responses = [ - reply.text("## Objective\n- Preserve the updated task", "text-summary-2"), - reply.text("Continued again", "text-final-2"), - ] - yield* admit(session, "Newest exact request ".repeat(180)) - yield* session.resume(sessionID) + yield* TestLLM.push( + TestLLM.text("## Objective\n- Preserve the updated task", "text-summary-2"), + TestLLM.text("Continued again", "text-final-2"), + ) + yield* runPrompt(session, "Newest exact request ".repeat(180)) expect(requests).toHaveLength(2) expect(userTexts(requests[0])[0]).toContain( @@ -2052,14 +1963,12 @@ describe("SessionRunnerLLM", () => { Effect.gen(function* () { const session = yield* setup currentModel = fullOutputModel - response = reply.textWithUsage("Earlier answer", "text-full-output-first", 9_500) - yield* admit(session, "Earlier question") - yield* session.resume(sessionID) + yield* TestLLM.push(TestLLM.textWithUsage("Earlier answer", "text-full-output-first", 9_500)) + yield* runPrompt(session, "Earlier question") requests.length = 0 - response = reply.text("Continued", "text-full-output-final") - yield* admit(session, "Continue") - yield* session.resume(sessionID) + yield* TestLLM.push(TestLLM.text("Continued", "text-full-output-final")) + yield* runPrompt(session, "Continue") expect(requests).toHaveLength(1) expect(userTexts(requests[0])).toContain("Continue") @@ -2070,16 +1979,15 @@ describe("SessionRunnerLLM", () => { it.effect("stops after required automatic compaction fails", () => Effect.gen(function* () { const session = yield* setup - response = reply.textWithUsage("Earlier answer", "text-before-failed-compaction", 3_950) - yield* admit(session, "Earlier question ".repeat(180)) - yield* session.resume(sessionID) + yield* TestLLM.push(TestLLM.textWithUsage("Earlier answer", "text-before-failed-compaction", 3_950)) + yield* runPrompt(session, "Earlier question ".repeat(180)) currentModel = compactModel requests.length = 0 - responses = [ + yield* TestLLM.push( [LLMEvent.providerError({ message: "Unsupported parameter: max_output_tokens" })], - reply.text("Must not run", "text-after-failed-compaction"), - ] + TestLLM.text("Must not run", "text-after-failed-compaction"), + ) yield* admit(session, "Recent exact request ".repeat(180)) expect(yield* Effect.exit(session.resume(sessionID))).toMatchObject({ _tag: "Failure" }) @@ -2099,16 +2007,15 @@ describe("SessionRunnerLLM", () => { it.effect("forces one compaction and retries after provider context overflow", () => Effect.gen(function* () { const session = yield* setupOverflowRecovery - responses = [ + yield* TestLLM.push( [ LLMEvent.stepStart({ index: 0 }), LLMEvent.providerError({ message: "prompt too long", classification: "context-overflow" }), ], - reply.text("## Objective\n- Recover overflow", "text-summary"), - reply.text("Recovered", "text-final"), - ] - yield* admit(session, "Continue") - yield* session.resume(sessionID) + TestLLM.text("## Objective\n- Recover overflow", "text-summary"), + TestLLM.text("Recovered", "text-final"), + ) + yield* runPrompt(session, "Continue") expect(requests).toHaveLength(3) expect(userTexts(requests[1])[0]).toContain("## Objective") @@ -2129,13 +2036,12 @@ describe("SessionRunnerLLM", () => { Effect.gen(function* () { const session = yield* setupOverflowRecovery currentModel = model - responses = [ + yield* TestLLM.push( [LLMEvent.providerError({ message: "prompt too long", classification: "context-overflow" })], - reply.text("## Objective\n- Recover unknown limit", "text-summary-unknown-limit"), - reply.text("Recovered", "text-final-unknown-limit"), - ] - yield* admit(session, "Continue") - yield* session.resume(sessionID) + TestLLM.text("## Objective\n- Recover unknown limit", "text-summary-unknown-limit"), + TestLLM.text("Recovered", "text-final-unknown-limit"), + ) + yield* runPrompt(session, "Continue") expect(requests).toHaveLength(3) expect(yield* session.context(sessionID)).toMatchObject([ @@ -2149,13 +2055,12 @@ describe("SessionRunnerLLM", () => { Effect.gen(function* () { const session = yield* setupOverflowRecovery currentModel = undersizedContextModel - responses = [ + yield* TestLLM.push( [LLMEvent.providerError({ message: "prompt too long", classification: "context-overflow" })], - reply.text("## Objective\n- Recover undersized limit", "text-summary-undersized-limit"), - reply.text("Recovered", "text-final-undersized-limit"), - ] - yield* admit(session, "Continue") - yield* session.resume(sessionID) + TestLLM.text("## Objective\n- Recover undersized limit", "text-summary-undersized-limit"), + TestLLM.text("Recovered", "text-final-undersized-limit"), + ) + yield* runPrompt(session, "Continue") expect(requests).toHaveLength(3) expect(yield* session.context(sessionID)).toMatchObject([ @@ -2172,7 +2077,7 @@ describe("SessionRunnerLLM", () => { LLMEvent.stepStart({ index: 0 }), LLMEvent.providerError({ message: "prompt too long", classification: "context-overflow" }), ] - responses = [overflow(), reply.text("## Objective\n- Recover once", "text-summary"), overflow()] + yield* TestLLM.push(overflow(), TestLLM.text("## Objective\n- Recover once", "text-summary"), overflow()) yield* admit(session, "Continue") expect((yield* session.resume(sessionID).pipe(Effect.flip)).message).toBe("prompt too long") @@ -2187,22 +2092,23 @@ describe("SessionRunnerLLM", () => { it.effect("recovers once from a raw context overflow failure", () => Effect.gen(function* () { const session = yield* setupOverflowRecovery - responseStream = Stream.fail( - new LLMError({ - module: "test", - method: "stream", - reason: new InvalidRequestReason({ - message: "prompt too long", - classification: "context-overflow", + yield* TestLLM.push( + Stream.fail( + new LLMError({ + module: "test", + method: "stream", + reason: new InvalidRequestReason({ + message: "prompt too long", + classification: "context-overflow", + }), }), - }), + ), ) - responses = [ - reply.text("## Objective\n- Recover raw overflow", "text-summary"), - reply.text("Recovered", "text-final"), - ] - yield* admit(session, "Continue") - yield* session.resume(sessionID) + yield* TestLLM.push( + TestLLM.text("## Objective\n- Recover raw overflow", "text-summary"), + TestLLM.text("Recovered", "text-final"), + ) + yield* runPrompt(session, "Continue") expect(requests).toHaveLength(3) expect(yield* session.context(sessionID)).toMatchObject([ @@ -2215,10 +2121,10 @@ describe("SessionRunnerLLM", () => { it.effect("publishes the original overflow when recovery summarization fails", () => Effect.gen(function* () { const session = yield* setupOverflowRecovery - responses = [ + yield* TestLLM.push( [LLMEvent.providerError({ message: "prompt too long", classification: "context-overflow" })], [LLMEvent.providerError({ message: "summary unavailable" })], - ] + ) yield* admit(session, "Continue") expect((yield* session.resume(sessionID).pipe(Effect.flip)).message).toBe("prompt too long") @@ -2243,26 +2149,29 @@ describe("SessionRunnerLLM", () => { it.effect("interrupts overflow recovery while the summary provider is running", () => Effect.gen(function* () { const session = yield* setupOverflowRecovery - responses = [ + yield* TestLLM.push( [LLMEvent.providerError({ message: "prompt too long", classification: "context-overflow" })], - reply.text("## Objective\n- Interrupted", "text-summary"), - ] - const firstGate = yield* Deferred.make() - const summaryGate = yield* Deferred.make() - streamGate = firstGate + TestLLM.text("## Objective\n- Interrupted", "text-summary"), + ) + const first = yield* TestLLM.gate yield* admit(session, "Continue") const run = yield* session.resume(sessionID).pipe(Effect.forkChild) - while (requests.length < 1) yield* Effect.yieldNow - streamGate = summaryGate - yield* Deferred.succeed(firstGate, undefined) - while (requests.length < 2) yield* Effect.yieldNow + yield* first.started + + const summary = yield* TestLLM.gate + yield* first.release + yield* summary.started yield* session.interrupt(sessionID) - expect(yield* Fiber.await(run)).toMatchObject({ _tag: "Failure" }) - streamGate = undefined - expect(requests).toHaveLength(2) + const exit = yield* Fiber.await(run) + expect(Exit.isFailure(exit) && Cause.hasInterruptsOnly(exit.cause)).toBeTrue() expect(yield* session.context(sessionID)).toContainEqual( - expect.objectContaining({ type: "compaction", status: "failed", reason: "auto" }), + expect.objectContaining({ + type: "compaction", + status: "failed", + reason: "auto", + error: { type: "compaction.interrupted", message: "Compaction was interrupted" }, + }), ) }), ) @@ -2271,12 +2180,9 @@ describe("SessionRunnerLLM", () => { Effect.gen(function* () { const session = yield* setup const bus = yield* Bus.Service - yield* admit(session, "First") - - yield* session.resume(sessionID) + yield* runPrompt(session, "First") systemBaseline = "Changed context" - yield* admit(session, "Second") - yield* session.resume(sessionID) + yield* runPrompt(session, "Second") yield* bus.publish(SessionEvent.Compaction.Started, { sessionID, reason: "manual", @@ -2289,8 +2195,7 @@ describe("SessionRunnerLLM", () => { recent: "", }) systemUnavailable = true - yield* admit(session, "Third") - yield* session.resume(sessionID) + yield* runPrompt(session, "Third") // Compaction already moved current values into the new epoch before the unavailable read. expect(requests.at(-1)?.system.map((part) => part.text)).toEqual([defaultSystem, "Changed context"]) @@ -2303,50 +2208,49 @@ describe("SessionRunnerLLM", () => { const session = yield* setup yield* admit(session, "Use tools") - response = [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.reasoningStart({ id: "reasoning-1" }), - LLMEvent.reasoningDelta({ id: "reasoning-1", text: "Think" }), - LLMEvent.reasoningEnd({ id: "reasoning-1" }), - LLMEvent.toolInputStart({ id: "call-error", name: "write" }), - LLMEvent.toolInputDelta({ id: "call-error", name: "write", text: '{"path":"README.md"}' }), - LLMEvent.toolInputEnd({ id: "call-error", name: "write" }), - LLMEvent.toolCall({ id: "call-error", name: "write", input: { path: "README.md" }, providerExecuted: true }), - LLMEvent.toolError({ id: "call-error", name: "write", message: "Denied" }), - LLMEvent.toolResult({ id: "call-error", name: "write", result: { type: "error", value: "Denied" } }), - LLMEvent.toolCall({ - id: "call-provider", - name: "web_search", - input: { query: "hello" }, - providerExecuted: true, - providerMetadata: { openai: { source: "provider" } }, - }), - LLMEvent.toolResult({ - id: "call-provider", - name: "web_search", - result: { - type: "content", - value: [ - { type: "text", text: "Hello" }, - { type: "file", uri: "data:image/png;base64,aGVsbG8=", mime: "image/png", name: "hello.png" }, - ], + yield* TestLLM.push( + TestLLM.complete( + { + reason: { normalized: "tool-calls" }, + usage: { + inputTokens: 10, + nonCachedInputTokens: 8, + outputTokens: 4, + reasoningTokens: 1, + cacheReadInputTokens: 2, + }, }, - providerExecuted: true, - providerMetadata: { openai: { source: "provider" } }, - }), - LLMEvent.stepFinish({ - index: 0, - reason: { normalized: "tool-calls" }, - usage: { - inputTokens: 10, - nonCachedInputTokens: 8, - outputTokens: 4, - reasoningTokens: 1, - cacheReadInputTokens: 2, - }, - }), - LLMEvent.finish({ reason: { normalized: "tool-calls" } }), - ] + LLMEvent.reasoningStart({ id: "reasoning-1" }), + LLMEvent.reasoningDelta({ id: "reasoning-1", text: "Think" }), + LLMEvent.reasoningEnd({ id: "reasoning-1" }), + LLMEvent.toolInputStart({ id: "call-error", name: "write" }), + LLMEvent.toolInputDelta({ id: "call-error", name: "write", text: '{"path":"README.md"}' }), + LLMEvent.toolInputEnd({ id: "call-error", name: "write" }), + LLMEvent.toolCall({ id: "call-error", name: "write", input: { path: "README.md" }, providerExecuted: true }), + LLMEvent.toolError({ id: "call-error", name: "write", message: "Denied" }), + LLMEvent.toolResult({ id: "call-error", name: "write", result: { type: "error", value: "Denied" } }), + LLMEvent.toolCall({ + id: "call-provider", + name: "web_search", + input: { query: "hello" }, + providerExecuted: true, + providerMetadata: { openai: { source: "provider" } }, + }), + LLMEvent.toolResult({ + id: "call-provider", + name: "web_search", + result: { + type: "content", + value: [ + { type: "text", text: "Hello" }, + { type: "file", uri: "data:image/png;base64,aGVsbG8=", mime: "image/png", name: "hello.png" }, + ], + }, + providerExecuted: true, + providerMetadata: { openai: { source: "provider" } }, + }), + ), + ) yield* session.resume(sessionID) @@ -2398,12 +2302,12 @@ describe("SessionRunnerLLM", () => { const session = yield* setup yield* admit(session, "Echo this") - responses = [reply.tool("call-echo", "echo", { text: "hello" }), reply.text("Done", "text-final")] + yield* TestLLM.push(TestLLM.tool("call-echo", "echo", { text: "hello" }), TestLLM.text("Done", "text-final")) yield* session.resume(sessionID) expect(requests).toHaveLength(2) - expect(requests[1]?.messages.map((message) => message.role)).toEqual(["user", "assistant", "tool"]) + expect(messageRoles(requests[1])).toEqual(["user", "assistant", "tool"]) expect(authorizations).toMatchObject([{ sessionID, callID: "call-echo" }]) expect(executions).toEqual(["hello"]) const context = yield* session.context(sessionID) @@ -2428,7 +2332,7 @@ describe("SessionRunnerLLM", () => { { type: "assistant", finish: "stop", content: [{ type: "text", text: "Done" }] }, ]) const assistant = requireAssistant(context) - expect((yield* recordedStepSettlementEvents(sessionID, assistant.id)).map((event) => event.type)).toEqual([ + expect(yield* recordedStepSettlementTypes(sessionID, assistant.id)).toEqual([ "session.step.started.1", "session.tool.called.1", "session.tool.success.2", @@ -2443,18 +2347,16 @@ describe("SessionRunnerLLM", () => { const bus = yield* Bus.Service yield* admit(session, "Echo this") - responses = [reply.tool("call-echo", "echo", { text: "hello" }), reply.stop()] - toolExecutionGate = yield* Deferred.make() - toolExecutionsStarted = yield* Deferred.make() - toolExecutionsReady = 1 + yield* TestLLM.push(TestLLM.tool("call-echo", "echo", { text: "hello" }), TestLLM.stop()) + const tools = yield* blockTools() const run = yield* Effect.forkChild(session.resume(sessionID)) - yield* Deferred.await(toolExecutionsStarted) + yield* tools.started yield* bus.publish(SessionEvent.ModelSelected, { sessionID, model: { id: ID.make("replacement"), providerID: Provider.ID.make("fake") }, }) systemBaseline = "Replacement context" - yield* Deferred.succeed(toolExecutionGate, undefined) + yield* tools.release yield* Fiber.join(run) expect(requests.map((request) => request.model)).toEqual([model, replacementModel]) @@ -2462,7 +2364,7 @@ describe("SessionRunnerLLM", () => { [defaultSystem, "Initial context"], [defaultSystem, "Initial context"], ]) - expect(systemTexts(requests[1]!)).toContain("Replacement context") + expect(systemTexts(requests[1])).toContain("Replacement context") }), ) @@ -2471,32 +2373,31 @@ describe("SessionRunnerLLM", () => { const session = yield* setup yield* admit(session, "Think first") - response = [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.reasoningStart({ id: "reasoning-anthropic" }), - LLMEvent.reasoningDelta({ id: "reasoning-anthropic", text: "Signed thought" }), - LLMEvent.reasoningEnd({ - id: "reasoning-anthropic", - providerMetadata: { openai: { signature: "sig_1" }, anthropic: { ignored: true } }, - }), - LLMEvent.reasoningStart({ - id: "reasoning-openai", - providerMetadata: { - openai: { itemId: "rs_1", reasoningEncryptedContent: null }, - anthropic: { ignored: true }, - }, - }), - LLMEvent.reasoningDelta({ id: "reasoning-openai", text: "Encrypted thought" }), - LLMEvent.reasoningEnd({ - id: "reasoning-openai", - providerMetadata: { - openai: { itemId: "rs_1", reasoningEncryptedContent: "encrypted-state" }, - anthropic: { ignored: true }, - }, - }), - LLMEvent.stepFinish({ index: 0, reason: { normalized: "stop" } }), - LLMEvent.finish({ reason: { normalized: "stop" } }), - ] + yield* TestLLM.push( + TestLLM.stop( + LLMEvent.reasoningStart({ id: "reasoning-anthropic" }), + LLMEvent.reasoningDelta({ id: "reasoning-anthropic", text: "Signed thought" }), + LLMEvent.reasoningEnd({ + id: "reasoning-anthropic", + providerMetadata: { openai: { signature: "sig_1" }, anthropic: { ignored: true } }, + }), + LLMEvent.reasoningStart({ + id: "reasoning-openai", + providerMetadata: { + openai: { itemId: "rs_1", reasoningEncryptedContent: null }, + anthropic: { ignored: true }, + }, + }), + LLMEvent.reasoningDelta({ id: "reasoning-openai", text: "Encrypted thought" }), + LLMEvent.reasoningEnd({ + id: "reasoning-openai", + providerMetadata: { + openai: { itemId: "rs_1", reasoningEncryptedContent: "encrypted-state" }, + anthropic: { ignored: true }, + }, + }), + ), + ) yield* session.resume(sessionID) yield* replaySessionProjection(sessionID) @@ -2520,7 +2421,7 @@ describe("SessionRunnerLLM", () => { ]) yield* admit(session, "Continue") - response = [] + yield* TestLLM.push([]) yield* session.resume(sessionID) expect(requests[1]?.messages[1]?.content).toEqual([ @@ -2543,17 +2444,16 @@ describe("SessionRunnerLLM", () => { const session = yield* setup yield* admit(session, "Check first") - response = [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.textStart({ id: "commentary", providerMetadata: { openai: { phase: "commentary" } } }), - LLMEvent.textDelta({ id: "commentary", text: "Checking." }), - LLMEvent.textEnd({ - id: "commentary", - providerMetadata: { openai: { phase: "commentary" }, anthropic: { ignored: true } }, - }), - LLMEvent.stepFinish({ index: 0, reason: { normalized: "stop" } }), - LLMEvent.finish({ reason: { normalized: "stop" } }), - ] + yield* TestLLM.push( + TestLLM.stop( + LLMEvent.textStart({ id: "commentary", providerMetadata: { openai: { phase: "commentary" } } }), + LLMEvent.textDelta({ id: "commentary", text: "Checking." }), + LLMEvent.textEnd({ + id: "commentary", + providerMetadata: { openai: { phase: "commentary" }, anthropic: { ignored: true } }, + }), + ), + ) yield* session.resume(sessionID) yield* replaySessionProjection(sessionID) @@ -2566,7 +2466,7 @@ describe("SessionRunnerLLM", () => { ]) yield* admit(session, "Continue") - response = [] + yield* TestLLM.push([]) yield* session.resume(sessionID) expect(requests[1]?.messages[1]?.content).toEqual([ @@ -2584,33 +2484,32 @@ describe("SessionRunnerLLM", () => { const session = yield* setup yield* admit(session, "Search first") - response = [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.toolCall({ - id: "hosted-search", - name: "web_search", - input: { query: "Effect" }, - providerExecuted: true, - providerMetadata: { openai: { itemId: "hosted-search" }, fake: { ignored: true } }, - }), - LLMEvent.toolResult({ - id: "hosted-search", - name: "web_search", - result: { type: "json", value: [{ title: "Effect" }] }, - providerExecuted: true, - providerMetadata: { openai: { blockType: "web_search_tool_result" }, anthropic: { ignored: true } }, - }), - LLMEvent.stepFinish({ index: 0, reason: { normalized: "stop" } }), - LLMEvent.finish({ reason: { normalized: "stop" } }), - ] + yield* TestLLM.push( + TestLLM.stop( + LLMEvent.toolCall({ + id: "hosted-search", + name: "web_search", + input: { query: "Effect" }, + providerExecuted: true, + providerMetadata: { openai: { itemId: "hosted-search" }, fake: { ignored: true } }, + }), + LLMEvent.toolResult({ + id: "hosted-search", + name: "web_search", + result: { type: "json", value: [{ title: "Effect" }] }, + providerExecuted: true, + providerMetadata: { openai: { blockType: "web_search_tool_result" }, anthropic: { ignored: true } }, + }), + ), + ) yield* session.resume(sessionID) yield* replaySessionProjection(sessionID) yield* admit(session, "Continue") - response = [] + yield* TestLLM.push([]) yield* session.resume(sessionID) - expect(requests[1]?.messages.map((message) => message.role)).toEqual(["user", "assistant", "user"]) + expect(messageRoles(requests[1])).toEqual(["user", "assistant", "user"]) expect(requests[1]?.messages[1]?.content).toMatchObject([ { type: "tool-call", @@ -2638,8 +2537,7 @@ describe("SessionRunnerLLM", () => { const session = yield* setup yield* admit(session, "Echo five times") - toolExecutionGate = yield* Deferred.make() - toolExecutionsStarted = yield* Deferred.make() + const tools = yield* blockTools(5) const providerGate = yield* Deferred.make() const initial = Stream.fromIterable([ LLMEvent.stepStart({ index: 0 }), @@ -2651,16 +2549,15 @@ describe("SessionRunnerLLM", () => { LLMEvent.stepFinish({ index: 0, reason: { normalized: "tool-calls" } }), LLMEvent.finish({ reason: { normalized: "tool-calls" } }), ]) - responseStream = Stream.concat( - initial, - Stream.fromEffect(Deferred.await(providerGate)).pipe(Stream.flatMap(() => final)), + yield* TestLLM.push( + Stream.concat(initial, Stream.fromEffect(Deferred.await(providerGate)).pipe(Stream.flatMap(() => final))), ) const run = yield* session.resume(sessionID).pipe(Effect.forkChild) - yield* Deferred.await(toolExecutionsStarted) + yield* tools.started expect(executions).toHaveLength(5) - expect(maxActiveToolExecutions).toBe(5) + expect(yield* tools.maxActive).toBe(5) expect(yield* session.context(sessionID)).toMatchObject([ { type: "user", text: "Echo five times" }, { @@ -2677,13 +2574,11 @@ describe("SessionRunnerLLM", () => { yield* Effect.yieldNow expect(requests).toHaveLength(1) - yield* Deferred.succeed(toolExecutionGate, undefined) + yield* tools.release yield* Fiber.join(run) - toolExecutionGate = undefined - toolExecutionsStarted = undefined expect(executions).toHaveLength(5) - expect(maxActiveToolExecutions).toBe(5) + expect(yield* tools.maxActive).toBe(5) expect(requests).toHaveLength(2) }), ) @@ -2693,17 +2588,15 @@ describe("SessionRunnerLLM", () => { const session = yield* setup yield* admit(session, "Echo twice") - responses = [ - reply.tool("tool_0", "echo", { text: "first" }), - reply.tool("tool_0", "echo", { text: "second" }), + yield* TestLLM.push( + TestLLM.tool("tool_0", "echo", { text: "first" }), + TestLLM.tool("tool_0", "echo", { text: "second" }), [], - ] + ) yield* session.resume(sessionID) - expect(executions).toEqual(["first", "second"]) - expect(requests).toHaveLength(3) - expect(yield* session.context(sessionID)).toMatchObject([ + const expected = [ { type: "user", text: "Echo twice" }, { type: "assistant", @@ -2725,33 +2618,14 @@ describe("SessionRunnerLLM", () => { }, ], }, - ]) + ] + expect(executions).toEqual(["first", "second"]) + expect(requests).toHaveLength(3) + expect(yield* session.context(sessionID)).toMatchObject(expected) yield* replaySessionProjection(sessionID) - expect(yield* session.context(sessionID)).toMatchObject([ - { type: "user", text: "Echo twice" }, - { - type: "assistant", - content: [ - { - type: "tool", - id: "tool_0", - state: { status: "completed", content: [{ type: "text", text: "first" }] }, - }, - ], - }, - { - type: "assistant", - content: [ - { - type: "tool", - id: "tool_0", - state: { status: "completed", content: [{ type: "text", text: "second" }] }, - }, - ], - }, - ]) + expect(yield* session.context(sessionID)).toMatchObject(expected) }), ) @@ -2760,21 +2634,18 @@ describe("SessionRunnerLLM", () => { const session = yield* setup yield* admit(session, "Run once") - response = reply.text("Once", "text-once") - streamGate = yield* Deferred.make() - streamStarted = yield* Deferred.make() + yield* TestLLM.push(TestLLM.text("Once", "text-once")) + const stream = yield* TestLLM.gate const first = yield* session.resume(sessionID).pipe(Effect.forkChild) - yield* Deferred.await(streamStarted) + yield* stream.started const second = yield* session.resume(sessionID).pipe(Effect.forkChild) yield* Effect.yieldNow expect(requests).toHaveLength(1) - yield* Deferred.succeed(streamGate, undefined) + yield* stream.release yield* Fiber.join(first) yield* Fiber.join(second) - streamGate = undefined - streamStarted = undefined expect(requests).toHaveLength(1) expect(yield* session.context(sessionID)).toMatchObject([ @@ -2789,22 +2660,19 @@ describe("SessionRunnerLLM", () => { const session = yield* setup yield* admit(session, "Start working") - responses = [reply.stop(), reply.stop()] - streamGate = yield* Deferred.make() - streamStarted = yield* Deferred.make() + yield* TestLLM.push(TestLLM.stop(), TestLLM.stop()) + const stream = yield* TestLLM.gate const first = yield* session.resume(sessionID).pipe(Effect.forkChild) - yield* Deferred.await(streamStarted) + yield* stream.started yield* session.prompt({ sessionID, text: "Change direction" }) - yield* Deferred.succeed(streamGate, undefined) + yield* stream.release yield* Fiber.join(first) - streamGate = undefined - streamStarted = undefined yield* Effect.yieldNow expect(requests).toHaveLength(2) - expect(userTexts(requests[0]!)).toEqual(["Start working"]) - expect(userTexts(requests[1]!)).toEqual(["Start working", "Change direction"]) + expect(userTexts(requests[0])).toEqual(["Start working"]) + expect(userTexts(requests[1])).toEqual(["Start working", "Change direction"]) expect((yield* session.context(sessionID)).map((message) => message.type)).toEqual([ "user", "assistant", @@ -2819,26 +2687,23 @@ describe("SessionRunnerLLM", () => { const session = yield* setup yield* admit(session, "Start working") - responses = [reply.tool("call-echo", "echo", { text: "hello" }), reply.stop(), reply.stop()] - streamGate = yield* Deferred.make() - streamStarted = yield* Deferred.make() + yield* TestLLM.push(TestLLM.tool("call-echo", "echo", { text: "hello" }), TestLLM.stop(), TestLLM.stop()) + const stream = yield* TestLLM.gate const first = yield* session.resume(sessionID).pipe(Effect.forkChild) - yield* Deferred.await(streamStarted) + yield* stream.started yield* session.prompt({ sessionID, text: "Wait until continuation ends", delivery: "queue", }) - yield* Deferred.succeed(streamGate, undefined) + yield* stream.release yield* Fiber.join(first) - streamGate = undefined - streamStarted = undefined expect(requests).toHaveLength(3) - expect(userTexts(requests[0]!)).toEqual(["Start working"]) - expect(userTexts(requests[1]!)).toEqual(["Start working"]) - expect(userTexts(requests[2]!)).toEqual(["Start working", "Wait until continuation ends"]) + expect(userTexts(requests[0])).toEqual(["Start working"]) + expect(userTexts(requests[1])).toEqual(["Start working"]) + expect(userTexts(requests[2])).toEqual(["Start working", "Wait until continuation ends"]) }), ) @@ -2848,12 +2713,11 @@ describe("SessionRunnerLLM", () => { const { db } = yield* Database.Service yield* admit(session, "Interrupt current work") - responses = [[], reply.stop()] - streamGate = yield* Deferred.make() - streamStarted = yield* Deferred.make() + yield* TestLLM.push([], TestLLM.stop()) + const stream = yield* TestLLM.gate const run = yield* session.resume(sessionID).pipe(Effect.forkChild) - yield* Deferred.await(streamStarted) + yield* stream.started yield* session.prompt({ sessionID, text: "Run after interrupt", @@ -2864,15 +2728,13 @@ describe("SessionRunnerLLM", () => { expect(requests).toHaveLength(1) expect(yield* SessionPending.has(db, sessionID, "queue")).toBe(true) const resumed = yield* session.resume(sessionID).pipe(Effect.forkChild) - while (requests.length < 2) yield* Effect.yieldNow - yield* Deferred.succeed(streamGate, undefined) + yield* stream.started + yield* stream.release yield* Fiber.join(resumed) - streamGate = undefined - streamStarted = undefined expect(requests).toHaveLength(2) - expect(userTexts(requests[0]!)).toEqual(["Interrupt current work"]) - expect(userTexts(requests[1]!)).toEqual(["Interrupt current work", "Run after interrupt"]) + expect(userTexts(requests[0])).toEqual(["Interrupt current work"]) + expect(userTexts(requests[1])).toEqual(["Interrupt current work", "Run after interrupt"]) }), ) @@ -2882,12 +2744,11 @@ describe("SessionRunnerLLM", () => { const { db } = yield* Database.Service yield* admit(session, "Interrupt current work") - responses = [[], reply.stop()] - streamGate = yield* Deferred.make() - streamStarted = yield* Deferred.make() + yield* TestLLM.push([], TestLLM.stop()) + const stream = yield* TestLLM.gate const run = yield* session.resume(sessionID).pipe(Effect.forkChild) - yield* Deferred.await(streamStarted) + yield* stream.started yield* session.prompt({ sessionID, text: "Steer after interrupt", @@ -2898,15 +2759,13 @@ describe("SessionRunnerLLM", () => { expect(yield* SessionPending.has(db, sessionID, "steer")).toBe(true) const resumed = yield* session.resume(sessionID).pipe(Effect.forkChild) - while (requests.length < 2) yield* Effect.yieldNow - yield* Deferred.succeed(streamGate, undefined) + yield* stream.started + yield* stream.release yield* Fiber.join(resumed) - streamGate = undefined - streamStarted = undefined expect(requests).toHaveLength(2) - expect(userTexts(requests[0]!)).toEqual(["Interrupt current work"]) - expect(userTexts(requests[1]!)).toEqual(["Interrupt current work", "Steer after interrupt"]) + expect(userTexts(requests[0])).toEqual(["Interrupt current work"]) + expect(userTexts(requests[1])).toEqual(["Interrupt current work", "Steer after interrupt"]) }), ) @@ -2915,23 +2774,20 @@ describe("SessionRunnerLLM", () => { const session = yield* setup yield* admit(session, "Start working") - responses = [reply.stop(), reply.stop(), reply.stop()] - streamGate = yield* Deferred.make() - streamStarted = yield* Deferred.make() + yield* TestLLM.push(TestLLM.stop(), TestLLM.stop(), TestLLM.stop()) + const stream = yield* TestLLM.gate const first = yield* session.resume(sessionID).pipe(Effect.forkChild) - yield* Deferred.await(streamStarted) + yield* stream.started yield* session.prompt({ sessionID, text: "Queue first", delivery: "queue" }) yield* session.prompt({ sessionID, text: "Queue second", delivery: "queue" }) - yield* Deferred.succeed(streamGate, undefined) + yield* stream.release yield* Fiber.join(first) - streamGate = undefined - streamStarted = undefined expect(requests).toHaveLength(3) - expect(userTexts(requests[0]!)).toEqual(["Start working"]) - expect(userTexts(requests[1]!)).toEqual(["Start working", "Queue first"]) - expect(userTexts(requests[2]!)).toEqual(["Start working", "Queue first", "Queue second"]) + expect(userTexts(requests[0])).toEqual(["Start working"]) + expect(userTexts(requests[1])).toEqual(["Start working", "Queue first"]) + expect(userTexts(requests[2])).toEqual(["Start working", "Queue first", "Queue second"]) }), ) @@ -2946,13 +2802,13 @@ describe("SessionRunnerLLM", () => { resume: false, }) - responses = [reply.stop(), reply.stop()] + yield* TestLLM.push(TestLLM.stop(), TestLLM.stop()) yield* session.resume(sessionID) expect(requests).toHaveLength(2) - expect(userTexts(requests[0]!)).toEqual(["Start steering"]) - expect(userTexts(requests[1]!)).toEqual(["Start steering", "Queue for later"]) + expect(userTexts(requests[0])).toEqual(["Start steering"]) + expect(userTexts(requests[1])).toEqual(["Start steering", "Queue for later"]) }), ) @@ -2961,39 +2817,36 @@ describe("SessionRunnerLLM", () => { const session = yield* setup yield* admit(session, "Start working") - responses = [reply.stop(), reply.stop(), reply.stop(), reply.stop()] - const firstGate = yield* Deferred.make() - const secondGate = yield* Deferred.make() - streamGate = firstGate + yield* TestLLM.push(TestLLM.stop(), TestLLM.stop(), TestLLM.stop(), TestLLM.stop()) + const firstStream = yield* TestLLM.gate const first = yield* session.resume(sessionID).pipe(Effect.forkChild) - while (requests.length < 1) yield* Effect.yieldNow + yield* firstStream.started yield* session.prompt({ sessionID, text: "Queue first", delivery: "queue" }) yield* session.prompt({ sessionID, text: "Queue second", delivery: "queue" }) - streamGate = secondGate - yield* Deferred.succeed(firstGate, undefined) - while (requests.length < 2) yield* Effect.yieldNow + const secondStream = yield* TestLLM.gate + yield* firstStream.release + yield* secondStream.started yield* session.prompt({ sessionID, text: "Steer before next queued input" }) yield* session.prompt({ sessionID, text: "Also steer before next queued input", }) yield* session.synthetic({ sessionID, text: "Background completion before next queued input" }) - yield* Deferred.succeed(secondGate, undefined) + yield* secondStream.release yield* Fiber.join(first) - streamGate = undefined expect(requests).toHaveLength(4) - expect(userTexts(requests[0]!)).toEqual(["Start working"]) - expect(userTexts(requests[1]!)).toEqual(["Start working", "Queue first"]) - expect(userTexts(requests[2]!)).toEqual([ + expect(userTexts(requests[0])).toEqual(["Start working"]) + expect(userTexts(requests[1])).toEqual(["Start working", "Queue first"]) + expect(userTexts(requests[2])).toEqual([ "Start working", "Queue first", "Steer before next queued input", "Also steer before next queued input", "Background completion before next queued input", ]) - expect(userTexts(requests[3]!)).toEqual([ + expect(userTexts(requests[3])).toEqual([ "Start working", "Queue first", "Steer before next queued input", @@ -3009,22 +2862,19 @@ describe("SessionRunnerLLM", () => { const session = yield* setup yield* admit(session, "Start working") - responses = [reply.stop(), reply.stop()] - streamGate = yield* Deferred.make() - streamStarted = yield* Deferred.make() + yield* TestLLM.push(TestLLM.stop(), TestLLM.stop()) + const stream = yield* TestLLM.gate const first = yield* session.resume(sessionID).pipe(Effect.forkChild) - yield* Deferred.await(streamStarted) + yield* stream.started yield* session.prompt({ sessionID, text: "First steer" }) yield* session.prompt({ sessionID, text: "Second steer" }) - yield* Deferred.succeed(streamGate, undefined) + yield* stream.release yield* Fiber.join(first) - streamGate = undefined - streamStarted = undefined yield* Effect.yieldNow expect(requests).toHaveLength(2) - expect(userTexts(requests[1]!)).toEqual(["Start working", "First steer", "Second steer"]) + expect(userTexts(requests[1])).toEqual(["Start working", "First steer", "Second steer"]) yield* (yield* SessionExecution.Service).wake(sessionID) yield* Effect.yieldNow expect(requests).toHaveLength(2) @@ -3036,23 +2886,21 @@ describe("SessionRunnerLLM", () => { const session = yield* setup yield* admit(session, "Start working") - streamFailure = invalidRequest() - streamGate = yield* Deferred.make() - streamStarted = yield* Deferred.make() + const failure = invalidRequest() + yield* TestLLM.push(Stream.fail(failure)) + const stream = yield* TestLLM.gate const first = yield* session.resume(sessionID).pipe(Effect.forkChild) - yield* Deferred.await(streamStarted) + yield* stream.started yield* session.prompt({ sessionID, text: "Recover with this" }) - yield* Deferred.succeed(streamGate, undefined) - expect(yield* Fiber.join(first).pipe(Effect.flip)).toBe(streamFailure) + yield* stream.release + expect(yield* Fiber.join(first).pipe(Effect.flip)).toBe(failure) - streamFailure = undefined - streamGate = undefined - streamStarted = undefined + yield* TestLLM.push([]) yield* session.wait(sessionID) expect(requests).toHaveLength(2) - expect(userTexts(requests[1]!)).toEqual(["Start working", "Recover with this"]) + expect(userTexts(requests[1])).toEqual(["Start working", "Recover with this"]) }), ) @@ -3089,11 +2937,11 @@ describe("SessionRunnerLLM", () => { executed: false, }) requests.length = 0 - response = [] + yield* TestLLM.push([]) yield* session.resume(sessionID) expect(requests).toHaveLength(1) - expect(requests[0]?.messages.map((message) => message.role)).toEqual(["user", "assistant", "tool"]) + expect(messageRoles(requests[0])).toEqual(["user", "assistant", "tool"]) expect(yield* session.context(sessionID)).toMatchObject([ { type: "user", text: "Recover interrupted tool" }, { @@ -3147,11 +2995,11 @@ describe("SessionRunnerLLM", () => { state: { itemId: "call-hosted-interrupted" }, }) requests.length = 0 - response = [] + yield* TestLLM.push([]) yield* session.resume(sessionID) expect(requests).toHaveLength(1) - expect(requests[0]?.messages.map((message) => message.role)).toEqual(["user", "assistant"]) + expect(messageRoles(requests[0])).toEqual(["user", "assistant"]) expect(requests[0]?.messages[1]?.content).toMatchObject([ { type: "tool-call", @@ -3184,11 +3032,11 @@ describe("SessionRunnerLLM", () => { name: "echo", }) requests.length = 0 - response = [] + yield* TestLLM.push([]) yield* session.resume(sessionID) expect(requests).toHaveLength(1) - expect(requests[0]?.messages.map((message) => message.role)).toEqual(["user", "assistant", "tool"]) + expect(messageRoles(requests[0])).toEqual(["user", "assistant", "tool"]) expect(yield* session.context(sessionID)).toMatchObject([ { type: "user", text: "Recover interrupted tool input" }, { type: "assistant", content: [{ type: "tool", id: "call-pending-interrupted", state: { status: "error" } }] }, @@ -3206,11 +3054,13 @@ describe("SessionRunnerLLM", () => { resume: false, }) + const stream = yield* TestLLM.gate yield* (yield* SessionExecution.Service).wake(sessionID) - while (requests.length === 0) yield* Effect.yieldNow + yield* stream.started + yield* stream.release expect(requests).toHaveLength(1) - expect(userTexts(requests[0]!)).toEqual(["Wait in queue"]) + expect(userTexts(requests[0])).toEqual(["Wait in queue"]) }), ) @@ -3226,12 +3076,14 @@ describe("SessionRunnerLLM", () => { expect(yield* session.resume(sessionID).pipe(Effect.catchDefect(Effect.succeed))).toBe(defect) fail = false requests.length = 0 - response = reply.stop() + yield* TestLLM.push(TestLLM.stop()) + const stream = yield* TestLLM.gate yield* (yield* SessionExecution.Service).wake(sessionID) - while (requests.length === 0) yield* Effect.yieldNow + yield* stream.started + yield* stream.release - expect(userTexts(requests[0]!)).toEqual(["Recover promoted input"]) + expect(userTexts(requests[0])).toEqual(["Recover promoted input"]) }), ) @@ -3244,21 +3096,17 @@ describe("SessionRunnerLLM", () => { ? Effect.die("fail after prompt promotion commits") : Effect.void, ) - yield* admit(session, "Run committed promotion") - - yield* session.resume(sessionID) + yield* runPrompt(session, "Run committed promotion") expect(requests).toHaveLength(1) - expect(userTexts(requests[0]!)).toEqual(["Run committed promotion"]) + expect(userTexts(requests[0])).toEqual(["Run committed promotion"]) }), ) it.effect("adds session correlation headers to model requests", () => Effect.gen(function* () { const session = yield* setup - yield* admit(session, "Run correlated request") - - yield* session.resume(sessionID) + yield* runPrompt(session, "Run correlated request") expect(requests[0]?.http?.headers).toEqual({ "x-session-affinity": sessionID, @@ -3282,9 +3130,7 @@ describe("SessionRunnerLLM", () => { .where(eq(SessionTable.id, sessionID)) .run() .pipe(Effect.orDie) - yield* admit(session, "Run child request") - - yield* session.resume(sessionID) + yield* runPrompt(session, "Run child request") expect(requests[0]?.http?.headers?.["x-parent-session-id"]).toBe(parentID) }), @@ -3301,25 +3147,21 @@ describe("SessionRunnerLLM", () => { resume: false, }) - streamGate = yield* Deferred.make() - streamStarted = yield* Deferred.make() + const stream = yield* TestLLM.gate const first = yield* session.resume(sessionID).pipe(Effect.forkChild) - yield* Deferred.await(streamStarted) - streamStarted = yield* Deferred.make() + yield* stream.started const second = yield* session.resume(otherSessionID).pipe(Effect.forkChild) - yield* Deferred.await(streamStarted) + yield* stream.started expect(requests).toHaveLength(2) expect(requests.map((request) => request.providerOptions?.openai?.promptCacheKey)).toEqual([ sessionID, otherSessionID, ]) - yield* Deferred.succeed(streamGate, undefined) + yield* stream.release yield* Fiber.join(first) yield* Fiber.join(second) - streamGate = undefined - streamStarted = undefined }), ) @@ -3356,23 +3198,20 @@ describe("SessionRunnerLLM", () => { const session = yield* setup yield* admit(session, "Retry after failure") - streamFailure = invalidRequest() - streamGate = yield* Deferred.make() - streamStarted = yield* Deferred.make() + yield* TestLLM.push(Stream.fail(invalidRequest())) + const stream = yield* TestLLM.gate const first = yield* session.resume(sessionID).pipe(Effect.forkChild) - yield* Deferred.await(streamStarted) + yield* stream.started const second = yield* session.resume(sessionID).pipe(Effect.forkChild) yield* Effect.yieldNow expect(requests).toHaveLength(1) - yield* Deferred.succeed(streamGate, undefined) + yield* stream.release const [firstExit, secondExit] = yield* Effect.all([Fiber.await(first), Fiber.await(second)]) expect(secondExit).toEqual(firstExit) - streamFailure = undefined - streamGate = undefined - streamStarted = undefined + yield* TestLLM.push([]) yield* session.resume(sessionID) expect(requests).toHaveLength(2) }), @@ -3383,7 +3222,7 @@ describe("SessionRunnerLLM", () => { const session = yield* setup yield* admit(session, "Call missing") - responses = [reply.tool("call-missing", "missing", {}), reply.text("Recovered", "text-after-error")] + yield* TestLLM.push(TestLLM.tool("call-missing", "missing", {}), TestLLM.text("Recovered", "text-after-error")) yield* session.resume(sessionID) expect(requests).toHaveLength(2) @@ -3412,12 +3251,12 @@ describe("SessionRunnerLLM", () => { const session = yield* setup yield* admit(session, "Call defect") - responses = [reply.tool("call-defect", "defect", {}), reply.text("Recovered", "text-after-defect")] + yield* TestLLM.push(TestLLM.tool("call-defect", "defect", {}), TestLLM.text("Recovered", "text-after-defect")) yield* session.resume(sessionID) expect(requests).toHaveLength(2) - expect(requests[1]?.messages.map((message) => message.role)).toEqual(["user", "assistant", "tool"]) + expect(messageRoles(requests[1])).toEqual(["user", "assistant", "tool"]) const context = yield* session.context(sessionID) expect(context).toMatchObject([ { type: "user", text: "Call defect" }, @@ -3437,7 +3276,7 @@ describe("SessionRunnerLLM", () => { { type: "assistant", finish: "stop", content: [{ type: "text", text: "Recovered" }] }, ]) const assistant = requireAssistant(context) - expect((yield* recordedStepSettlementEvents(sessionID, assistant.id)).map((event) => event.type)).toEqual([ + expect(yield* recordedStepSettlementTypes(sessionID, assistant.id)).toEqual([ "session.step.started.1", "session.tool.called.1", "session.tool.failed.2", @@ -3450,9 +3289,10 @@ describe("SessionRunnerLLM", () => { Effect.gen(function* () { const session = yield* setup const registry = yield* Tool.Service - yield* transformTools(registry, + yield* transformTools( + registry, { - blocked: ({ + blocked: { name: "blocked", description: "Fail because policy blocked execution", input: Schema.Struct({}), @@ -3461,13 +3301,13 @@ describe("SessionRunnerLLM", () => { Effect.fail(new Permission.BlockedError({ rules: [], permission: "blocked", resources: ["*"] })).pipe( Effect.mapError(() => new Tool.Error({ message: "Permission blocked" })), ), - }), + }, }, { codemode: false }, ) yield* admit(session, "Call blocked") - responses = [reply.tool("call-blocked", "blocked", {}), reply.stop()] + yield* TestLLM.push(TestLLM.tool("call-blocked", "blocked", {}), TestLLM.stop()) yield* session.resume(sessionID) @@ -3489,21 +3329,22 @@ describe("SessionRunnerLLM", () => { Effect.gen(function* () { const session = yield* setup const registry = yield* Tool.Service - yield* transformTools(registry, + yield* transformTools( + registry, { - declined: ({ + declined: { name: "declined", description: "Fail because the user declined approval", input: Schema.Struct({}), output: Schema.Struct({}), execute: () => Effect.die(new Permission.DeclinedError()), - }), + }, }, { codemode: false }, ) yield* admit(session, "Call declined") - response = reply.tool("call-declined", "declined", {}) + yield* TestLLM.push(TestLLM.tool("call-declined", "declined", {})) const exit = yield* session.resume(sessionID).pipe(Effect.exit) @@ -3530,9 +3371,10 @@ describe("SessionRunnerLLM", () => { Effect.gen(function* () { const session = yield* setup const registry = yield* Tool.Service - yield* transformTools(registry, + yield* transformTools( + registry, { - corrected: ({ + corrected: { name: "corrected", description: "Fail with user correction feedback", input: Schema.Struct({}), @@ -3541,13 +3383,13 @@ describe("SessionRunnerLLM", () => { Effect.fail(new Permission.CorrectedError({ feedback: "Use another tool" })).pipe( Effect.mapError(() => new Tool.Error({ message: "Use another tool" })), ), - }), + }, }, { codemode: false }, ) yield* admit(session, "Call corrected") - responses = [reply.tool("call-corrected", "corrected", {}), reply.stop()] + yield* TestLLM.push(TestLLM.tool("call-corrected", "corrected", {}), TestLLM.stop()) yield* session.resume(sessionID) @@ -3571,10 +3413,10 @@ describe("SessionRunnerLLM", () => { const registry = yield* Tool.Service yield* transformTools(registry, { permissionfail: permissionFail }, { codemode: false }) yield* admit(session, "Reject permission") - responses = [ - reply.tool("call-permission", "permissionfail", {}), - [LLMEvent.stepStart({ index: 0 }), LLMEvent.stepFinish({ index: 0, reason: { normalized: "stop" } })], - ] + yield* TestLLM.push(TestLLM.tool("call-permission", "permissionfail", {}), [ + LLMEvent.stepStart({ index: 0 }), + LLMEvent.stepFinish({ index: 0, reason: { normalized: "stop" } }), + ]) yield* session.resume(sessionID) @@ -3607,21 +3449,22 @@ describe("SessionRunnerLLM", () => { Effect.gen(function* () { const session = yield* setup const registry = yield* Tool.Service - yield* transformTools(registry, + yield* transformTools( + registry, { - question: ({ + question: { name: "question", description: "Ask the user", input: Schema.Struct({}), output: Schema.Struct({}), execute: () => Effect.die(new QuestionTool.CancelledError()), - }), + }, }, { codemode: false }, ) yield* admit(session, "Ask then stop") - responses = [reply.tool("call-question", "question", {}), []] + yield* TestLLM.push(TestLLM.tool("call-question", "question", {}), []) const run = yield* session.resume(sessionID).pipe(Effect.exit, Effect.forkChild) const exit = yield* Fiber.join(run) @@ -3650,21 +3493,19 @@ describe("SessionRunnerLLM", () => { const session = yield* setup yield* admit(session, "Settle before failing") const failure = providerUnavailable() - toolExecutionGate = yield* Deferred.make() - responseStream = Stream.concat( - Stream.fromIterable([ + const tools = yield* blockTools() + yield* TestLLM.push( + TestLLM.failAfter( + failure, LLMEvent.stepStart({ index: 0 }), LLMEvent.toolCall({ id: "call-before-failure", name: "echo", input: { text: "settle" } }), - ]), - Stream.fail(failure), + ), ) const run = yield* session.resume(sessionID).pipe(Effect.forkChild) - while (executions.length === 0) yield* Effect.yieldNow - yield* Effect.yieldNow - yield* Deferred.succeed(toolExecutionGate, undefined) + yield* tools.started + yield* tools.release expect(yield* Fiber.join(run).pipe(Effect.flip)).toBe(failure) - toolExecutionGate = undefined const context = yield* session.context(sessionID) expect(context).toMatchObject([ @@ -3681,7 +3522,7 @@ describe("SessionRunnerLLM", () => { }, ]) const assistant = requireAssistant(context) - expect((yield* recordedStepSettlementEvents(sessionID, assistant.id)).map((event) => event.type)).toEqual([ + expect(yield* recordedStepSettlementTypes(sessionID, assistant.id)).toEqual([ "session.step.started.1", "session.tool.called.1", "session.tool.success.2", @@ -3694,19 +3535,17 @@ describe("SessionRunnerLLM", () => { Effect.gen(function* () { const session = yield* setup yield* admit(session, "Interrupt blocked tool") - toolExecutionGate = yield* Deferred.make() - responseStream = Stream.concat( - Stream.fromIterable([ + const tools = yield* blockTools() + yield* TestLLM.push( + TestLLM.hangAfter( LLMEvent.stepStart({ index: 0 }), LLMEvent.toolCall({ id: "call-before-interrupt", name: "echo", input: { text: "blocked" } }), - ]), - Stream.never, + ), ) const run = yield* session.resume(sessionID).pipe(Effect.forkChild) - while (executions.length === 0) yield* Effect.yieldNow + yield* tools.started yield* session.interrupt(sessionID) - toolExecutionGate = undefined expect(yield* Fiber.await(run)).toMatchObject({ _tag: "Failure" }) yield* session.interrupt(sessionID) @@ -3725,7 +3564,7 @@ describe("SessionRunnerLLM", () => { }, ]) const assistant = requireAssistant(context) - expect((yield* recordedStepSettlementEvents(sessionID, assistant.id)).map((event) => event.type)).toEqual([ + expect(yield* recordedStepSettlementTypes(sessionID, assistant.id)).toEqual([ "session.step.started.1", "session.tool.called.1", "session.tool.failed.2", @@ -3739,10 +3578,9 @@ describe("SessionRunnerLLM", () => { { type: "assistant", content: [{ type: "tool", id: "call-before-interrupt", state: { status: "error" } }] }, ]) requests.length = 0 - responseStream = undefined - response = [] + yield* TestLLM.push([]) yield* session.resume(sessionID) - expect(requests[0]?.messages.map((message) => message.role)).toEqual(["user", "assistant", "tool"]) + expect(messageRoles(requests[0])).toEqual(["user", "assistant", "tool"]) }), ) @@ -3750,15 +3588,12 @@ describe("SessionRunnerLLM", () => { Effect.gen(function* () { const session = yield* setup yield* admit(session, "Interrupt provider") - streamGate = yield* Deferred.make() - streamStarted = yield* Deferred.make() + const stream = yield* TestLLM.gate const run = yield* session.resume(sessionID).pipe(Effect.forkChild) - yield* Deferred.await(streamStarted) + yield* stream.started yield* session.interrupt(sessionID) const exit = yield* Fiber.await(run) - streamGate = undefined - streamStarted = undefined expect(Exit.isFailure(exit) && Cause.hasInterruptsOnly(exit.cause)).toBeTrue() expect(requests).toHaveLength(1) @@ -3775,16 +3610,13 @@ describe("SessionRunnerLLM", () => { Effect.gen(function* () { const session = yield* setup yield* admit(session, "Interrupt tool settlement") - toolExecutionGate = yield* Deferred.make() - toolExecutionsStarted = yield* Deferred.make() - toolExecutionsReady = 1 - response = reply.tool("call-await-interrupt", "echo", { text: "blocked" }) + const tools = yield* blockTools() + yield* TestLLM.push(TestLLM.tool("call-await-interrupt", "echo", { text: "blocked" })) const runner = yield* SessionRunner.Service const run = yield* runner.drain({ sessionID, force: true }).pipe(Effect.forkChild) - yield* Deferred.await(toolExecutionsStarted) + yield* tools.started yield* Fiber.interrupt(run) - toolExecutionGate = undefined expect(yield* Fiber.await(run)).toMatchObject({ _tag: "Failure" }) expect(yield* session.context(sessionID)).toMatchObject([ @@ -3819,10 +3651,10 @@ describe("SessionRunnerLLM", () => { ) yield* admit(session, "Finish at the limit") - responses = [ - reply.tool("call-terminal", "echo", { text: "done" }), - reply.tool("call-forbidden", "echo", { text: "forbidden" }), - ] + yield* TestLLM.push( + TestLLM.tool("call-terminal", "echo", { text: "done" }), + TestLLM.tool("call-forbidden", "echo", { text: "forbidden" }), + ) yield* session.resume(sessionID) @@ -3855,21 +3687,18 @@ describe("SessionRunnerLLM", () => { ) yield* admit(session, "Start work") - responses = [ - reply.tool("call-before-steer", "echo", { text: "before" }), - reply.tool("call-after-steer", "echo", { text: "after" }), - reply.stop(), - ] - streamGate = yield* Deferred.make() - streamStarted = yield* Deferred.make() + yield* TestLLM.push( + TestLLM.tool("call-before-steer", "echo", { text: "before" }), + TestLLM.tool("call-after-steer", "echo", { text: "after" }), + TestLLM.stop(), + ) + const stream = yield* TestLLM.gate const run = yield* session.resume(sessionID).pipe(Effect.forkChild) - yield* Deferred.await(streamStarted) + yield* stream.started yield* session.prompt({ sessionID, text: "Change direction" }) - yield* Deferred.succeed(streamGate, undefined) + yield* stream.release yield* Fiber.join(run) - streamGate = undefined - streamStarted = undefined expect(requests).toHaveLength(3) expect(requests[1]?.toolChoice).toBeUndefined() @@ -3882,11 +3711,12 @@ describe("SessionRunnerLLM", () => { it.effect("projects provider errors as terminal assistant step failures", () => Effect.gen(function* () { const session = yield* setup - yield* admit(session, "Fail durably") + yield* TestLLM.push([ + LLMEvent.stepStart({ index: 0 }), + LLMEvent.providerError({ message: "Provider unavailable" }), + ]) - response = [LLMEvent.stepStart({ index: 0 }), LLMEvent.providerError({ message: "Provider unavailable" })] - - expect((yield* session.resume(sessionID).pipe(Effect.flip)).message).toBe("Provider unavailable") + expect((yield* runPrompt(session, "Fail durably").pipe(Effect.flip)).message).toBe("Provider unavailable") expect(requests).toHaveLength(1) expect(yield* session.context(sessionID)).toMatchObject([ @@ -3899,11 +3729,9 @@ describe("SessionRunnerLLM", () => { it.effect("projects provider errors emitted before assistant step start", () => Effect.gen(function* () { const session = yield* setup - yield* admit(session, "Fail before step") + yield* TestLLM.push([LLMEvent.providerError({ message: "Provider unavailable" })]) - response = [LLMEvent.providerError({ message: "Provider unavailable" })] - - expect((yield* session.resume(sessionID).pipe(Effect.flip)).message).toBe("Provider unavailable") + expect((yield* runPrompt(session, "Fail before step").pipe(Effect.flip)).message).toBe("Provider unavailable") expect(requests).toHaveLength(1) expect(yield* session.context(sessionID)).toMatchObject([ @@ -3916,20 +3744,20 @@ describe("SessionRunnerLLM", () => { it.effect("projects content-filter finishes as visible terminal failures", () => Effect.gen(function* () { const session = yield* setup - yield* admit(session, "Blocked response") - response = [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.textStart({ id: "partial" }), - LLMEvent.textDelta({ id: "partial", text: "Partial" }), - LLMEvent.stepFinish({ - index: 0, - reason: { normalized: "content-filter" }, - usage: { nonCachedInputTokens: 8, outputTokens: 3, reasoningTokens: 1 }, - }), - LLMEvent.finish({ reason: { normalized: "content-filter" } }), - ] + yield* TestLLM.push( + TestLLM.complete( + { + reason: { normalized: "content-filter" }, + usage: { nonCachedInputTokens: 8, outputTokens: 3, reasoningTokens: 1 }, + }, + LLMEvent.textStart({ id: "partial" }), + LLMEvent.textDelta({ id: "partial", text: "Partial" }), + ), + ) - expect((yield* session.resume(sessionID).pipe(Effect.flip)).message).toBe("Provider blocked the response") + expect((yield* runPrompt(session, "Blocked response").pipe(Effect.flip)).message).toBe( + "Provider blocked the response", + ) expect(yield* session.context(sessionID)).toMatchObject([ { type: "user" }, { @@ -3953,22 +3781,18 @@ describe("SessionRunnerLLM", () => { Effect.gen(function* () { const session = yield* setup yield* admit(session, "Tool before blocked response") - toolExecutionGate = yield* Deferred.make() - toolExecutionsStarted = yield* Deferred.make() - toolExecutionsReady = 1 - response = [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.toolCall({ id: "call-before-content-filter", name: "echo", input: { text: "settled" } }), - LLMEvent.stepFinish({ index: 0, reason: { normalized: "content-filter" } }), - LLMEvent.finish({ reason: { normalized: "content-filter" } }), - ] + const tools = yield* blockTools() + yield* TestLLM.push( + TestLLM.complete( + { reason: { normalized: "content-filter" } }, + LLMEvent.toolCall({ id: "call-before-content-filter", name: "echo", input: { text: "settled" } }), + ), + ) const run = yield* session.resume(sessionID).pipe(Effect.forkChild) - yield* Deferred.await(toolExecutionsStarted) - yield* Deferred.succeed(toolExecutionGate, undefined) + yield* tools.started + yield* tools.release expect((yield* Fiber.join(run).pipe(Effect.flip)).message).toBe("Provider blocked the response") - toolExecutionGate = undefined - toolExecutionsStarted = undefined const assistant = requireAssistant(yield* session.context(sessionID)) const bus = yield* recordedStepSettlementEvents(sessionID, assistant.id) @@ -3987,16 +3811,14 @@ describe("SessionRunnerLLM", () => { it.effect("does not recover context overflow after durable assistant output", () => Effect.gen(function* () { const session = yield* setup - yield* admit(session, "Fail after output") - - response = [ + yield* TestLLM.push([ LLMEvent.stepStart({ index: 0 }), LLMEvent.textStart({ id: "text-partial" }), LLMEvent.textDelta({ id: "text-partial", text: "Partial" }), LLMEvent.textEnd({ id: "text-partial" }), LLMEvent.providerError({ message: "prompt too long", classification: "context-overflow" }), - ] - expect((yield* session.resume(sessionID).pipe(Effect.flip)).message).toBe("prompt too long") + ]) + expect((yield* runPrompt(session, "Fail after output").pipe(Effect.flip)).message).toBe("prompt too long") expect(requests).toHaveLength(1) expect(yield* session.context(sessionID)).toMatchObject([ @@ -4014,11 +3836,10 @@ describe("SessionRunnerLLM", () => { it.effect("projects raw provider stream failures as terminal assistant step failures", () => Effect.gen(function* () { const session = yield* setup - yield* admit(session, "Fail raw stream durably") const failure = invalidRequest() - responseStream = Stream.fail(failure) + yield* TestLLM.push(Stream.fail(failure)) - expect(yield* session.resume(sessionID).pipe(Effect.flip)).toBe(failure) + expect(yield* runPrompt(session, "Fail raw stream durably").pipe(Effect.flip)).toBe(failure) yield* replaySessionProjection(sessionID) expect(yield* session.context(sessionID)).toMatchObject([ { type: "user", text: "Fail raw stream durably" }, @@ -4031,11 +3852,11 @@ describe("SessionRunnerLLM", () => { Effect.gen(function* () { const session = yield* setup yield* admit(session, "Retry transport") - responseStream = Stream.fail(providerUnavailable()) - response = reply.text("Recovered", "retry-success") + yield* TestLLM.push(Stream.fail(providerUnavailable())) + yield* TestLLM.push(TestLLM.text("Recovered", "retry-success")) const run = yield* session.resume(sessionID).pipe(Effect.forkChild) - while (requests.length < 1) yield* Effect.yieldNow + yield* TestLLM.wait(1) yield* TestClock.adjust("1999 millis") expect(requests).toHaveLength(1) yield* TestClock.adjust("1 millis") @@ -4058,11 +3879,11 @@ describe("SessionRunnerLLM", () => { Effect.gen(function* () { const session = yield* setup yield* admit(session, "Retry rate limit") - responseStream = Stream.fail(rateLimited(5_000)) - response = reply.text("Recovered", "retry-after-success") + yield* TestLLM.push(Stream.fail(rateLimited(5_000))) + yield* TestLLM.push(TestLLM.text("Recovered", "retry-after-success")) const run = yield* session.resume(sessionID).pipe(Effect.forkChild) - while (requests.length < 1) yield* Effect.yieldNow + yield* TestLLM.wait(1) yield* TestClock.adjust("4999 millis") expect(requests).toHaveLength(1) yield* TestClock.adjust("1 millis") @@ -4074,15 +3895,17 @@ describe("SessionRunnerLLM", () => { it.effect("does not retry eligible failures after observable output", () => Effect.gen(function* () { const session = yield* setup - yield* admit(session, "Do not replay partial output") const failure = rateLimited() - responseStream = Stream.fromIterable([ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.textStart({ id: "partial-rate-limit" }), - LLMEvent.textDelta({ id: "partial-rate-limit", text: "Partial" }), - ]).pipe(Stream.concat(Stream.fail(failure))) + yield* TestLLM.push( + TestLLM.failAfter( + failure, + LLMEvent.stepStart({ index: 0 }), + LLMEvent.textStart({ id: "partial-rate-limit" }), + LLMEvent.textDelta({ id: "partial-rate-limit", text: "Partial" }), + ), + ) - expect(yield* session.resume(sessionID).pipe(Effect.flip)).toBe(failure) + expect(yield* runPrompt(session, "Do not replay partial output").pipe(Effect.flip)).toBe(failure) expect(requests).toHaveLength(1) expect(yield* recordedEventTypes(sessionID)).not.toContain("session.retry.scheduled.1") expect(yield* session.context(sessionID)).toMatchObject([ @@ -4101,15 +3924,16 @@ describe("SessionRunnerLLM", () => { Effect.gen(function* () { const session = yield* setup yield* admit(session, "Exhaust retries") - streamFailure = providerUnavailable() + const failure = providerUnavailable() + yield* TestLLM.always(Stream.fail(failure)) const run = yield* session.resume(sessionID).pipe(Effect.forkChild) - while (requests.length < 1) yield* Effect.yieldNow + yield* TestLLM.wait(1) for (const [index, delay] of [2_000, 4_000, 8_000, 16_000].entries()) { yield* TestClock.adjust(delay) - while (requests.length < index + 2) yield* Effect.yieldNow + yield* TestLLM.wait(index + 2) } - expect(yield* Fiber.join(run).pipe(Effect.flip)).toBe(streamFailure) + expect(yield* Fiber.join(run).pipe(Effect.flip)).toBe(failure) expect(requests).toHaveLength(5) const database = (yield* Database.Service).db @@ -4150,11 +3974,11 @@ describe("SessionRunnerLLM", () => { ) yield* admit(session, "Retry without consuming a step") const failure = providerUnavailable() - responseStream = Stream.fail(failure) - responses = [reply.tool("call-after-retry", "echo", { text: "recovered" }), reply.stop()] + yield* TestLLM.push(Stream.fail(failure)) + yield* TestLLM.push(TestLLM.tool("call-after-retry", "echo", { text: "recovered" }), TestLLM.stop()) const run = yield* session.resume(sessionID).pipe(Effect.forkChild) - while (requests.length < 1) yield* Effect.yieldNow + yield* TestLLM.wait(1) yield* TestClock.adjust("2 seconds") yield* Fiber.join(run) @@ -4185,11 +4009,10 @@ describe("SessionRunnerLLM", () => { it.effect("does not retry non-eligible provider failures", () => Effect.gen(function* () { const session = yield* setup - yield* admit(session, "Do not retry") const failure = invalidRequest() - streamFailure = failure + yield* TestLLM.push(Stream.fail(failure)) - expect(yield* session.resume(sessionID).pipe(Effect.flip)).toBe(failure) + expect(yield* runPrompt(session, "Do not retry").pipe(Effect.flip)).toBe(failure) expect(requests).toHaveLength(1) expect(yield* recordedEventTypes(sessionID)).not.toContain("session.retry.scheduled.1") }), @@ -4198,24 +4021,25 @@ describe("SessionRunnerLLM", () => { it.effect("settles malformed streamed tool input before the provider failure", () => Effect.gen(function* () { const session = yield* setup - yield* admit(session, "Call a malformed tool") const failure = new LLMError({ module: "test", method: "stream", reason: new InvalidProviderOutputReason({ message: "Invalid JSON input for tool call echo" }), }) - responseStream = Stream.fromIterable([ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.toolInputStart({ id: "call-malformed", name: "echo" }), - LLMEvent.toolInputDelta({ id: "call-malformed", name: "echo", text: '{"text":"partial' }), - ]).pipe(Stream.concat(Stream.fail(failure))) + yield* TestLLM.push( + TestLLM.failAfter( + failure, + LLMEvent.stepStart({ index: 0 }), + LLMEvent.toolInputStart({ id: "call-malformed", name: "echo" }), + LLMEvent.toolInputDelta({ id: "call-malformed", name: "echo", text: '{"text":"partial' }), + ), + ) - expect(yield* session.resume(sessionID).pipe(Effect.flip)).toBe(failure) + expect(yield* runPrompt(session, "Call a malformed tool").pipe(Effect.flip)).toBe(failure) const assistant = requireAssistant(yield* session.context(sessionID)) - response = reply.stop() - yield* admit(session, "Continue") - yield* session.resume(sessionID) + yield* TestLLM.push(TestLLM.stop()) + yield* runPrompt(session, "Continue") expect(yield* recordedStepSettlementEvents(sessionID, assistant.id)).toMatchObject([ { type: "session.step.started.1" }, @@ -4237,12 +4061,10 @@ describe("SessionRunnerLLM", () => { it.effect("continues after malformed local tool input without exposing raw arguments", () => Effect.gen(function* () { const session = yield* setup - yield* admit(session, "Recover malformed tool input") const marker = "raw-malformed-marker" const raw = `{"text":"${marker}` - responses = [ - [ - LLMEvent.stepStart({ index: 0 }), + yield* TestLLM.push( + TestLLM.toolCalls( LLMEvent.toolInputStart({ id: "call-malformed", name: "echo" }), LLMEvent.toolInputDelta({ id: "call-malformed", name: "echo", text: raw }), LLMEvent.toolInputEnd({ id: "call-malformed", name: "echo" }), @@ -4251,13 +4073,11 @@ describe("SessionRunnerLLM", () => { name: "echo", raw, }), - LLMEvent.stepFinish({ index: 0, reason: { normalized: "tool-calls" } }), - LLMEvent.finish({ reason: { normalized: "tool-calls" } }), - ], - reply.stop(), - ] + ), + TestLLM.stop(), + ) - yield* session.resume(sessionID) + yield* runPrompt(session, "Recover malformed tool input") expect(requests).toHaveLength(2) expect(executions).toEqual([]) @@ -4313,7 +4133,7 @@ describe("SessionRunnerLLM", () => { }) if (!failed) throw new Error("Malformed tool assistant missing") expect(failed.error).toBeUndefined() - expect((yield* recordedStepSettlementEvents(sessionID, failed.id)).map((event) => event.type)).toEqual([ + expect(yield* recordedStepSettlementTypes(sessionID, failed.id)).toEqual([ "session.step.started.1", "session.tool.failed.2", "session.step.ended.1", @@ -4336,31 +4156,24 @@ describe("SessionRunnerLLM", () => { Effect.gen(function* () { const session = yield* setup yield* admit(session, "Run parallel tools") - toolExecutionGate = yield* Deferred.make() - toolExecutionsStarted = yield* Deferred.make() - toolExecutionsReady = 1 - responses = [ - [ - LLMEvent.stepStart({ index: 0 }), + const tools = yield* blockTools() + yield* TestLLM.push( + TestLLM.toolCalls( LLMEvent.toolCall({ id: "call-valid", name: "echo", input: { text: "valid" } }), LLMEvent.toolInputError({ id: "call-malformed", name: "echo", raw: '{"text":"partial', }), - LLMEvent.stepFinish({ index: 0, reason: { normalized: "tool-calls" } }), - LLMEvent.finish({ reason: { normalized: "tool-calls" } }), - ], - reply.stop(), - ] + ), + TestLLM.stop(), + ) const run = yield* session.resume(sessionID).pipe(Effect.forkChild) - yield* Deferred.await(toolExecutionsStarted) + yield* tools.started expect(requests).toHaveLength(1) - yield* Deferred.succeed(toolExecutionGate, undefined) + yield* tools.release yield* Fiber.join(run) - toolExecutionGate = undefined - toolExecutionsStarted = undefined expect(requests).toHaveLength(2) expect(executions).toEqual(["valid"]) @@ -4379,23 +4192,20 @@ describe("SessionRunnerLLM", () => { Effect.gen(function* () { const session = yield* setup yield* admit(session, "Interrupt malformed recovery") - toolExecutionGate = yield* Deferred.make() - toolExecutionsStarted = yield* Deferred.make() - toolExecutionsReady = 1 - response = [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.toolCall({ id: "call-valid", name: "echo", input: { text: "blocked" } }), - LLMEvent.toolInputError({ - id: "call-malformed", - name: "echo", - raw: '{"text":"partial', - }), - LLMEvent.stepFinish({ index: 0, reason: { normalized: "tool-calls" } }), - LLMEvent.finish({ reason: { normalized: "tool-calls" } }), - ] + const tools = yield* blockTools() + yield* TestLLM.push( + TestLLM.toolCalls( + LLMEvent.toolCall({ id: "call-valid", name: "echo", input: { text: "blocked" } }), + LLMEvent.toolInputError({ + id: "call-malformed", + name: "echo", + raw: '{"text":"partial', + }), + ), + ) const run = yield* session.resume(sessionID).pipe(Effect.forkChild) - yield* Deferred.await(toolExecutionsStarted) + yield* tools.started while ( !(yield* session.context(sessionID)).some( (message) => @@ -4405,8 +4215,6 @@ describe("SessionRunnerLLM", () => { ) yield* Effect.yieldNow yield* session.interrupt(sessionID) - toolExecutionGate = undefined - toolExecutionsStarted = undefined expect(yield* Fiber.await(run)).toMatchObject({ _tag: "Failure" }) expect(requests).toHaveLength(1) @@ -4427,19 +4235,21 @@ describe("SessionRunnerLLM", () => { it.effect("records malformed provider-executed input as executed", () => Effect.gen(function* () { const session = yield* setup - yield* admit(session, "Fail malformed hosted input") const failure = new LLMError({ module: "test", method: "stream", reason: new InvalidProviderOutputReason({ message: "Invalid hosted tool input" }), }) - responseStream = Stream.fromIterable([ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.toolInputStart({ id: "call-hosted", name: "web_search", providerExecuted: true }), - LLMEvent.toolInputDelta({ id: "call-hosted", name: "web_search", text: '{"query":"partial' }), - ]).pipe(Stream.concat(Stream.fail(failure))) + yield* TestLLM.push( + TestLLM.failAfter( + failure, + LLMEvent.stepStart({ index: 0 }), + LLMEvent.toolInputStart({ id: "call-hosted", name: "web_search", providerExecuted: true }), + LLMEvent.toolInputDelta({ id: "call-hosted", name: "web_search", text: '{"query":"partial' }), + ), + ) - expect(yield* session.resume(sessionID).pipe(Effect.flip)).toBe(failure) + expect(yield* runPrompt(session, "Fail malformed hosted input").pipe(Effect.flip)).toBe(failure) expect(requireAssistant(yield* session.context(sessionID))).toMatchObject({ error: { type: "provider.invalid-output", message: "Invalid hosted tool input" }, content: [ @@ -4457,22 +4267,24 @@ describe("SessionRunnerLLM", () => { it.effect("records a provider failure after malformed input", () => Effect.gen(function* () { const session = yield* setup - yield* admit(session, "Fail after malformed input") const failure = new LLMError({ module: "test", method: "stream", reason: new InvalidProviderOutputReason({ message: "Provider failed after malformed input" }), }) - responseStream = Stream.fromIterable([ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.toolInputError({ - id: "call-malformed", - name: "echo", - raw: '{"text":"partial', - }), - ]).pipe(Stream.concat(Stream.fail(failure))) + yield* TestLLM.push( + TestLLM.failAfter( + failure, + LLMEvent.stepStart({ index: 0 }), + LLMEvent.toolInputError({ + id: "call-malformed", + name: "echo", + raw: '{"text":"partial', + }), + ), + ) - expect(yield* session.resume(sessionID).pipe(Effect.flip)).toBe(failure) + expect(yield* runPrompt(session, "Fail after malformed input").pipe(Effect.flip)).toBe(failure) expect(requireAssistant(yield* session.context(sessionID))).toMatchObject({ error: { type: "provider.invalid-output", message: "Provider failed after malformed input" }, content: [ @@ -4491,25 +4303,22 @@ describe("SessionRunnerLLM", () => { it.effect("continues after repeated malformed tool input", () => Effect.gen(function* () { const session = yield* setup - yield* admit(session, "Keep producing malformed tools") - const malformed = (id: string) => [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.toolInputError({ - id, - name: "echo", - raw: '{"text":"partial', - }), - LLMEvent.stepFinish({ index: 0, reason: { normalized: "tool-calls" } }), - LLMEvent.finish({ reason: { normalized: "tool-calls" } }), - ] - responses = [ + const malformed = (id: string) => + TestLLM.toolCalls( + LLMEvent.toolInputError({ + id, + name: "echo", + raw: '{"text":"partial', + }), + ) + yield* TestLLM.push( malformed("call-first"), - reply.tool("call-valid-between", "echo", { text: "valid" }), + TestLLM.tool("call-valid-between", "echo", { text: "valid" }), malformed("call-second"), - reply.stop(), - ] + TestLLM.stop(), + ) - yield* session.resume(sessionID) + yield* runPrompt(session, "Keep producing malformed tools") expect(requests).toHaveLength(4) expect(executions).toEqual(["valid"]) @@ -4526,20 +4335,17 @@ describe("SessionRunnerLLM", () => { agent.steps = 2 }), ) - yield* admit(session, "Stop malformed tools at the step limit") - const malformed = (id: string) => [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.toolInputError({ - id, - name: "echo", - raw: '{"text":"partial', - }), - LLMEvent.stepFinish({ index: 0, reason: { normalized: "tool-calls" } }), - LLMEvent.finish({ reason: { normalized: "tool-calls" } }), - ] - responses = [malformed("call-first"), malformed("call-at-limit")] + const malformed = (id: string) => + TestLLM.toolCalls( + LLMEvent.toolInputError({ + id, + name: "echo", + raw: '{"text":"partial', + }), + ) + yield* TestLLM.push(malformed("call-first"), malformed("call-at-limit")) - yield* session.resume(sessionID) + yield* runPrompt(session, "Stop malformed tools at the step limit") expect(requests).toHaveLength(2) expect(requests[0]?.toolChoice).toBeUndefined() @@ -4552,28 +4358,23 @@ describe("SessionRunnerLLM", () => { Effect.gen(function* () { const session = yield* setup yield* admit(session, "Do not continue failed provider") - - toolExecutionGate = yield* Deferred.make() - toolExecutionsStarted = yield* Deferred.make() - toolExecutionsReady = 1 - response = [ + const tools = yield* blockTools() + yield* TestLLM.push([ LLMEvent.stepStart({ index: 0 }), LLMEvent.toolCall({ id: "call-before-provider-error", name: "echo", input: { text: "settled" } }), LLMEvent.providerError({ message: "Provider unavailable" }), - ] + ]) const run = yield* session.resume(sessionID).pipe(Effect.forkChild) - yield* Deferred.await(toolExecutionsStarted) - yield* Deferred.succeed(toolExecutionGate, undefined) + yield* tools.started + yield* tools.release expect((yield* Fiber.join(run).pipe(Effect.flip)).message).toBe("Provider unavailable") - toolExecutionGate = undefined - toolExecutionsStarted = undefined expect(requests).toHaveLength(1) expect(executions).toEqual(["settled"]) const context = yield* session.context(sessionID) const assistant = requireAssistant(context) - expect((yield* recordedStepSettlementEvents(sessionID, assistant.id)).map((event) => event.type)).toEqual([ + expect(yield* recordedStepSettlementTypes(sessionID, assistant.id)).toEqual([ "session.step.started.1", "session.tool.called.1", "session.tool.success.2", @@ -4585,15 +4386,15 @@ describe("SessionRunnerLLM", () => { it.effect("durably fails a hosted tool when its provider errors before returning a result", () => Effect.gen(function* () { const session = yield* setup - yield* admit(session, "Fail hosted tool durably") - - response = [ + yield* TestLLM.push([ LLMEvent.stepStart({ index: 0 }), hostedCall("call-hosted-provider-error", "effect"), LLMEvent.providerError({ message: "Provider unavailable" }), - ] + ]) - expect((yield* session.resume(sessionID).pipe(Effect.flip)).message).toBe("Provider unavailable") + expect((yield* runPrompt(session, "Fail hosted tool durably").pipe(Effect.flip)).message).toBe( + "Provider unavailable", + ) expect(requests).toHaveLength(1) const context = yield* session.context(sessionID) @@ -4605,7 +4406,7 @@ describe("SessionRunnerLLM", () => { }, ]) const assistant = requireAssistant(context) - expect((yield* recordedStepSettlementEvents(sessionID, assistant.id)).map((event) => event.type)).toEqual([ + expect(yield* recordedStepSettlementTypes(sessionID, assistant.id)).toEqual([ "session.step.started.1", "session.tool.called.1", "session.tool.failed.2", @@ -4617,14 +4418,15 @@ describe("SessionRunnerLLM", () => { it.effect("preserves a tool defect before provider failure settlement", () => Effect.gen(function* () { const session = yield* setup - yield* admit(session, "Defect while provider fails") - response = [ + yield* TestLLM.push([ LLMEvent.stepStart({ index: 0 }), LLMEvent.toolCall({ id: "call-defect-provider-error", name: "defect", input: {} }), LLMEvent.providerError({ message: "Provider unavailable" }), - ] + ]) - expect((yield* session.resume(sessionID).pipe(Effect.flip)).message).toBe("Provider unavailable") + expect((yield* runPrompt(session, "Defect while provider fails").pipe(Effect.flip)).message).toBe( + "Provider unavailable", + ) const context = yield* session.context(sessionID) const assistant = requireAssistant(context) @@ -4643,13 +4445,15 @@ describe("SessionRunnerLLM", () => { Effect.gen(function* () { const session = yield* setup yield* admit(session, "Storage fails while provider fails") - response = [ + yield* TestLLM.push([ LLMEvent.stepStart({ index: 0 }), LLMEvent.toolCall({ id: "call-store-provider-error", name: "storefail", input: {} }), LLMEvent.providerError({ message: "Provider unavailable" }), - ] + ]) - expect(yield* session.resume(sessionID).pipe(Effect.exit)).toMatchObject({ _tag: "Failure" }) + expect(yield* session.resume(sessionID).pipe(Effect.exit)).toMatchObject({ + _tag: "Failure", + }) expect(requireAssistant(yield* session.context(sessionID))).toMatchObject({ error: { type: "provider.unknown", message: "Provider unavailable" }, @@ -4660,10 +4464,11 @@ describe("SessionRunnerLLM", () => { it.effect("durably fails a hosted tool left unresolved at normal provider EOF", () => Effect.gen(function* () { const session = yield* setup - yield* admit(session, "Fail hosted tool at EOF") - response = [LLMEvent.stepStart({ index: 0 }), hostedCall("call-hosted-eof", "effect")] + yield* TestLLM.push([LLMEvent.stepStart({ index: 0 }), hostedCall("call-hosted-eof", "effect")]) - expect((yield* session.resume(sessionID).pipe(Effect.flip)).message).toBe("Provider did not return a tool result") + expect((yield* runPrompt(session, "Fail hosted tool at EOF").pipe(Effect.flip)).message).toBe( + "Provider did not return a tool result", + ) const assistant = requireAssistant(yield* session.context(sessionID)) const bus = yield* recordedStepSettlementEvents(sessionID, assistant.id) expect(bus.map((event) => event.type)).toEqual([ @@ -4692,15 +4497,9 @@ describe("SessionRunnerLLM", () => { it.effect("fails an unresolved hosted tool before one clean step end", () => Effect.gen(function* () { const session = yield* setup - yield* admit(session, "Settle hosted tool before ending") - response = [ - LLMEvent.stepStart({ index: 0 }), - hostedCall("call-hosted-clean-end", "effect"), - LLMEvent.stepFinish({ index: 0, reason: { normalized: "stop" } }), - LLMEvent.finish({ reason: { normalized: "stop" } }), - ] + yield* TestLLM.push(TestLLM.stop(hostedCall("call-hosted-clean-end", "effect"))) - yield* session.resume(sessionID) + yield* runPrompt(session, "Settle hosted tool before ending") const assistant = requireAssistant(yield* session.context(sessionID)) const bus = yield* recordedStepSettlementEvents(sessionID, assistant.id) @@ -4722,21 +4521,24 @@ describe("SessionRunnerLLM", () => { yield* admit(session, "Fail unresolved tools") const failure = invalidRequest() const providerFailed = yield* Deferred.make() - toolExecutionGate = yield* Deferred.make() - responseStream = Stream.concat( - Stream.fromIterable([ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.toolCall({ id: "call-local-raw-failure", name: "defect", input: {} }), - hostedCall("call-hosted-raw-failure-pair", "effect"), - ]), - Stream.fromEffect(Deferred.succeed(providerFailed, undefined)).pipe(Stream.flatMap(() => Stream.fail(failure))), + const tools = yield* blockTools() + yield* TestLLM.push( + Stream.concat( + Stream.fromIterable([ + LLMEvent.stepStart({ index: 0 }), + LLMEvent.toolCall({ id: "call-local-raw-failure", name: "defect", input: {} }), + hostedCall("call-hosted-raw-failure-pair", "effect"), + ]), + Stream.fromEffect(Deferred.succeed(providerFailed, undefined)).pipe( + Stream.flatMap(() => Stream.fail(failure)), + ), + ), ) const run = yield* session.resume(sessionID).pipe(Effect.forkChild) yield* Deferred.await(providerFailed) - yield* Deferred.succeed(toolExecutionGate, undefined) + yield* tools.release expect(yield* Fiber.join(run).pipe(Effect.flip)).toBe(failure) - toolExecutionGate = undefined const assistant = requireAssistant(yield* session.context(sessionID)) const bus = yield* recordedStepSettlementEvents(sessionID, assistant.id) @@ -4757,14 +4559,15 @@ describe("SessionRunnerLLM", () => { it.effect("durably fails a hosted tool left unresolved by a raw provider stream failure", () => Effect.gen(function* () { const session = yield* setup - yield* admit(session, "Fail hosted tool on raw failure") const failure = providerUnavailable() - responseStream = Stream.concat( - Stream.fromIterable([LLMEvent.stepStart({ index: 0 }), hostedCall("call-hosted-raw-failure", "effect")]), - Stream.fail(failure), + yield* TestLLM.push( + Stream.concat( + Stream.fromIterable([LLMEvent.stepStart({ index: 0 }), hostedCall("call-hosted-raw-failure", "effect")]), + Stream.fail(failure), + ), ) - expect(yield* session.resume(sessionID).pipe(Effect.flip)).toBe(failure) + expect(yield* runPrompt(session, "Fail hosted tool on raw failure").pipe(Effect.flip)).toBe(failure) expect(requests).toHaveLength(1) const assistant = requireAssistant(yield* session.context(sessionID)) const bus = yield* recordedStepSettlementEvents(sessionID, assistant.id) @@ -4793,15 +4596,13 @@ describe("SessionRunnerLLM", () => { it.effect("rejects a second text start before the open fragment ends", () => Effect.gen(function* () { const session = yield* setup - yield* admit(session, "Two blocks") - - response = [ + yield* TestLLM.push([ LLMEvent.stepStart({ index: 0 }), LLMEvent.textStart({ id: "text-1" }), LLMEvent.textStart({ id: "text-2" }), - ] + ]) - const defect = yield* session.resume(sessionID).pipe(Effect.catchDefect(Effect.succeed)) + const defect = yield* runPrompt(session, "Two blocks").pipe(Effect.catchDefect(Effect.succeed)) expect(defect).toBeInstanceOf(Error) if (!(defect instanceof Error)) return expect(defect.message).toBe("text start before end: text-2") @@ -4811,21 +4612,18 @@ describe("SessionRunnerLLM", () => { it.effect("projects sequential text fragments as separate content parts", () => Effect.gen(function* () { const session = yield* setup - yield* admit(session, "Two blocks") + yield* TestLLM.push( + TestLLM.stop( + LLMEvent.textStart({ id: "text-1" }), + LLMEvent.textDelta({ id: "text-1", text: "First" }), + LLMEvent.textEnd({ id: "text-1" }), + LLMEvent.textStart({ id: "text-2" }), + LLMEvent.textDelta({ id: "text-2", text: "Second" }), + LLMEvent.textEnd({ id: "text-2" }), + ), + ) - response = [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.textStart({ id: "text-1" }), - LLMEvent.textDelta({ id: "text-1", text: "First" }), - LLMEvent.textEnd({ id: "text-1" }), - LLMEvent.textStart({ id: "text-2" }), - LLMEvent.textDelta({ id: "text-2", text: "Second" }), - LLMEvent.textEnd({ id: "text-2" }), - LLMEvent.stepFinish({ index: 0, reason: { normalized: "stop" } }), - LLMEvent.finish({ reason: { normalized: "stop" } }), - ] - - yield* session.resume(sessionID) + yield* runPrompt(session, "Two blocks") expect(yield* session.context(sessionID)).toMatchObject([ { type: "user", text: "Two blocks" }, @@ -4855,7 +4653,7 @@ describe("SessionRunnerLLM", () => { it.effect("rejects duplicate streamed text starts", () => Effect.gen(function* () { const session = yield* setup - response = [LLMEvent.textStart({ id: "text-1" }), LLMEvent.textStart({ id: "text-1" })] + yield* TestLLM.push([LLMEvent.textStart({ id: "text-1" }), LLMEvent.textStart({ id: "text-1" })]) const defect = yield* session.resume(sessionID).pipe(Effect.catchDefect(Effect.succeed)) expect(defect).toBeInstanceOf(Error) @@ -4867,19 +4665,16 @@ describe("SessionRunnerLLM", () => { it.effect("transitions streamed raw tool input to parsed called input", () => Effect.gen(function* () { const session = yield* setup - yield* admit(session, "Call provider tool") + yield* TestLLM.push( + TestLLM.stop( + LLMEvent.toolInputStart({ id: "call-parsed", name: "web_search" }), + LLMEvent.toolInputDelta({ id: "call-parsed", name: "web_search", text: '{"query":"hello"}' }), + LLMEvent.toolInputEnd({ id: "call-parsed", name: "web_search" }), + hostedCall("call-parsed", "hello"), + ), + ) - response = [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.toolInputStart({ id: "call-parsed", name: "web_search" }), - LLMEvent.toolInputDelta({ id: "call-parsed", name: "web_search", text: '{"query":"hello"}' }), - LLMEvent.toolInputEnd({ id: "call-parsed", name: "web_search" }), - hostedCall("call-parsed", "hello"), - LLMEvent.stepFinish({ index: 0, reason: { normalized: "stop" } }), - LLMEvent.finish({ reason: { normalized: "stop" } }), - ] - - yield* session.resume(sessionID) + yield* runPrompt(session, "Call provider tool") expect(yield* session.context(sessionID)).toMatchObject([ { type: "user", text: "Call provider tool" }, @@ -4894,7 +4689,7 @@ describe("SessionRunnerLLM", () => { it.effect("rejects malformed streamed tool input ordering", () => Effect.gen(function* () { const session = yield* setup - response = [LLMEvent.toolInputDelta({ id: "call-1", name: "read", text: "{}" })] + yield* TestLLM.push([LLMEvent.toolInputDelta({ id: "call-1", name: "read", text: "{}" })]) const defect = yield* session.resume(sessionID).pipe(Effect.catchDefect(Effect.succeed)) expect(defect).toBeInstanceOf(Error)