From ad40cebbd31cf8cd9667a22ba58816b091c3c2fe Mon Sep 17 00:00:00 2001 From: Kit Langton Date: Tue, 7 Jul 2026 11:56:56 -0400 Subject: [PATCH] test(core): migrate runner delivery scenarios --- packages/core/test/session-runner.test.ts | 817 ++++++++++------------ 1 file changed, 387 insertions(+), 430 deletions(-) diff --git a/packages/core/test/session-runner.test.ts b/packages/core/test/session-runner.test.ts index 6f3544116d0..d2f00931ac9 100644 --- a/packages/core/test/session-runner.test.ts +++ b/packages/core/test/session-runner.test.ts @@ -2226,38 +2226,22 @@ describe("SessionRunnerLLM", () => { }), ) - it.effect("joins concurrent resume calls into one active provider run", () => + scenarioIt("joins concurrent resume calls into one active provider run", (scenario) => Effect.gen(function* () { yield* setup const session = yield* SessionV2.Service yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Run once" }), resume: false }) - requests.length = 0 - responses = undefined - response = [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.textStart({ id: "text-once" }), - LLMEvent.textDelta({ id: "text-once", text: "Once" }), - LLMEvent.textEnd({ id: "text-once" }), - LLMEvent.stepFinish({ index: 0, reason: "stop" }), - LLMEvent.finish({ reason: "stop" }), - ] - streamGate = yield* Deferred.make() - streamStarted = yield* Deferred.make() + yield* scenario.run(function* () { + const call = yield* scenario.llm.next() + const second = yield* session.resume(sessionID).pipe(Effect.forkChild) + yield* Effect.yieldNow + expect(yield* scenario.llm.requests).toHaveLength(1) + yield* call.respond.text("Once", { id: "text-once" }) + yield* Fiber.join(second) + }) - const first = yield* session.resume(sessionID).pipe(Effect.forkChild) - yield* Deferred.await(streamStarted) - const second = yield* session.resume(sessionID).pipe(Effect.forkChild) - yield* Effect.yieldNow - - expect(requests).toHaveLength(1) - yield* Deferred.succeed(streamGate, undefined) - yield* Fiber.join(first) - yield* Fiber.join(second) - streamGate = undefined - streamStarted = undefined - - expect(requests).toHaveLength(1) + expect(yield* scenario.llm.requests).toHaveLength(1) expect(yield* session.context(sessionID)).toMatchObject([ { type: "user", text: "Run once" }, { type: "assistant", finish: "stop", content: [{ type: "text", text: "Once" }] }, @@ -2292,191 +2276,153 @@ describe("SessionRunnerLLM", () => { }), ) - it.effect("promotes queued input after continuation ends", () => + scenarioIt("promotes queued input after continuation ends", (scenario) => Effect.gen(function* () { yield* setup const session = yield* SessionV2.Service yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Start working" }), resume: false }) - requests.length = 0 - responses = [ - [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.toolCall({ id: "call-echo", name: "echo", input: { text: "hello" } }), - LLMEvent.stepFinish({ index: 0, reason: "tool-calls" }), - LLMEvent.finish({ reason: "tool-calls" }), - ], - [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.stepFinish({ index: 0, reason: "stop" }), - LLMEvent.finish({ reason: "stop" }), - ], - [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.stepFinish({ index: 0, reason: "stop" }), - LLMEvent.finish({ reason: "stop" }), - ], - ] - streamGate = yield* Deferred.make() - streamStarted = yield* Deferred.make() + yield* scenario.run(function* () { + const first = yield* scenario.llm.next() + expect(userTexts(first.request)).toEqual(["Start working"]) + yield* session.prompt({ + sessionID, + prompt: PromptInput.Prompt.make({ text: "Wait until continuation ends" }), + delivery: "queue", + }) + yield* first.respond.toolCall("echo", { text: "hello" }, { id: "call-echo" }) - const first = yield* session.resume(sessionID).pipe(Effect.forkChild) - yield* Deferred.await(streamStarted) - yield* session.prompt({ - sessionID, - prompt: PromptInput.Prompt.make({ text: "Wait until continuation ends" }), - delivery: "queue", + const second = yield* scenario.llm.next() + expect(userTexts(second.request)).toEqual(["Start working"]) + yield* second.respond.stop() + + const third = yield* scenario.llm.next() + expect(userTexts(third.request)).toEqual(["Start working", "Wait until continuation ends"]) + yield* third.respond.stop() }) - yield* Deferred.succeed(streamGate, undefined) - 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(yield* scenario.llm.requests).toHaveLength(3) }), ) - it.effect("preserves durable queued input for a later wake after interruption", () => + effectIt.effect("preserves durable queued input for a later wake after interruption", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service - const { db } = yield* Database.Service - yield* session.prompt({ - sessionID, - prompt: PromptInput.Prompt.make({ text: "Interrupt current work" }), - resume: false, - }) + const interrupted = yield* Deferred.make>() + const retry = yield* Deferred.make() + const scenario = yield* RunnerScenario.make(() => + SessionV2.Service.use((session) => + Effect.gen(function* () { + yield* Deferred.succeed(interrupted, yield* session.resume(sessionID).pipe(Effect.exit)) + yield* Deferred.await(retry) + yield* session.resume(sessionID) + }), + ), + ) + yield* Effect.gen(function* () { + yield* setup + const session = yield* SessionV2.Service + const { db } = yield* Database.Service + yield* session.prompt({ + sessionID, + prompt: PromptInput.Prompt.make({ text: "Interrupt current work" }), + resume: false, + }) - requests.length = 0 - responses = [ - [], - [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.stepFinish({ index: 0, reason: "stop" }), - LLMEvent.finish({ reason: "stop" }), - ], - ] - streamGate = yield* Deferred.make() - streamStarted = yield* Deferred.make() + yield* scenario.run(function* () { + const first = yield* scenario.llm.next() + expect(userTexts(first.request)).toEqual(["Interrupt current work"]) + yield* session.prompt({ + sessionID, + prompt: PromptInput.Prompt.make({ text: "Run after interrupt" }), + delivery: "queue", + }) + yield* session.interrupt(sessionID) + expect(yield* Deferred.await(interrupted)).toMatchObject({ _tag: "Failure" }) + expect(yield* SessionInput.hasPending(db, sessionID, "queue")).toBe(true) - const run = yield* session.resume(sessionID).pipe(Effect.forkChild) - yield* Deferred.await(streamStarted) - yield* session.prompt({ - sessionID, - prompt: PromptInput.Prompt.make({ text: "Run after interrupt" }), - delivery: "queue", - }) - yield* session.interrupt(sessionID) - expect(yield* Fiber.await(run)).toMatchObject({ _tag: "Failure" }) - expect(requests).toHaveLength(1) - expect(yield* SessionInput.hasPending(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* Fiber.join(resumed) - streamGate = undefined - streamStarted = undefined + yield* Deferred.succeed(retry, undefined) + const second = yield* scenario.llm.next() + expect(userTexts(second.request)).toEqual(["Interrupt current work", "Run after interrupt"]) + yield* second.respond.stop() + }) - expect(requests).toHaveLength(2) - expect(userTexts(requests[0]!)).toEqual(["Interrupt current work"]) - expect(userTexts(requests[1]!)).toEqual(["Interrupt current work", "Run after interrupt"]) + expect(yield* scenario.llm.requests).toHaveLength(2) + }).pipe(Effect.provide(testLayerWith(scenario.llm.layer))) }), ) - it.effect("preserves durable steering input for a later resume after interruption", () => + effectIt.effect("preserves durable steering input for a later resume after interruption", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service - const { db } = yield* Database.Service - yield* session.prompt({ - sessionID, - prompt: PromptInput.Prompt.make({ text: "Interrupt current work" }), - resume: false, - }) + const interrupted = yield* Deferred.make>() + const retry = yield* Deferred.make() + const scenario = yield* RunnerScenario.make(() => + SessionV2.Service.use((session) => + Effect.gen(function* () { + yield* Deferred.succeed(interrupted, yield* session.resume(sessionID).pipe(Effect.exit)) + yield* Deferred.await(retry) + yield* session.resume(sessionID) + }), + ), + ) + yield* Effect.gen(function* () { + yield* setup + const session = yield* SessionV2.Service + const { db } = yield* Database.Service + yield* session.prompt({ + sessionID, + prompt: PromptInput.Prompt.make({ text: "Interrupt current work" }), + resume: false, + }) - requests.length = 0 - responses = [ - [], - [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.stepFinish({ index: 0, reason: "stop" }), - LLMEvent.finish({ reason: "stop" }), - ], - ] - streamGate = yield* Deferred.make() - streamStarted = yield* Deferred.make() + yield* scenario.run(function* () { + const first = yield* scenario.llm.next() + expect(userTexts(first.request)).toEqual(["Interrupt current work"]) + yield* session.prompt({ + sessionID, + prompt: PromptInput.Prompt.make({ text: "Steer after interrupt" }), + }) + yield* session.interrupt(sessionID) + expect(yield* Deferred.await(interrupted)).toMatchObject({ _tag: "Failure" }) + expect(yield* SessionInput.hasPending(db, sessionID, "steer")).toBe(true) - const run = yield* session.resume(sessionID).pipe(Effect.forkChild) - yield* Deferred.await(streamStarted) - yield* session.prompt({ - sessionID, - prompt: PromptInput.Prompt.make({ text: "Steer after interrupt" }), - }) - yield* session.interrupt(sessionID) - expect(yield* Fiber.await(run)).toMatchObject({ _tag: "Failure" }) - expect(requests).toHaveLength(1) - expect(yield* SessionInput.hasPending(db, sessionID, "steer")).toBe(true) + yield* Deferred.succeed(retry, undefined) + const second = yield* scenario.llm.next() + expect(userTexts(second.request)).toEqual(["Interrupt current work", "Steer after interrupt"]) + yield* second.respond.stop() + }) - const resumed = yield* session.resume(sessionID).pipe(Effect.forkChild) - while (requests.length < 2) yield* Effect.yieldNow - yield* Deferred.succeed(streamGate, undefined) - 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(yield* scenario.llm.requests).toHaveLength(2) + }).pipe(Effect.provide(testLayerWith(scenario.llm.layer))) }), ) - it.effect("promotes queued inputs one at a time in FIFO order", () => + scenarioIt("promotes queued inputs one at a time in FIFO order", (scenario) => Effect.gen(function* () { yield* setup const session = yield* SessionV2.Service yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Start working" }), resume: false }) - requests.length = 0 - responses = [ - [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.stepFinish({ index: 0, reason: "stop" }), - LLMEvent.finish({ reason: "stop" }), - ], - [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.stepFinish({ index: 0, reason: "stop" }), - LLMEvent.finish({ reason: "stop" }), - ], - [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.stepFinish({ index: 0, reason: "stop" }), - LLMEvent.finish({ reason: "stop" }), - ], - ] - streamGate = yield* Deferred.make() - streamStarted = yield* Deferred.make() + yield* scenario.run(function* () { + const first = yield* scenario.llm.next() + expect(userTexts(first.request)).toEqual(["Start working"]) + yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Queue first" }), delivery: "queue" }) + yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Queue second" }), delivery: "queue" }) + yield* first.respond.stop() - const first = yield* session.resume(sessionID).pipe(Effect.forkChild) - yield* Deferred.await(streamStarted) - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Queue first" }), delivery: "queue" }) - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Queue second" }), delivery: "queue" }) - yield* Deferred.succeed(streamGate, undefined) - yield* Fiber.join(first) - streamGate = undefined - streamStarted = undefined + const second = yield* scenario.llm.next() + expect(userTexts(second.request)).toEqual(["Start working", "Queue first"]) + yield* second.respond.stop() - 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"]) + const third = yield* scenario.llm.next() + expect(userTexts(third.request)).toEqual(["Start working", "Queue first", "Queue second"]) + yield* third.respond.stop() + }) + + expect(yield* scenario.llm.requests).toHaveLength(3) }), ) - it.effect("promotes queued input after steering continuation ends", () => + scenarioIt("promotes queued input after steering continuation ends", (scenario) => Effect.gen(function* () { yield* setup const session = yield* SessionV2.Service @@ -2488,166 +2434,125 @@ describe("SessionRunnerLLM", () => { resume: false, }) - requests.length = 0 - responses = [ - [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.stepFinish({ index: 0, reason: "stop" }), - LLMEvent.finish({ reason: "stop" }), - ], - [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.stepFinish({ index: 0, reason: "stop" }), - LLMEvent.finish({ reason: "stop" }), - ], - ] + yield* scenario.run(function* () { + const first = yield* scenario.llm.next() + expect(userTexts(first.request)).toEqual(["Start steering"]) + yield* first.respond.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"]) - }), - ) - - it.effect("promotes steers before the next queued input", () => - Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Start working" }), resume: false }) - - requests.length = 0 - responses = [ - [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.stepFinish({ index: 0, reason: "stop" }), - LLMEvent.finish({ reason: "stop" }), - ], - [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.stepFinish({ index: 0, reason: "stop" }), - LLMEvent.finish({ reason: "stop" }), - ], - [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.stepFinish({ index: 0, reason: "stop" }), - LLMEvent.finish({ reason: "stop" }), - ], - [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.stepFinish({ index: 0, reason: "stop" }), - LLMEvent.finish({ reason: "stop" }), - ], - ] - const firstGate = yield* Deferred.make() - const secondGate = yield* Deferred.make() - streamGate = firstGate - - const first = yield* session.resume(sessionID).pipe(Effect.forkChild) - while (requests.length < 1) yield* Effect.yieldNow - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Queue first" }), delivery: "queue" }) - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Queue second" }), delivery: "queue" }) - streamGate = secondGate - yield* Deferred.succeed(firstGate, undefined) - while (requests.length < 2) yield* Effect.yieldNow - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Steer before next queued input" }) }) - yield* session.prompt({ - sessionID, - prompt: PromptInput.Prompt.make({ text: "Also steer before next queued input" }), + const second = yield* scenario.llm.next() + expect(userTexts(second.request)).toEqual(["Start steering", "Queue for later"]) + yield* second.respond.stop() }) - yield* Deferred.succeed(secondGate, undefined) - 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([ - "Start working", - "Queue first", - "Steer before next queued input", - "Also steer before next queued input", - ]) - expect(userTexts(requests[3]!)).toEqual([ - "Start working", - "Queue first", - "Steer before next queued input", - "Also steer before next queued input", - "Queue second", - ]) + expect(yield* scenario.llm.requests).toHaveLength(2) }), ) - it.effect("coalesces multiple active steering prompts into one continuation turn", () => + scenarioIt("promotes steers before the next queued input", (scenario) => Effect.gen(function* () { yield* setup const session = yield* SessionV2.Service yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Start working" }), resume: false }) - requests.length = 0 - responses = [ - [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.stepFinish({ index: 0, reason: "stop" }), - LLMEvent.finish({ reason: "stop" }), - ], - [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.stepFinish({ index: 0, reason: "stop" }), - LLMEvent.finish({ reason: "stop" }), - ], - ] - streamGate = yield* Deferred.make() - streamStarted = yield* Deferred.make() + yield* scenario.run(function* () { + const first = yield* scenario.llm.next() + expect(userTexts(first.request)).toEqual(["Start working"]) + yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Queue first" }), delivery: "queue" }) + yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Queue second" }), delivery: "queue" }) + yield* first.respond.stop() - const first = yield* session.resume(sessionID).pipe(Effect.forkChild) - yield* Deferred.await(streamStarted) - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "First steer" }) }) - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Second steer" }) }) - yield* Deferred.succeed(streamGate, undefined) - yield* Fiber.join(first) - streamGate = undefined - streamStarted = undefined - yield* Effect.yieldNow + const second = yield* scenario.llm.next() + expect(userTexts(second.request)).toEqual(["Start working", "Queue first"]) + yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Steer before next queued input" }) }) + yield* session.prompt({ + sessionID, + prompt: PromptInput.Prompt.make({ text: "Also steer before next queued input" }), + }) + yield* second.respond.stop() - expect(requests).toHaveLength(2) - expect(userTexts(requests[1]!)).toEqual(["Start working", "First steer", "Second steer"]) + const third = yield* scenario.llm.next() + expect(userTexts(third.request)).toEqual([ + "Start working", + "Queue first", + "Steer before next queued input", + "Also steer before next queued input", + ]) + yield* third.respond.stop() + + const fourth = yield* scenario.llm.next() + expect(userTexts(fourth.request)).toEqual([ + "Start working", + "Queue first", + "Steer before next queued input", + "Also steer before next queued input", + "Queue second", + ]) + yield* fourth.respond.stop() + }) + + expect(yield* scenario.llm.requests).toHaveLength(4) + }), + ) + + scenarioIt("coalesces multiple active steering prompts into one continuation turn", (scenario) => + Effect.gen(function* () { + yield* setup + const session = yield* SessionV2.Service + yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Start working" }), resume: false }) + + yield* scenario.run(function* () { + const first = yield* scenario.llm.next() + yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "First steer" }) }) + yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Second steer" }) }) + yield* first.respond.stop() + + const second = yield* scenario.llm.next() + expect(userTexts(second.request)).toEqual(["Start working", "First steer", "Second steer"]) + yield* second.respond.stop() + }) + + expect(yield* scenario.llm.requests).toHaveLength(2) yield* (yield* SessionExecution.Service).wake(sessionID) - yield* Effect.yieldNow - expect(requests).toHaveLength(2) + yield* session.wait(sessionID) + expect(yield* scenario.llm.requests).toHaveLength(2) }), ) - it.effect("runs steering input accepted while the active provider turn fails", () => + effectIt.effect("runs steering input accepted while the active provider turn fails", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Start working" }), resume: false }) + const failed = yield* Deferred.make>() + const scenario = yield* RunnerScenario.make(() => + SessionV2.Service.use((session) => + Effect.gen(function* () { + yield* Deferred.succeed(failed, yield* session.resume(sessionID).pipe(Effect.exit)) + yield* session.wait(sessionID) + }), + ), + ) + yield* Effect.gen(function* () { + yield* setup + const session = yield* SessionV2.Service + yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Start working" }), resume: false }) + const failure = invalidRequest() - requests.length = 0 - responses = undefined - response = [] - streamFailure = invalidRequest() - streamGate = yield* Deferred.make() - streamStarted = yield* Deferred.make() + yield* scenario.run(function* () { + const first = yield* scenario.llm.next() + yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Recover with this" }) }) + yield* first.respond.fail(failure) + const exit = yield* Deferred.await(failed) + expect(Exit.isFailure(exit) && Cause.squash(exit.cause)).toBe(failure) - const first = yield* session.resume(sessionID).pipe(Effect.forkChild) - yield* Deferred.await(streamStarted) - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Recover with this" }) }) - yield* Deferred.succeed(streamGate, undefined) - expect(yield* Fiber.join(first).pipe(Effect.flip)).toBe(streamFailure) + const second = yield* scenario.llm.next() + expect(userTexts(second.request)).toEqual(["Start working", "Recover with this"]) + yield* second.respond.stop() + }) - streamFailure = undefined - streamGate = undefined - streamStarted = undefined - yield* Effect.yieldNow - - expect(requests).toHaveLength(2) - expect(userTexts(requests[1]!)).toEqual(["Start working", "Recover with this"]) + expect(yield* scenario.llm.requests).toHaveLength(2) + }).pipe(Effect.provide(testLayerWith(scenario.llm.layer))) }), ) - it.effect("durably fails local tools left running by a prior process before continuing", () => + scenarioIt("durably fails local tools left running by a prior process before continuing", (scenario) => Effect.gen(function* () { yield* setup const session = yield* SessionV2.Service @@ -2684,12 +2589,13 @@ describe("SessionRunnerLLM", () => { input: { text: "stale" }, executed: false, }) - requests.length = 0 - response = [] - yield* session.resume(sessionID) + yield* scenario.run(function* () { + const call = yield* scenario.llm.next() + expect(call.request.messages.map((message) => message.role)).toEqual(["user", "assistant", "tool"]) + yield* call.respond.events() + }) - expect(requests).toHaveLength(1) - expect(requests[0]?.messages.map((message) => message.role)).toEqual(["user", "assistant", "tool"]) + expect(yield* scenario.llm.requests).toHaveLength(1) expect(yield* session.context(sessionID)).toMatchObject([ { type: "user", text: "Recover interrupted tool" }, { @@ -2709,7 +2615,7 @@ describe("SessionRunnerLLM", () => { }), ) - it.effect("durably fails hosted tools left running by a prior process before continuing inline", () => + scenarioIt("durably fails hosted tools left running by a prior process before continuing inline", (scenario) => Effect.gen(function* () { yield* setup const session = yield* SessionV2.Service @@ -2747,25 +2653,26 @@ describe("SessionRunnerLLM", () => { executed: true, state: { itemId: "call-hosted-interrupted" }, }) - requests.length = 0 - response = [] - yield* session.resume(sessionID) + yield* scenario.run(function* () { + const call = yield* scenario.llm.next() + expect(call.request.messages.map((message) => message.role)).toEqual(["user", "assistant"]) + expect(call.request.messages[1]?.content).toMatchObject([ + { + type: "tool-call", + id: "call-hosted-interrupted", + providerExecuted: true, + providerMetadata: { fake: { itemId: "call-hosted-interrupted" } }, + }, + { type: "tool-result", id: "call-hosted-interrupted", providerExecuted: true, result: { type: "error" } }, + ]) + yield* call.respond.events() + }) - expect(requests).toHaveLength(1) - expect(requests[0]?.messages.map((message) => message.role)).toEqual(["user", "assistant"]) - expect(requests[0]?.messages[1]?.content).toMatchObject([ - { - type: "tool-call", - id: "call-hosted-interrupted", - providerExecuted: true, - providerMetadata: { fake: { itemId: "call-hosted-interrupted" } }, - }, - { type: "tool-result", id: "call-hosted-interrupted", providerExecuted: true, result: { type: "error" } }, - ]) + expect(yield* scenario.llm.requests).toHaveLength(1) }), ) - it.effect("durably fails pending tool input left by a prior process before continuing", () => + scenarioIt("durably fails pending tool input left by a prior process before continuing", (scenario) => Effect.gen(function* () { yield* setup const session = yield* SessionV2.Service @@ -2789,12 +2696,13 @@ describe("SessionRunnerLLM", () => { callID: "call-pending-interrupted", name: "echo", }) - requests.length = 0 - response = [] - yield* session.resume(sessionID) + yield* scenario.run(function* () { + const call = yield* scenario.llm.next() + expect(call.request.messages.map((message) => message.role)).toEqual(["user", "assistant", "tool"]) + yield* call.respond.events() + }) - expect(requests).toHaveLength(1) - expect(requests[0]?.messages.map((message) => message.role)).toEqual(["user", "assistant", "tool"]) + expect(yield* scenario.llm.requests).toHaveLength(1) 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" } }] }, @@ -2802,57 +2710,79 @@ describe("SessionRunnerLLM", () => { }), ) - it.effect("promotes the first queued input when woken while idle", () => + effectIt.effect("promotes the first queued input when woken while idle", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service - yield* session.prompt({ - sessionID, - prompt: PromptInput.Prompt.make({ text: "Wait in queue" }), - delivery: "queue", - resume: false, - }) + const scenario = yield* RunnerScenario.make(() => + Effect.gen(function* () { + const execution = yield* SessionExecution.Service + yield* execution.wake(sessionID) + yield* execution.awaitIdle(sessionID) + }), + ) + yield* Effect.gen(function* () { + yield* setup + const session = yield* SessionV2.Service + yield* session.prompt({ + sessionID, + prompt: PromptInput.Prompt.make({ text: "Wait in queue" }), + delivery: "queue", + resume: false, + }) - requests.length = 0 - yield* (yield* SessionExecution.Service).wake(sessionID) - yield* Effect.yieldNow + yield* scenario.run(function* () { + const call = yield* scenario.llm.next() + expect(userTexts(call.request)).toEqual(["Wait in queue"]) + yield* call.respond.events() + }) - expect(requests).toHaveLength(1) - expect(userTexts(requests[0]!)).toEqual(["Wait in queue"]) + expect(yield* scenario.llm.requests).toHaveLength(1) + }).pipe(Effect.provide(testLayerWith(scenario.llm.layer))) }), ) - it.effect("retries inbox input after prompt projection rolls back", () => + effectIt.effect("retries inbox input after prompt projection rolls back", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service - const events = yield* EventV2.Service const defect = new Error("fail after prompt promotion") let fail = true - yield* events.project(SessionEvent.PromptPromoted, () => (fail ? Effect.die(defect) : Effect.void)) - yield* session.prompt({ - sessionID, - prompt: PromptInput.Prompt.make({ text: "Recover promoted input" }), - resume: false, - }) + const rolledBack = yield* Deferred.make() + const retry = yield* Deferred.make() + const scenario = yield* RunnerScenario.make(() => + Effect.gen(function* () { + const session = yield* SessionV2.Service + yield* Deferred.succeed( + rolledBack, + yield* session.resume(sessionID).pipe(Effect.catchDefect(Effect.succeed)), + ) + yield* Deferred.await(retry) + const execution = yield* SessionExecution.Service + yield* execution.wake(sessionID) + yield* execution.awaitIdle(sessionID) + }), + ) + yield* Effect.gen(function* () { + yield* setup + const session = yield* SessionV2.Service + const events = yield* EventV2.Service + yield* events.project(SessionEvent.PromptPromoted, () => (fail ? Effect.die(defect) : Effect.void)) + yield* session.prompt({ + sessionID, + prompt: PromptInput.Prompt.make({ text: "Recover promoted input" }), + resume: false, + }) - expect(yield* session.resume(sessionID).pipe(Effect.catchDefect(Effect.succeed))).toBe(defect) - fail = false - requests.length = 0 - response = [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.stepFinish({ index: 0, reason: "stop" }), - LLMEvent.finish({ reason: "stop" }), - ] - - yield* (yield* SessionExecution.Service).wake(sessionID) - while (requests.length === 0) yield* Effect.yieldNow - - expect(userTexts(requests[0]!)).toEqual(["Recover promoted input"]) + yield* scenario.run(function* () { + expect(yield* Deferred.await(rolledBack)).toBe(defect) + fail = false + yield* Deferred.succeed(retry, undefined) + const call = yield* scenario.llm.next() + expect(userTexts(call.request)).toEqual(["Recover promoted input"]) + yield* call.respond.stop() + }) + }).pipe(Effect.provide(testLayerWith(scenario.llm.layer))) }), ) - it.effect("does not strand a committed promotion when a post-commit listener defects", () => + scenarioIt("does not strand a committed promotion when a post-commit listener defects", (scenario) => Effect.gen(function* () { yield* setup const session = yield* SessionV2.Service @@ -2868,11 +2798,13 @@ describe("SessionRunnerLLM", () => { resume: false, }) - requests.length = 0 - yield* session.resume(sessionID) + yield* scenario.run(function* () { + const call = yield* scenario.llm.next() + expect(userTexts(call.request)).toEqual(["Run committed promotion"]) + yield* call.respond.events() + }) - expect(requests).toHaveLength(1) - expect(userTexts(requests[0]!)).toEqual(["Run committed promotion"]) + expect(yield* scenario.llm.requests).toHaveLength(1) }), ) @@ -2913,68 +2845,93 @@ describe("SessionRunnerLLM", () => { }), ) - it.effect("bounds 64-character session prompt cache keys", () => + effectIt.effect("bounds 64-character session prompt cache keys", () => Effect.gen(function* () { - yield* setup const longSessionID = SessionV2.ID.make(`ses_${"a".repeat(64)}`) const otherLongSessionID = SessionV2.ID.make(`ses_${"b".repeat(64)}`) - yield* insertSession(longSessionID) - yield* insertSession(otherLongSessionID) - const session = yield* SessionV2.Service - yield* session.prompt({ - sessionID: longSessionID, - prompt: PromptInput.Prompt.make({ text: "Run long session" }), - resume: false, - }) - yield* session.prompt({ - sessionID: otherLongSessionID, - prompt: PromptInput.Prompt.make({ text: "Run other long session" }), - resume: false, - }) + const scenario = yield* RunnerScenario.make(() => + SessionV2.Service.use((session) => + session.resume(longSessionID).pipe(Effect.andThen(session.resume(otherLongSessionID))), + ), + ) + yield* Effect.gen(function* () { + yield* setup + yield* insertSession(longSessionID) + yield* insertSession(otherLongSessionID) + const session = yield* SessionV2.Service + yield* session.prompt({ + sessionID: longSessionID, + prompt: PromptInput.Prompt.make({ text: "Run long session" }), + resume: false, + }) + yield* session.prompt({ + sessionID: otherLongSessionID, + prompt: PromptInput.Prompt.make({ text: "Run other long session" }), + resume: false, + }) - requests.length = 0 - yield* session.resume(longSessionID) - yield* session.resume(otherLongSessionID) + yield* scenario.run(function* () { + yield* (yield* scenario.llm.next()).respond.events() + yield* (yield* scenario.llm.next()).respond.events() + }) - const keys = requests.map((request) => request.providerOptions?.openai?.promptCacheKey) - expect(keys).toEqual([longSessionID.slice(4), otherLongSessionID.slice(4)]) - expect(keys.every((key) => typeof key === "string" && key.length === 64)).toBe(true) - expect(keys[0]).not.toBe(keys[1]) + const keys = (yield* scenario.llm.requests).map( + (request) => request.providerOptions?.openai?.promptCacheKey, + ) + expect(keys).toEqual([longSessionID.slice(4), otherLongSessionID.slice(4)]) + expect(keys.every((key) => typeof key === "string" && key.length === 64)).toBe(true) + expect(keys[0]).not.toBe(keys[1]) + }).pipe(Effect.provide(testLayerWith(scenario.llm.layer))) }), ) - it.effect("fans out one failed run and allows a later retry", () => + effectIt.effect("fans out one failed run and allows a later retry", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service - yield* session.prompt({ - sessionID, - prompt: PromptInput.Prompt.make({ text: "Retry after failure" }), - resume: false, - }) + const join = yield* Deferred.make() + const joined = yield* Deferred.make< + readonly [ + Exit.Exit, + Exit.Exit, + ] + >() + const retry = yield* Deferred.make() + const scenario = yield* RunnerScenario.make(() => + SessionV2.Service.use((session) => + Effect.gen(function* () { + const first = yield* session.resume(sessionID).pipe(Effect.forkChild) + yield* Deferred.await(join) + const second = yield* session.resume(sessionID).pipe(Effect.forkChild) + yield* Deferred.succeed(joined, yield* Effect.all([Fiber.await(first), Fiber.await(second)])) + yield* Deferred.await(retry) + yield* session.resume(sessionID) + }), + ), + ) + yield* Effect.gen(function* () { + yield* setup + const session = yield* SessionV2.Service + yield* session.prompt({ + sessionID, + prompt: PromptInput.Prompt.make({ text: "Retry after failure" }), + resume: false, + }) + const failure = invalidRequest() - requests.length = 0 - responses = undefined - response = [] - streamFailure = invalidRequest() - streamGate = yield* Deferred.make() - streamStarted = yield* Deferred.make() + yield* scenario.run(function* () { + const first = yield* scenario.llm.next() + yield* Deferred.succeed(join, undefined) + yield* Effect.yieldNow + expect(yield* scenario.llm.requests).toHaveLength(1) + yield* first.respond.fail(failure) + const [firstExit, secondExit] = yield* Deferred.await(joined) + expect(secondExit).toEqual(firstExit) - const first = yield* session.resume(sessionID).pipe(Effect.forkChild) - yield* Deferred.await(streamStarted) - const second = yield* session.resume(sessionID).pipe(Effect.forkChild) - yield* Effect.yieldNow + yield* Deferred.succeed(retry, undefined) + yield* (yield* scenario.llm.next()).respond.events() + }) - expect(requests).toHaveLength(1) - yield* Deferred.succeed(streamGate, undefined) - const [firstExit, secondExit] = yield* Effect.all([Fiber.await(first), Fiber.await(second)]) - expect(secondExit).toEqual(firstExit) - - streamFailure = undefined - streamGate = undefined - streamStarted = undefined - yield* session.resume(sessionID) - expect(requests).toHaveLength(2) + expect(yield* scenario.llm.requests).toHaveLength(2) + }).pipe(Effect.provide(testLayerWith(scenario.llm.layer))) }), )