diff --git a/packages/core/test/session-runner.test.ts b/packages/core/test/session-runner.test.ts index d2f00931ac9..7a8ec837ad6 100644 --- a/packages/core/test/session-runner.test.ts +++ b/packages/core/test/session-runner.test.ts @@ -3633,21 +3633,26 @@ describe("SessionRunnerLLM", () => { }), ) - it.effect("projects provider errors as terminal assistant step failures", () => + scenarioIt("projects provider errors as terminal assistant step failures", (scenario) => Effect.gen(function* () { yield* setup const session = yield* SessionV2.Service yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Fail durably" }), resume: false }) - requests.length = 0 - responses = undefined - streamGate = undefined - streamStarted = undefined - response = [LLMEvent.stepStart({ index: 0 }), LLMEvent.providerError({ message: "Provider unavailable" })] + expect( + ( + yield* scenario + .run(function* () { + yield* (yield* scenario.llm.next()).respond.events( + LLMEvent.stepStart({ index: 0 }), + LLMEvent.providerError({ message: "Provider unavailable" }), + ) + }) + .pipe(Effect.flip) + ).message, + ).toBe("Provider unavailable") - expect((yield* session.resume(sessionID).pipe(Effect.flip)).message).toBe("Provider unavailable") - - expect(requests).toHaveLength(1) + expect(yield* scenario.llm.requests).toHaveLength(1) expect(yield* session.context(sessionID)).toMatchObject([ { type: "user", text: "Fail durably" }, { type: "assistant", finish: "error", error: { type: "provider.unknown", message: "Provider unavailable" } }, @@ -3655,18 +3660,25 @@ describe("SessionRunnerLLM", () => { }), ) - it.effect("projects provider errors emitted before assistant step start", () => + scenarioIt("projects provider errors emitted before assistant step start", (scenario) => Effect.gen(function* () { yield* setup const session = yield* SessionV2.Service yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Fail before step" }), resume: false }) - requests.length = 0 - response = [LLMEvent.providerError({ message: "Provider unavailable" })] + expect( + ( + yield* scenario + .run(function* () { + yield* (yield* scenario.llm.next()).respond.events( + LLMEvent.providerError({ message: "Provider unavailable" }), + ) + }) + .pipe(Effect.flip) + ).message, + ).toBe("Provider unavailable") - expect((yield* session.resume(sessionID).pipe(Effect.flip)).message).toBe("Provider unavailable") - - expect(requests).toHaveLength(1) + expect(yield* scenario.llm.requests).toHaveLength(1) expect(yield* session.context(sessionID)).toMatchObject([ { type: "user", text: "Fail before step" }, { type: "assistant", finish: "error", error: { type: "provider.unknown", message: "Provider unavailable" } }, @@ -3674,24 +3686,30 @@ describe("SessionRunnerLLM", () => { }), ) - it.effect("projects content-filter finishes as visible terminal failures", () => + scenarioIt("projects content-filter finishes as visible terminal failures", (scenario) => Effect.gen(function* () { yield* setup const session = yield* SessionV2.Service yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Blocked response" }), resume: false }) - response = [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.textStart({ id: "partial" }), - LLMEvent.textDelta({ id: "partial", text: "Partial" }), - LLMEvent.stepFinish({ - index: 0, - reason: "content-filter", - usage: { nonCachedInputTokens: 8, outputTokens: 3, reasoningTokens: 1 }, - }), - LLMEvent.finish({ reason: "content-filter" }), - ] - - expect((yield* session.resume(sessionID).pipe(Effect.flip)).message).toBe("Provider blocked the response") + expect( + ( + yield* scenario + .run(function* () { + yield* (yield* scenario.llm.next()).respond.events( + LLMEvent.stepStart({ index: 0 }), + LLMEvent.textStart({ id: "partial" }), + LLMEvent.textDelta({ id: "partial", text: "Partial" }), + LLMEvent.stepFinish({ + index: 0, + reason: "content-filter", + usage: { nonCachedInputTokens: 8, outputTokens: 3, reasoningTokens: 1 }, + }), + LLMEvent.finish({ reason: "content-filter" }), + ) + }) + .pipe(Effect.flip) + ).message, + ).toBe("Provider blocked the response") expect(yield* session.context(sessionID)).toMatchObject([ { type: "user" }, { @@ -3711,7 +3729,7 @@ describe("SessionRunnerLLM", () => { }), ) - it.effect("settles a local tool before one content-filter step failure", () => + scenarioIt("settles a local tool before one content-filter step failure", (scenario) => Effect.gen(function* () { yield* setup const session = yield* SessionV2.Service @@ -3720,20 +3738,25 @@ describe("SessionRunnerLLM", () => { prompt: PromptInput.Prompt.make({ text: "Tool before blocked response" }), resume: false, }) - toolExecutionGate = yield* Deferred.make() - toolExecutionsStarted = yield* Deferred.make() + const gate = (toolExecutionGate = yield* Deferred.make()) + const started = (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: "content-filter" }), - LLMEvent.finish({ reason: "content-filter" }), - ] - - const run = yield* session.resume(sessionID).pipe(Effect.forkChild) - yield* Deferred.await(toolExecutionsStarted) - yield* Deferred.succeed(toolExecutionGate, undefined) - expect((yield* Fiber.join(run).pipe(Effect.flip)).message).toBe("Provider blocked the response") + expect( + ( + yield* scenario + .run(function* () { + yield* (yield* scenario.llm.next()).respond.events( + LLMEvent.stepStart({ index: 0 }), + LLMEvent.toolCall({ id: "call-before-content-filter", name: "echo", input: { text: "settled" } }), + LLMEvent.stepFinish({ index: 0, reason: "content-filter" }), + LLMEvent.finish({ reason: "content-filter" }), + ) + yield* Deferred.await(started) + yield* Deferred.succeed(gate, undefined) + }) + .pipe(Effect.flip) + ).message, + ).toBe("Provider blocked the response") toolExecutionGate = undefined toolExecutionsStarted = undefined @@ -3751,7 +3774,7 @@ describe("SessionRunnerLLM", () => { }), ) - it.effect("does not recover context overflow after durable assistant output", () => + scenarioIt("does not recover context overflow after durable assistant output", (scenario) => Effect.gen(function* () { yield* setup const session = yield* SessionV2.Service @@ -3761,17 +3784,23 @@ describe("SessionRunnerLLM", () => { resume: false, }) - requests.length = 0 - response = [ - 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* scenario + .run(function* () { + yield* (yield* scenario.llm.next()).respond.events( + 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" }), + ) + }) + .pipe(Effect.flip) + ).message, + ).toBe("prompt too long") - expect(requests).toHaveLength(1) + expect(yield* scenario.llm.requests).toHaveLength(1) expect(yield* session.context(sessionID)).toMatchObject([ { type: "user", text: "Fail after output" }, { @@ -3810,23 +3839,23 @@ describe("SessionRunnerLLM", () => { }), ) - it.effect("retries eligible pre-output failures after exponential backoff", () => + scenarioIt("retries eligible pre-output failures after exponential backoff", (scenario) => Effect.gen(function* () { yield* setup const session = yield* SessionV2.Service yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Retry transport" }), resume: false }) - requests.length = 0 - responseStream = Stream.fail(providerUnavailable()) - response = fragmentFixture("text", "retry-success", ["Recovered"]).completeEvents - const run = yield* session.resume(sessionID).pipe(Effect.forkChild) - while (requests.length < 1) yield* Effect.yieldNow - yield* TestClock.adjust("1999 millis") - expect(requests).toHaveLength(1) - yield* TestClock.adjust("1 millis") - yield* Fiber.join(run) + yield* scenario.run(function* () { + yield* (yield* scenario.llm.next()).respond.fail(providerUnavailable()) + yield* TestClock.adjust("1999 millis") + expect(yield* scenario.llm.requests).toHaveLength(1) + yield* TestClock.adjust("1 millis") + yield* (yield* scenario.llm.next()).respond.events( + ...fragmentFixture("text", "retry-success", ["Recovered"]).completeEvents, + ) + }) - expect(requests).toHaveLength(2) + expect(yield* scenario.llm.requests).toHaveLength(2) const eventTypes = yield* recordedEventTypes(sessionID) expect(eventTypes).toContain("session.retry.scheduled.1") expect(eventTypes.filter((type) => type === "session.step.started.1")).toHaveLength(2) @@ -3839,41 +3868,44 @@ describe("SessionRunnerLLM", () => { }), ) - it.effect("uses a larger provider retry-after delay", () => + scenarioIt("uses a larger provider retry-after delay", (scenario) => Effect.gen(function* () { yield* setup const session = yield* SessionV2.Service yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Retry rate limit" }), resume: false }) - requests.length = 0 - responseStream = Stream.fail(rateLimited(5_000)) - response = fragmentFixture("text", "retry-after-success", ["Recovered"]).completeEvents - const run = yield* session.resume(sessionID).pipe(Effect.forkChild) - while (requests.length < 1) yield* Effect.yieldNow - yield* TestClock.adjust("4999 millis") - expect(requests).toHaveLength(1) - yield* TestClock.adjust("1 millis") - yield* Fiber.join(run) - expect(requests).toHaveLength(2) + yield* scenario.run(function* () { + yield* (yield* scenario.llm.next()).respond.fail(rateLimited(5_000)) + yield* TestClock.adjust("4999 millis") + expect(yield* scenario.llm.requests).toHaveLength(1) + yield* TestClock.adjust("1 millis") + yield* (yield* scenario.llm.next()).respond.events( + ...fragmentFixture("text", "retry-after-success", ["Recovered"]).completeEvents, + ) + }) + expect(yield* scenario.llm.requests).toHaveLength(2) }), ) - it.effect("stops after five total retry attempts", () => + scenarioIt("stops after five total retry attempts", (scenario) => Effect.gen(function* () { yield* setup const session = yield* SessionV2.Service yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Exhaust retries" }), resume: false }) - requests.length = 0 - streamFailure = providerUnavailable() + const failure = providerUnavailable() - const run = yield* session.resume(sessionID).pipe(Effect.forkChild) - while (requests.length < 1) yield* Effect.yieldNow - 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 - } - expect(yield* Fiber.join(run).pipe(Effect.flip)).toBe(streamFailure) - expect(requests).toHaveLength(5) + expect( + yield* scenario + .run(function* () { + for (const delay of [2_000, 4_000, 8_000, 16_000]) { + yield* (yield* scenario.llm.next()).respond.fail(failure) + yield* TestClock.adjust(delay) + } + yield* (yield* scenario.llm.next()).respond.fail(failure) + }) + .pipe(Effect.flip), + ).toBe(failure) + expect(yield* scenario.llm.requests).toHaveLength(5) const database = (yield* Database.Service).db const retries = yield* database @@ -3894,7 +3926,7 @@ describe("SessionRunnerLLM", () => { }), ) - it.effect("counts retry attempts against the agent step allowance", () => + scenarioIt("counts retry attempts against the agent step allowance", (scenario) => Effect.gen(function* () { yield* setup const agents = yield* AgentV2.Service @@ -3909,17 +3941,19 @@ describe("SessionRunnerLLM", () => { prompt: PromptInput.Prompt.make({ text: "Bound retries by steps" }), resume: false, }) - requests.length = 0 const failure = providerUnavailable() - responseStream = Stream.fail(failure) - streamFailure = failure - const run = yield* session.resume(sessionID).pipe(Effect.forkChild) - while (requests.length < 1) yield* Effect.yieldNow - yield* TestClock.adjust("2 seconds") - expect(yield* Fiber.join(run).pipe(Effect.flip)).toBe(failure) + expect( + yield* scenario + .run(function* () { + yield* (yield* scenario.llm.next()).respond.fail(failure) + yield* TestClock.adjust("2 seconds") + yield* (yield* scenario.llm.next()).respond.fail(failure) + }) + .pipe(Effect.flip), + ).toBe(failure) - expect(requests).toHaveLength(2) + expect(yield* scenario.llm.requests).toHaveLength(2) const eventTypes = yield* recordedEventTypes(sessionID) expect(eventTypes.filter((type) => type === "session.step.started.1")).toHaveLength(2) expect(eventTypes.filter((type) => type === "session.retry.scheduled.1")).toHaveLength(1) @@ -3927,22 +3961,26 @@ describe("SessionRunnerLLM", () => { }), ) - it.effect("does not retry non-eligible provider failures", () => + scenarioIt("does not retry non-eligible provider failures", (scenario) => Effect.gen(function* () { yield* setup const session = yield* SessionV2.Service yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Do not retry" }), resume: false }) - requests.length = 0 const failure = invalidRequest() - streamFailure = failure - expect(yield* session.resume(sessionID).pipe(Effect.flip)).toBe(failure) - expect(requests).toHaveLength(1) + expect( + yield* scenario + .run(function* () { + yield* (yield* scenario.llm.next()).respond.fail(failure) + }) + .pipe(Effect.flip), + ).toBe(failure) + expect(yield* scenario.llm.requests).toHaveLength(1) expect(yield* recordedEventTypes(sessionID)).not.toContain("session.retry.scheduled.1") }), ) - it.effect("does not continue automatically after a provider error follows a local tool call", () => + scenarioIt("does not continue automatically after a provider error follows a local tool call", (scenario) => Effect.gen(function* () { yield* setup const session = yield* SessionV2.Service @@ -3952,25 +3990,29 @@ describe("SessionRunnerLLM", () => { resume: false, }) - requests.length = 0 const executionCount = executions.length - toolExecutionGate = yield* Deferred.make() - toolExecutionsStarted = yield* Deferred.make() + const gate = (toolExecutionGate = yield* Deferred.make()) + const started = (toolExecutionsStarted = yield* Deferred.make()) toolExecutionsReady = 1 - response = [ - 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) - expect((yield* Fiber.join(run).pipe(Effect.flip)).message).toBe("Provider unavailable") + expect( + ( + yield* scenario + .run(function* () { + yield* (yield* scenario.llm.next()).respond.events( + LLMEvent.stepStart({ index: 0 }), + LLMEvent.toolCall({ id: "call-before-provider-error", name: "echo", input: { text: "settled" } }), + LLMEvent.providerError({ message: "Provider unavailable" }), + ) + yield* Deferred.await(started) + yield* Deferred.succeed(gate, undefined) + }) + .pipe(Effect.flip) + ).message, + ).toBe("Provider unavailable") toolExecutionGate = undefined toolExecutionsStarted = undefined - expect(requests).toHaveLength(1) + expect(yield* scenario.llm.requests).toHaveLength(1) expect(executions.slice(executionCount)).toEqual(["settled"]) const context = yield* session.context(sessionID) const assistant = requireAssistant(context) @@ -3983,7 +4025,7 @@ describe("SessionRunnerLLM", () => { }), ) - it.effect("durably fails a hosted tool when its provider errors before returning a result", () => + scenarioIt("durably fails a hosted tool when its provider errors before returning a result", (scenario) => Effect.gen(function* () { yield* setup const session = yield* SessionV2.Service @@ -3993,21 +4035,26 @@ describe("SessionRunnerLLM", () => { resume: false, }) - requests.length = 0 - response = [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.toolCall({ - id: "call-hosted-provider-error", - name: "web_search", - input: { query: "effect" }, - providerExecuted: true, - }), - LLMEvent.providerError({ message: "Provider unavailable" }), - ] + expect( + ( + yield* scenario + .run(function* () { + yield* (yield* scenario.llm.next()).respond.events( + LLMEvent.stepStart({ index: 0 }), + LLMEvent.toolCall({ + id: "call-hosted-provider-error", + name: "web_search", + input: { query: "effect" }, + providerExecuted: true, + }), + LLMEvent.providerError({ message: "Provider unavailable" }), + ) + }) + .pipe(Effect.flip) + ).message, + ).toBe("Provider unavailable") - expect((yield* session.resume(sessionID).pipe(Effect.flip)).message).toBe("Provider unavailable") - - expect(requests).toHaveLength(1) + expect(yield* scenario.llm.requests).toHaveLength(1) const context = yield* session.context(sessionID) expect(context).toMatchObject([ { type: "user", text: "Fail hosted tool durably" }, @@ -4026,7 +4073,7 @@ describe("SessionRunnerLLM", () => { }), ) - it.effect("preserves a tool defect before provider failure settlement", () => + scenarioIt("preserves a tool defect before provider failure settlement", (scenario) => Effect.gen(function* () { yield* setup const session = yield* SessionV2.Service @@ -4035,13 +4082,20 @@ describe("SessionRunnerLLM", () => { prompt: PromptInput.Prompt.make({ text: "Defect while provider fails" }), resume: false, }) - response = [ - 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* scenario + .run(function* () { + yield* (yield* scenario.llm.next()).respond.events( + LLMEvent.stepStart({ index: 0 }), + LLMEvent.toolCall({ id: "call-defect-provider-error", name: "defect", input: {} }), + LLMEvent.providerError({ message: "Provider unavailable" }), + ) + }) + .pipe(Effect.flip) + ).message, + ).toBe("Provider unavailable") const context = yield* session.context(sessionID) const assistant = requireAssistant(context) @@ -4056,7 +4110,7 @@ describe("SessionRunnerLLM", () => { }), ) - it.effect("durably fails a hosted tool left unresolved at normal provider EOF", () => + scenarioIt("durably fails a hosted tool left unresolved at normal provider EOF", (scenario) => Effect.gen(function* () { yield* setup const session = yield* SessionV2.Service @@ -4065,17 +4119,24 @@ describe("SessionRunnerLLM", () => { prompt: PromptInput.Prompt.make({ text: "Fail hosted tool at EOF" }), resume: false, }) - response = [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.toolCall({ - id: "call-hosted-eof", - name: "web_search", - input: { query: "effect" }, - providerExecuted: true, - }), - ] - expect((yield* session.resume(sessionID).pipe(Effect.flip)).message).toBe("Provider did not return a tool result") + expect( + ( + yield* scenario + .run(function* () { + yield* (yield* scenario.llm.next()).respond.events( + LLMEvent.stepStart({ index: 0 }), + LLMEvent.toolCall({ + id: "call-hosted-eof", + name: "web_search", + input: { query: "effect" }, + providerExecuted: true, + }), + ) + }) + .pipe(Effect.flip) + ).message, + ).toBe("Provider did not return a tool result") const assistant = requireAssistant(yield* session.context(sessionID)) const events = yield* recordedStepSettlementEvents(sessionID, assistant.id) expect(events.map((event) => event.type)).toEqual([ @@ -4101,7 +4162,7 @@ describe("SessionRunnerLLM", () => { }), ) - it.effect("fails an unresolved hosted tool before one clean step end", () => + scenarioIt("fails an unresolved hosted tool before one clean step end", (scenario) => Effect.gen(function* () { yield* setup const session = yield* SessionV2.Service @@ -4110,19 +4171,20 @@ describe("SessionRunnerLLM", () => { prompt: PromptInput.Prompt.make({ text: "Settle hosted tool before ending" }), resume: false, }) - response = [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.toolCall({ - id: "call-hosted-clean-end", - name: "web_search", - input: { query: "effect" }, - providerExecuted: true, - }), - LLMEvent.stepFinish({ index: 0, reason: "stop" }), - LLMEvent.finish({ reason: "stop" }), - ] - yield* session.resume(sessionID) + yield* scenario.run(function* () { + yield* (yield* scenario.llm.next()).respond.events( + LLMEvent.stepStart({ index: 0 }), + LLMEvent.toolCall({ + id: "call-hosted-clean-end", + name: "web_search", + input: { query: "effect" }, + providerExecuted: true, + }), + LLMEvent.stepFinish({ index: 0, reason: "stop" }), + LLMEvent.finish({ reason: "stop" }), + ) + }) const assistant = requireAssistant(yield* session.context(sessionID)) const events = yield* recordedStepSettlementEvents(sessionID, assistant.id) @@ -4138,7 +4200,7 @@ describe("SessionRunnerLLM", () => { }), ) - it.effect("settles unresolved local and hosted tools before one raw provider failure", () => + scenarioIt("settles unresolved local and hosted tools before one raw provider failure", (scenario) => Effect.gen(function* () { yield* setup const session = yield* SessionV2.Service @@ -4149,25 +4211,33 @@ describe("SessionRunnerLLM", () => { }) 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: {} }), - LLMEvent.toolCall({ - id: "call-hosted-raw-failure-pair", - name: "web_search", - input: { query: "effect" }, - providerExecuted: true, - }), - ]), - Stream.fromEffect(Deferred.succeed(providerFailed, undefined)).pipe(Stream.flatMap(() => Stream.fail(failure))), - ) + const gate = (toolExecutionGate = yield* Deferred.make()) - const run = yield* session.resume(sessionID).pipe(Effect.forkChild) - yield* Deferred.await(providerFailed) - yield* Deferred.succeed(toolExecutionGate, undefined) - expect(yield* Fiber.join(run).pipe(Effect.flip)).toBe(failure) + expect( + yield* scenario + .run(function* () { + yield* (yield* scenario.llm.next()).respond.stream( + Stream.concat( + Stream.fromIterable([ + LLMEvent.stepStart({ index: 0 }), + LLMEvent.toolCall({ id: "call-local-raw-failure", name: "defect", input: {} }), + LLMEvent.toolCall({ + id: "call-hosted-raw-failure-pair", + name: "web_search", + input: { query: "effect" }, + providerExecuted: true, + }), + ]), + Stream.fromEffect(Deferred.succeed(providerFailed, undefined)).pipe( + Stream.flatMap(() => Stream.fail(failure)), + ), + ), + ) + yield* Deferred.await(providerFailed) + yield* Deferred.succeed(gate, undefined) + }) + .pipe(Effect.flip), + ).toBe(failure) toolExecutionGate = undefined const assistant = requireAssistant(yield* session.context(sessionID)) @@ -4186,7 +4256,7 @@ describe("SessionRunnerLLM", () => { }), ) - it.effect("durably fails a hosted tool left unresolved by a raw provider stream failure", () => + scenarioIt("durably fails a hosted tool left unresolved by a raw provider stream failure", (scenario) => Effect.gen(function* () { yield* setup const session = yield* SessionV2.Service @@ -4195,23 +4265,29 @@ describe("SessionRunnerLLM", () => { prompt: PromptInput.Prompt.make({ text: "Fail hosted tool on raw failure" }), resume: false, }) - requests.length = 0 const failure = providerUnavailable() - responseStream = Stream.concat( - Stream.fromIterable([ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.toolCall({ - id: "call-hosted-raw-failure", - name: "web_search", - input: { query: "effect" }, - providerExecuted: true, - }), - ]), - Stream.fail(failure), - ) - expect(yield* session.resume(sessionID).pipe(Effect.flip)).toBe(failure) - expect(requests).toHaveLength(1) + expect( + yield* scenario + .run(function* () { + yield* (yield* scenario.llm.next()).respond.stream( + Stream.concat( + Stream.fromIterable([ + LLMEvent.stepStart({ index: 0 }), + LLMEvent.toolCall({ + id: "call-hosted-raw-failure", + name: "web_search", + input: { query: "effect" }, + providerExecuted: true, + }), + ]), + Stream.fail(failure), + ), + ) + }) + .pipe(Effect.flip), + ).toBe(failure) + expect(yield* scenario.llm.requests).toHaveLength(1) const assistant = requireAssistant(yield* session.context(sessionID)) const events = yield* recordedStepSettlementEvents(sessionID, assistant.id) expect(events.map((event) => event.type)).toEqual([ @@ -4236,50 +4312,46 @@ describe("SessionRunnerLLM", () => { }), ) - it.effect("rejects a second text start before the open fragment ends", () => + scenarioIt("rejects a second text start before the open fragment ends", (scenario) => Effect.gen(function* () { yield* setup const session = yield* SessionV2.Service yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Two blocks" }), resume: false }) - responses = undefined - streamGate = undefined - streamStarted = undefined - response = [ - 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* scenario + .run(function* () { + yield* (yield* scenario.llm.next()).respond.events( + LLMEvent.stepStart({ index: 0 }), + LLMEvent.textStart({ id: "text-1" }), + LLMEvent.textStart({ id: "text-2" }), + ) + }) + .pipe(Effect.catchDefect(Effect.succeed)) expect(defect).toBeInstanceOf(Error) if (!(defect instanceof Error)) return expect(defect.message).toBe("text start before end: text-2") }), ) - it.effect("projects sequential text fragments as separate content parts", () => + scenarioIt("projects sequential text fragments as separate content parts", (scenario) => Effect.gen(function* () { yield* setup const session = yield* SessionV2.Service yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Two blocks" }), resume: false }) - responses = undefined - streamGate = undefined - streamStarted = undefined - 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: "stop" }), - LLMEvent.finish({ reason: "stop" }), - ] - - yield* session.resume(sessionID) + yield* scenario.run(function* () { + yield* (yield* scenario.llm.next()).respond.events( + 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: "stop" }), + LLMEvent.finish({ reason: "stop" }), + ) + }) expect(yield* session.context(sessionID)).toMatchObject([ { type: "user", text: "Two blocks" }, @@ -4295,34 +4367,128 @@ describe("SessionRunnerLLM", () => { ) for (const kind of fragmentKinds) { - it.effect(`broadcasts provider ${kind} deltas without storing projection rewrites`, () => - verifyEphemeralDeltas(kind), + scenarioIt(`broadcasts provider ${kind} deltas without storing projection rewrites`, (scenario) => + Effect.gen(function* () { + yield* setup + const session = yield* SessionV2.Service + const prompt = `Stream ${kind}` + const chunks = Array.from({ length: 32 }, (_, index) => `${index},`) + const fixture = fragmentFixture(kind, fragmentID(kind, "many"), chunks) + const expectedContext = [{ type: "user", text: prompt }, fixture.expectedAssistant] + yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: prompt }), resume: false }) + const events = yield* EventV2.Service + const live = yield* events.subscribe(fixture.delta).pipe(Stream.take(32), Stream.runCollect, Effect.forkScoped) + yield* Effect.yieldNow + + yield* scenario.run(function* () { + yield* (yield* scenario.llm.next()).respond.events(...fixture.completeEvents) + }) + + const { db } = yield* Database.Service + const deltas = yield* db + .select({ type: EventTable.type }) + .from(EventTable) + .where(eq(EventTable.type, EventV2.versionedType(fixture.delta.type, 1))) + .all() + .pipe(Effect.orDie) + expect(Array.from(yield* Fiber.join(live))).toHaveLength(32) + expect(deltas).toHaveLength(0) + expect(yield* session.context(sessionID)).toMatchObject(expectedContext) + + yield* replaySessionProjection(sessionID) + + expect(yield* session.context(sessionID)).toMatchObject(expectedContext) + }), ) - it.effect(`durably closes partial ${kind} when the provider stream fails`, () => verifyPartialFlushOnFailure(kind)) + scenarioIt(`durably closes partial ${kind} when the provider stream fails`, (scenario) => + Effect.gen(function* () { + yield* setup + const session = yield* SessionV2.Service + const prompt = `Fail after ${kind}` + const fixture = fragmentFixture(kind, fragmentID(kind, "partial"), ["Partial"]) + const failure = providerUnavailable() + yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: prompt }), resume: false }) - it.effect(`durably closes partial ${kind} when the provider stream is interrupted`, () => - verifyPartialFlushOnInterruption(kind), + expect( + yield* scenario + .run(function* () { + yield* (yield* scenario.llm.next()).respond.stream( + Stream.concat(Stream.fromIterable(fixture.partialEvents), Stream.fail(failure)), + ) + }) + .pipe(Effect.flip), + ).toBe(failure) + expect(yield* session.context(sessionID)).toMatchObject([ + { type: "user", text: prompt }, + { + type: "assistant", + finish: "error", + error: { type: "provider.transport", message: "Provider unavailable" }, + content: [fixture.expectedContent], + }, + ]) + expect(yield* scenario.llm.requests).toHaveLength(1) + }), + ) + + scenarioIt(`durably closes partial ${kind} when the provider stream is interrupted`, (scenario) => + Effect.gen(function* () { + yield* setup + const session = yield* SessionV2.Service + const prompt = `Interrupt after ${kind}` + const fixture = fragmentFixture(kind, fragmentID(kind, "interrupted"), ["Partial"]) + const streamed = yield* Deferred.make() + yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: prompt }), resume: false }) + + yield* scenario + .run(function* () { + yield* (yield* scenario.llm.next()).respond.stream( + Stream.concat( + Stream.fromIterable(fixture.partialEvents), + Stream.fromEffect(Deferred.succeed(streamed, undefined)).pipe(Stream.flatMap(() => Stream.never)), + ), + ) + yield* Deferred.await(streamed) + yield* session.interrupt(sessionID) + }) + .pipe(Effect.exit) + expect(yield* session.context(sessionID)).toMatchObject([ + { type: "user", text: prompt }, + { + type: "assistant", + finish: "error", + error: { type: "aborted", message: "Step interrupted" }, + content: [ + kind === "tool input" + ? { type: "tool", id: fragmentID(kind, "interrupted"), state: { status: "error" } } + : fixture.expectedContent, + ], + }, + ]) + }), ) } - it.effect("rejects duplicate streamed text starts", () => + scenarioIt("rejects duplicate streamed text starts", (scenario) => Effect.gen(function* () { yield* setup - const session = yield* SessionV2.Service - responses = undefined - streamGate = undefined - streamStarted = undefined - response = [LLMEvent.textStart({ id: "text-1" }), LLMEvent.textStart({ id: "text-1" })] - const defect = yield* session.resume(sessionID).pipe(Effect.catchDefect(Effect.succeed)) + const defect = yield* scenario + .run(function* () { + yield* (yield* scenario.llm.next()).respond.events( + LLMEvent.textStart({ id: "text-1" }), + LLMEvent.textStart({ id: "text-1" }), + ) + }) + .pipe(Effect.catchDefect(Effect.succeed)) expect(defect).toBeInstanceOf(Error) if (!(defect instanceof Error)) return expect(defect.message).toBe("Duplicate text start: text-1") }), ) - it.effect("transitions streamed raw tool input to parsed called input", () => + scenarioIt("transitions streamed raw tool input to parsed called input", (scenario) => Effect.gen(function* () { yield* setup const session = yield* SessionV2.Service @@ -4332,20 +4498,17 @@ describe("SessionRunnerLLM", () => { resume: false, }) - responses = undefined - streamGate = undefined - streamStarted = undefined - 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" }), - LLMEvent.toolCall({ id: "call-parsed", name: "web_search", input: { query: "hello" }, providerExecuted: true }), - LLMEvent.stepFinish({ index: 0, reason: "stop" }), - LLMEvent.finish({ reason: "stop" }), - ] - - yield* session.resume(sessionID) + yield* scenario.run(function* () { + yield* (yield* scenario.llm.next()).respond.events( + 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" }), + LLMEvent.toolCall({ id: "call-parsed", name: "web_search", input: { query: "hello" }, providerExecuted: true }), + LLMEvent.stepFinish({ index: 0, reason: "stop" }), + LLMEvent.finish({ reason: "stop" }), + ) + }) expect(yield* session.context(sessionID)).toMatchObject([ { type: "user", text: "Call provider tool" }, @@ -4357,16 +4520,17 @@ describe("SessionRunnerLLM", () => { }), ) - it.effect("rejects malformed streamed tool input ordering", () => + scenarioIt("rejects malformed streamed tool input ordering", (scenario) => Effect.gen(function* () { yield* setup - const session = yield* SessionV2.Service - responses = undefined - streamGate = undefined - streamStarted = undefined - response = [LLMEvent.toolInputDelta({ id: "call-1", name: "read", text: "{}" })] - const defect = yield* session.resume(sessionID).pipe(Effect.catchDefect(Effect.succeed)) + const defect = yield* scenario + .run(function* () { + yield* (yield* scenario.llm.next()).respond.events( + LLMEvent.toolInputDelta({ id: "call-1", name: "read", text: "{}" }), + ) + }) + .pipe(Effect.catchDefect(Effect.succeed)) expect(defect).toBeInstanceOf(Error) if (!(defect instanceof Error)) return expect(defect.message).toBe("Tool input delta before start: call-1")