test(core): migrate runner provider scenarios

This commit is contained in:
Kit Langton 2026-07-07 12:06:13 -04:00
parent ad40cebbd3
commit ddfa605c7c

View file

@ -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<void>()
toolExecutionsStarted = yield* Deferred.make<void>()
const gate = (toolExecutionGate = yield* Deferred.make<void>())
const started = (toolExecutionsStarted = yield* Deferred.make<void>())
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<void>()
toolExecutionsStarted = yield* Deferred.make<void>()
const gate = (toolExecutionGate = yield* Deferred.make<void>())
const started = (toolExecutionsStarted = yield* Deferred.make<void>())
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<void>()
toolExecutionGate = yield* Deferred.make<void>()
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<void>())
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<void>()
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")