test(core): migrate runner delivery scenarios

This commit is contained in:
Kit Langton 2026-07-07 11:56:56 -04:00
commit ad40cebbd3

View file

@ -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* () { Effect.gen(function* () {
yield* setup yield* setup
const session = yield* SessionV2.Service const session = yield* SessionV2.Service
yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Run once" }), resume: false }) yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Run once" }), resume: false })
requests.length = 0 yield* scenario.run(function* () {
responses = undefined const call = yield* scenario.llm.next()
response = [ const second = yield* session.resume(sessionID).pipe(Effect.forkChild)
LLMEvent.stepStart({ index: 0 }), yield* Effect.yieldNow
LLMEvent.textStart({ id: "text-once" }), expect(yield* scenario.llm.requests).toHaveLength(1)
LLMEvent.textDelta({ id: "text-once", text: "Once" }), yield* call.respond.text("Once", { id: "text-once" })
LLMEvent.textEnd({ id: "text-once" }), yield* Fiber.join(second)
LLMEvent.stepFinish({ index: 0, reason: "stop" }), })
LLMEvent.finish({ reason: "stop" }),
]
streamGate = yield* Deferred.make<void>()
streamStarted = yield* Deferred.make<void>()
const first = yield* session.resume(sessionID).pipe(Effect.forkChild) expect(yield* scenario.llm.requests).toHaveLength(1)
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* session.context(sessionID)).toMatchObject([ expect(yield* session.context(sessionID)).toMatchObject([
{ type: "user", text: "Run once" }, { type: "user", text: "Run once" },
{ type: "assistant", finish: "stop", content: [{ type: "text", text: "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* () { Effect.gen(function* () {
yield* setup yield* setup
const session = yield* SessionV2.Service const session = yield* SessionV2.Service
yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Start working" }), resume: false }) yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Start working" }), resume: false })
requests.length = 0 yield* scenario.run(function* () {
responses = [ const first = yield* scenario.llm.next()
[ expect(userTexts(first.request)).toEqual(["Start working"])
LLMEvent.stepStart({ index: 0 }), yield* session.prompt({
LLMEvent.toolCall({ id: "call-echo", name: "echo", input: { text: "hello" } }), sessionID,
LLMEvent.stepFinish({ index: 0, reason: "tool-calls" }), prompt: PromptInput.Prompt.make({ text: "Wait until continuation ends" }),
LLMEvent.finish({ reason: "tool-calls" }), delivery: "queue",
], })
[ yield* first.respond.toolCall("echo", { text: "hello" }, { id: "call-echo" })
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<void>()
streamStarted = yield* Deferred.make<void>()
const first = yield* session.resume(sessionID).pipe(Effect.forkChild) const second = yield* scenario.llm.next()
yield* Deferred.await(streamStarted) expect(userTexts(second.request)).toEqual(["Start working"])
yield* session.prompt({ yield* second.respond.stop()
sessionID,
prompt: PromptInput.Prompt.make({ text: "Wait until continuation ends" }), const third = yield* scenario.llm.next()
delivery: "queue", 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(yield* scenario.llm.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"])
}), }),
) )
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* () { Effect.gen(function* () {
yield* setup const interrupted = yield* Deferred.make<Exit.Exit<void, SessionRunner.RunError | SessionV2.NotFoundError>>()
const session = yield* SessionV2.Service const retry = yield* Deferred.make<void>()
const { db } = yield* Database.Service const scenario = yield* RunnerScenario.make(() =>
yield* session.prompt({ SessionV2.Service.use((session) =>
sessionID, Effect.gen(function* () {
prompt: PromptInput.Prompt.make({ text: "Interrupt current work" }), yield* Deferred.succeed(interrupted, yield* session.resume(sessionID).pipe(Effect.exit))
resume: false, 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 yield* scenario.run(function* () {
responses = [ const first = yield* scenario.llm.next()
[], expect(userTexts(first.request)).toEqual(["Interrupt current work"])
[ yield* session.prompt({
LLMEvent.stepStart({ index: 0 }), sessionID,
LLMEvent.stepFinish({ index: 0, reason: "stop" }), prompt: PromptInput.Prompt.make({ text: "Run after interrupt" }),
LLMEvent.finish({ reason: "stop" }), delivery: "queue",
], })
] yield* session.interrupt(sessionID)
streamGate = yield* Deferred.make<void>() expect(yield* Deferred.await(interrupted)).toMatchObject({ _tag: "Failure" })
streamStarted = yield* Deferred.make<void>() expect(yield* SessionInput.hasPending(db, sessionID, "queue")).toBe(true)
const run = yield* session.resume(sessionID).pipe(Effect.forkChild) yield* Deferred.succeed(retry, undefined)
yield* Deferred.await(streamStarted) const second = yield* scenario.llm.next()
yield* session.prompt({ expect(userTexts(second.request)).toEqual(["Interrupt current work", "Run after interrupt"])
sessionID, yield* second.respond.stop()
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
expect(requests).toHaveLength(2) expect(yield* scenario.llm.requests).toHaveLength(2)
expect(userTexts(requests[0]!)).toEqual(["Interrupt current work"]) }).pipe(Effect.provide(testLayerWith(scenario.llm.layer)))
expect(userTexts(requests[1]!)).toEqual(["Interrupt current work", "Run after interrupt"])
}), }),
) )
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* () { Effect.gen(function* () {
yield* setup const interrupted = yield* Deferred.make<Exit.Exit<void, SessionRunner.RunError | SessionV2.NotFoundError>>()
const session = yield* SessionV2.Service const retry = yield* Deferred.make<void>()
const { db } = yield* Database.Service const scenario = yield* RunnerScenario.make(() =>
yield* session.prompt({ SessionV2.Service.use((session) =>
sessionID, Effect.gen(function* () {
prompt: PromptInput.Prompt.make({ text: "Interrupt current work" }), yield* Deferred.succeed(interrupted, yield* session.resume(sessionID).pipe(Effect.exit))
resume: false, 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 yield* scenario.run(function* () {
responses = [ const first = yield* scenario.llm.next()
[], expect(userTexts(first.request)).toEqual(["Interrupt current work"])
[ yield* session.prompt({
LLMEvent.stepStart({ index: 0 }), sessionID,
LLMEvent.stepFinish({ index: 0, reason: "stop" }), prompt: PromptInput.Prompt.make({ text: "Steer after interrupt" }),
LLMEvent.finish({ reason: "stop" }), })
], yield* session.interrupt(sessionID)
] expect(yield* Deferred.await(interrupted)).toMatchObject({ _tag: "Failure" })
streamGate = yield* Deferred.make<void>() expect(yield* SessionInput.hasPending(db, sessionID, "steer")).toBe(true)
streamStarted = yield* Deferred.make<void>()
const run = yield* session.resume(sessionID).pipe(Effect.forkChild) yield* Deferred.succeed(retry, undefined)
yield* Deferred.await(streamStarted) const second = yield* scenario.llm.next()
yield* session.prompt({ expect(userTexts(second.request)).toEqual(["Interrupt current work", "Steer after interrupt"])
sessionID, yield* second.respond.stop()
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)
const resumed = yield* session.resume(sessionID).pipe(Effect.forkChild) expect(yield* scenario.llm.requests).toHaveLength(2)
while (requests.length < 2) yield* Effect.yieldNow }).pipe(Effect.provide(testLayerWith(scenario.llm.layer)))
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"])
}), }),
) )
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* () { Effect.gen(function* () {
yield* setup yield* setup
const session = yield* SessionV2.Service const session = yield* SessionV2.Service
yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Start working" }), resume: false }) yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Start working" }), resume: false })
requests.length = 0 yield* scenario.run(function* () {
responses = [ const first = yield* scenario.llm.next()
[ expect(userTexts(first.request)).toEqual(["Start working"])
LLMEvent.stepStart({ index: 0 }), yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Queue first" }), delivery: "queue" })
LLMEvent.stepFinish({ index: 0, reason: "stop" }), yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Queue second" }), delivery: "queue" })
LLMEvent.finish({ reason: "stop" }), yield* first.respond.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<void>()
streamStarted = yield* Deferred.make<void>()
const first = yield* session.resume(sessionID).pipe(Effect.forkChild) const second = yield* scenario.llm.next()
yield* Deferred.await(streamStarted) expect(userTexts(second.request)).toEqual(["Start working", "Queue first"])
yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Queue first" }), delivery: "queue" }) yield* second.respond.stop()
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
expect(requests).toHaveLength(3) const third = yield* scenario.llm.next()
expect(userTexts(requests[0]!)).toEqual(["Start working"]) expect(userTexts(third.request)).toEqual(["Start working", "Queue first", "Queue second"])
expect(userTexts(requests[1]!)).toEqual(["Start working", "Queue first"]) yield* third.respond.stop()
expect(userTexts(requests[2]!)).toEqual(["Start working", "Queue first", "Queue second"]) })
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* () { Effect.gen(function* () {
yield* setup yield* setup
const session = yield* SessionV2.Service const session = yield* SessionV2.Service
@ -2488,166 +2434,125 @@ describe("SessionRunnerLLM", () => {
resume: false, resume: false,
}) })
requests.length = 0 yield* scenario.run(function* () {
responses = [ const first = yield* scenario.llm.next()
[ expect(userTexts(first.request)).toEqual(["Start steering"])
LLMEvent.stepStart({ index: 0 }), yield* first.respond.stop()
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* session.resume(sessionID) const second = yield* scenario.llm.next()
expect(userTexts(second.request)).toEqual(["Start steering", "Queue for later"])
expect(requests).toHaveLength(2) yield* second.respond.stop()
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<void>()
const secondGate = yield* Deferred.make<void>()
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" }),
}) })
yield* Deferred.succeed(secondGate, undefined)
yield* Fiber.join(first)
streamGate = undefined
expect(requests).toHaveLength(4) expect(yield* scenario.llm.requests).toHaveLength(2)
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",
])
}), }),
) )
it.effect("coalesces multiple active steering prompts into one continuation turn", () => scenarioIt("promotes steers before the next queued input", (scenario) =>
Effect.gen(function* () { Effect.gen(function* () {
yield* setup yield* setup
const session = yield* SessionV2.Service const session = yield* SessionV2.Service
yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Start working" }), resume: false }) yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Start working" }), resume: false })
requests.length = 0 yield* scenario.run(function* () {
responses = [ const first = yield* scenario.llm.next()
[ expect(userTexts(first.request)).toEqual(["Start working"])
LLMEvent.stepStart({ index: 0 }), yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Queue first" }), delivery: "queue" })
LLMEvent.stepFinish({ index: 0, reason: "stop" }), yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Queue second" }), delivery: "queue" })
LLMEvent.finish({ reason: "stop" }), yield* first.respond.stop()
],
[
LLMEvent.stepStart({ index: 0 }),
LLMEvent.stepFinish({ index: 0, reason: "stop" }),
LLMEvent.finish({ reason: "stop" }),
],
]
streamGate = yield* Deferred.make<void>()
streamStarted = yield* Deferred.make<void>()
const first = yield* session.resume(sessionID).pipe(Effect.forkChild) const second = yield* scenario.llm.next()
yield* Deferred.await(streamStarted) expect(userTexts(second.request)).toEqual(["Start working", "Queue first"])
yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "First steer" }) }) yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Steer before next queued input" }) })
yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Second steer" }) }) yield* session.prompt({
yield* Deferred.succeed(streamGate, undefined) sessionID,
yield* Fiber.join(first) prompt: PromptInput.Prompt.make({ text: "Also steer before next queued input" }),
streamGate = undefined })
streamStarted = undefined yield* second.respond.stop()
yield* Effect.yieldNow
expect(requests).toHaveLength(2) const third = yield* scenario.llm.next()
expect(userTexts(requests[1]!)).toEqual(["Start working", "First steer", "Second steer"]) 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* (yield* SessionExecution.Service).wake(sessionID)
yield* Effect.yieldNow yield* session.wait(sessionID)
expect(requests).toHaveLength(2) 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* () { Effect.gen(function* () {
yield* setup const failed = yield* Deferred.make<Exit.Exit<void, SessionRunner.RunError | SessionV2.NotFoundError>>()
const session = yield* SessionV2.Service const scenario = yield* RunnerScenario.make(() =>
yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Start working" }), resume: false }) 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 yield* scenario.run(function* () {
responses = undefined const first = yield* scenario.llm.next()
response = [] yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Recover with this" }) })
streamFailure = invalidRequest() yield* first.respond.fail(failure)
streamGate = yield* Deferred.make<void>() const exit = yield* Deferred.await(failed)
streamStarted = yield* Deferred.make<void>() expect(Exit.isFailure(exit) && Cause.squash(exit.cause)).toBe(failure)
const first = yield* session.resume(sessionID).pipe(Effect.forkChild) const second = yield* scenario.llm.next()
yield* Deferred.await(streamStarted) expect(userTexts(second.request)).toEqual(["Start working", "Recover with this"])
yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Recover with this" }) }) yield* second.respond.stop()
yield* Deferred.succeed(streamGate, undefined) })
expect(yield* Fiber.join(first).pipe(Effect.flip)).toBe(streamFailure)
streamFailure = undefined expect(yield* scenario.llm.requests).toHaveLength(2)
streamGate = undefined }).pipe(Effect.provide(testLayerWith(scenario.llm.layer)))
streamStarted = undefined
yield* Effect.yieldNow
expect(requests).toHaveLength(2)
expect(userTexts(requests[1]!)).toEqual(["Start working", "Recover with this"])
}), }),
) )
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* () { Effect.gen(function* () {
yield* setup yield* setup
const session = yield* SessionV2.Service const session = yield* SessionV2.Service
@ -2684,12 +2589,13 @@ describe("SessionRunnerLLM", () => {
input: { text: "stale" }, input: { text: "stale" },
executed: false, executed: false,
}) })
requests.length = 0 yield* scenario.run(function* () {
response = [] const call = yield* scenario.llm.next()
yield* session.resume(sessionID) expect(call.request.messages.map((message) => message.role)).toEqual(["user", "assistant", "tool"])
yield* call.respond.events()
})
expect(requests).toHaveLength(1) expect(yield* scenario.llm.requests).toHaveLength(1)
expect(requests[0]?.messages.map((message) => message.role)).toEqual(["user", "assistant", "tool"])
expect(yield* session.context(sessionID)).toMatchObject([ expect(yield* session.context(sessionID)).toMatchObject([
{ type: "user", text: "Recover interrupted tool" }, { 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* () { Effect.gen(function* () {
yield* setup yield* setup
const session = yield* SessionV2.Service const session = yield* SessionV2.Service
@ -2747,25 +2653,26 @@ describe("SessionRunnerLLM", () => {
executed: true, executed: true,
state: { itemId: "call-hosted-interrupted" }, state: { itemId: "call-hosted-interrupted" },
}) })
requests.length = 0 yield* scenario.run(function* () {
response = [] const call = yield* scenario.llm.next()
yield* session.resume(sessionID) 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(yield* scenario.llm.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" } },
])
}), }),
) )
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* () { Effect.gen(function* () {
yield* setup yield* setup
const session = yield* SessionV2.Service const session = yield* SessionV2.Service
@ -2789,12 +2696,13 @@ describe("SessionRunnerLLM", () => {
callID: "call-pending-interrupted", callID: "call-pending-interrupted",
name: "echo", name: "echo",
}) })
requests.length = 0 yield* scenario.run(function* () {
response = [] const call = yield* scenario.llm.next()
yield* session.resume(sessionID) expect(call.request.messages.map((message) => message.role)).toEqual(["user", "assistant", "tool"])
yield* call.respond.events()
})
expect(requests).toHaveLength(1) expect(yield* scenario.llm.requests).toHaveLength(1)
expect(requests[0]?.messages.map((message) => message.role)).toEqual(["user", "assistant", "tool"])
expect(yield* session.context(sessionID)).toMatchObject([ expect(yield* session.context(sessionID)).toMatchObject([
{ type: "user", text: "Recover interrupted tool input" }, { type: "user", text: "Recover interrupted tool input" },
{ type: "assistant", content: [{ type: "tool", id: "call-pending-interrupted", state: { status: "error" } }] }, { 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* () { Effect.gen(function* () {
yield* setup const scenario = yield* RunnerScenario.make(() =>
const session = yield* SessionV2.Service Effect.gen(function* () {
yield* session.prompt({ const execution = yield* SessionExecution.Service
sessionID, yield* execution.wake(sessionID)
prompt: PromptInput.Prompt.make({ text: "Wait in queue" }), yield* execution.awaitIdle(sessionID)
delivery: "queue", }),
resume: false, )
}) 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* scenario.run(function* () {
yield* (yield* SessionExecution.Service).wake(sessionID) const call = yield* scenario.llm.next()
yield* Effect.yieldNow expect(userTexts(call.request)).toEqual(["Wait in queue"])
yield* call.respond.events()
})
expect(requests).toHaveLength(1) expect(yield* scenario.llm.requests).toHaveLength(1)
expect(userTexts(requests[0]!)).toEqual(["Wait in queue"]) }).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* () { Effect.gen(function* () {
yield* setup
const session = yield* SessionV2.Service
const events = yield* EventV2.Service
const defect = new Error("fail after prompt promotion") const defect = new Error("fail after prompt promotion")
let fail = true let fail = true
yield* events.project(SessionEvent.PromptPromoted, () => (fail ? Effect.die(defect) : Effect.void)) const rolledBack = yield* Deferred.make<unknown>()
yield* session.prompt({ const retry = yield* Deferred.make<void>()
sessionID, const scenario = yield* RunnerScenario.make(() =>
prompt: PromptInput.Prompt.make({ text: "Recover promoted input" }), Effect.gen(function* () {
resume: false, 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) yield* scenario.run(function* () {
fail = false expect(yield* Deferred.await(rolledBack)).toBe(defect)
requests.length = 0 fail = false
response = [ yield* Deferred.succeed(retry, undefined)
LLMEvent.stepStart({ index: 0 }), const call = yield* scenario.llm.next()
LLMEvent.stepFinish({ index: 0, reason: "stop" }), expect(userTexts(call.request)).toEqual(["Recover promoted input"])
LLMEvent.finish({ reason: "stop" }), yield* call.respond.stop()
] })
}).pipe(Effect.provide(testLayerWith(scenario.llm.layer)))
yield* (yield* SessionExecution.Service).wake(sessionID)
while (requests.length === 0) yield* Effect.yieldNow
expect(userTexts(requests[0]!)).toEqual(["Recover promoted input"])
}), }),
) )
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* () { Effect.gen(function* () {
yield* setup yield* setup
const session = yield* SessionV2.Service const session = yield* SessionV2.Service
@ -2868,11 +2798,13 @@ describe("SessionRunnerLLM", () => {
resume: false, resume: false,
}) })
requests.length = 0 yield* scenario.run(function* () {
yield* session.resume(sessionID) const call = yield* scenario.llm.next()
expect(userTexts(call.request)).toEqual(["Run committed promotion"])
yield* call.respond.events()
})
expect(requests).toHaveLength(1) expect(yield* scenario.llm.requests).toHaveLength(1)
expect(userTexts(requests[0]!)).toEqual(["Run committed promotion"])
}), }),
) )
@ -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* () { Effect.gen(function* () {
yield* setup
const longSessionID = SessionV2.ID.make(`ses_${"a".repeat(64)}`) const longSessionID = SessionV2.ID.make(`ses_${"a".repeat(64)}`)
const otherLongSessionID = SessionV2.ID.make(`ses_${"b".repeat(64)}`) const otherLongSessionID = SessionV2.ID.make(`ses_${"b".repeat(64)}`)
yield* insertSession(longSessionID) const scenario = yield* RunnerScenario.make(() =>
yield* insertSession(otherLongSessionID) SessionV2.Service.use((session) =>
const session = yield* SessionV2.Service session.resume(longSessionID).pipe(Effect.andThen(session.resume(otherLongSessionID))),
yield* session.prompt({ ),
sessionID: longSessionID, )
prompt: PromptInput.Prompt.make({ text: "Run long session" }), yield* Effect.gen(function* () {
resume: false, yield* setup
}) yield* insertSession(longSessionID)
yield* session.prompt({ yield* insertSession(otherLongSessionID)
sessionID: otherLongSessionID, const session = yield* SessionV2.Service
prompt: PromptInput.Prompt.make({ text: "Run other long session" }), yield* session.prompt({
resume: false, 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* scenario.run(function* () {
yield* session.resume(longSessionID) yield* (yield* scenario.llm.next()).respond.events()
yield* session.resume(otherLongSessionID) yield* (yield* scenario.llm.next()).respond.events()
})
const keys = requests.map((request) => request.providerOptions?.openai?.promptCacheKey) const keys = (yield* scenario.llm.requests).map(
expect(keys).toEqual([longSessionID.slice(4), otherLongSessionID.slice(4)]) (request) => request.providerOptions?.openai?.promptCacheKey,
expect(keys.every((key) => typeof key === "string" && key.length === 64)).toBe(true) )
expect(keys[0]).not.toBe(keys[1]) 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* () { Effect.gen(function* () {
yield* setup const join = yield* Deferred.make<void>()
const session = yield* SessionV2.Service const joined = yield* Deferred.make<
yield* session.prompt({ readonly [
sessionID, Exit.Exit<void, SessionRunner.RunError | SessionV2.NotFoundError>,
prompt: PromptInput.Prompt.make({ text: "Retry after failure" }), Exit.Exit<void, SessionRunner.RunError | SessionV2.NotFoundError>,
resume: false, ]
}) >()
const retry = yield* Deferred.make<void>()
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 yield* scenario.run(function* () {
responses = undefined const first = yield* scenario.llm.next()
response = [] yield* Deferred.succeed(join, undefined)
streamFailure = invalidRequest() yield* Effect.yieldNow
streamGate = yield* Deferred.make<void>() expect(yield* scenario.llm.requests).toHaveLength(1)
streamStarted = yield* Deferred.make<void>() 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.succeed(retry, undefined)
yield* Deferred.await(streamStarted) yield* (yield* scenario.llm.next()).respond.events()
const second = yield* session.resume(sessionID).pipe(Effect.forkChild) })
yield* Effect.yieldNow
expect(requests).toHaveLength(1) expect(yield* scenario.llm.requests).toHaveLength(2)
yield* Deferred.succeed(streamGate, undefined) }).pipe(Effect.provide(testLayerWith(scenario.llm.layer)))
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)
}), }),
) )