diff --git a/packages/core/test/session-runner.test.ts b/packages/core/test/session-runner.test.ts index b98c5ef106..da8d9539c7 100644 --- a/packages/core/test/session-runner.test.ts +++ b/packages/core/test/session-runner.test.ts @@ -103,6 +103,20 @@ const client = Layer.succeed( generate: () => Effect.die("unused"), }), ) +const reply = { + stop: () => [ + LLMEvent.stepStart({ index: 0 }), + LLMEvent.stepFinish({ index: 0, reason: "stop" }), + LLMEvent.finish({ reason: "stop" }), + ], + text: (text: string, id: string) => fragmentFixture("text", id, [text]).completeEvents, + tool: (id: string, name: string, input: unknown) => [ + LLMEvent.stepStart({ index: 0 }), + LLMEvent.toolCall({ id, name, input }), + LLMEvent.stepFinish({ index: 0, reason: "tool-calls" }), + LLMEvent.finish({ reason: "tool-calls" }), + ], +} const model = Model.make({ id: "fake-model", provider: "fake", route: OpenAIChat.route }) const defaultSystem = SessionRunnerSystemPrompt.provider(model) const replacementModel = Model.make({ id: "replacement", provider: "fake", route: OpenAIChat.route }) @@ -339,6 +353,8 @@ const it = testEffect( ) const sessionID = SessionV2.ID.make("ses_runner_test") const otherSessionID = SessionV2.ID.make("ses_runner_other") +const admit = (session: SessionV2.Interface, text: string) => + session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text }), resume: false }) const insertSession = (id: SessionV2.ID) => Effect.gen(function* () { @@ -360,6 +376,9 @@ const insertSession = (id: SessionV2.ID) => const setup = Effect.gen(function* () { const { db } = yield* Database.Service + requests.length = 0 + authorizations.length = 0 + executions.length = 0 response = [] systemBaseline = "Initial context" systemRemoved = false @@ -385,6 +404,7 @@ const setup = Effect.gen(function* () { .run() .pipe(Effect.orDie) yield* insertSession(sessionID) + return yield* SessionV2.Service }) const providerUnavailable = () => @@ -409,14 +429,9 @@ const rateLimited = (retryAfterMs?: number) => }) const setupOverflowRecovery = Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service - response = fragmentFixture("text", "text-earlier", ["Earlier answer"]).completeEvents - yield* session.prompt({ - sessionID, - prompt: PromptInput.Prompt.make({ text: "Earlier question ".repeat(700) }), - resume: false, - }) + const session = yield* setup + response = reply.text("Earlier answer", "text-earlier") + yield* admit(session, "Earlier question ".repeat(700)) yield* session.resume(sessionID) currentModel = recoveryModel requests.length = 0 @@ -467,6 +482,9 @@ const recordedStepSettlementEvents = (id: SessionV2.ID, assistantMessageID: Sess ) }) +const hostedCall = (id: string, query: string) => + LLMEvent.toolCall({ id, name: "web_search", input: { query }, providerExecuted: true }) + const requireAssistant = (messages: readonly SessionMessage.Message[]) => { const assistant = messages.find((message) => message.type === "assistant") if (!assistant) throw new Error("Assistant message missing") @@ -577,14 +595,12 @@ const fragmentFixture = (kind: FragmentKind, id: string, chunks: readonly string const verifyEphemeralDeltas = (kind: FragmentKind) => Effect.gen(function* () { - yield* setup - requests.length = 0 - const session = yield* SessionV2.Service + const session = yield* setup 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 }) + yield* admit(session, prompt) const events = yield* EventV2.Service const live = yield* events.subscribe(fixture.delta).pipe(Stream.take(32), Stream.runCollect, Effect.forkScoped) yield* Effect.yieldNow @@ -610,13 +626,11 @@ const verifyEphemeralDeltas = (kind: FragmentKind) => const verifyPartialFlushOnFailure = (kind: FragmentKind) => Effect.gen(function* () { - yield* setup - requests.length = 0 - const session = yield* SessionV2.Service + const session = yield* setup 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 }) + yield* admit(session, prompt) responseStream = Stream.concat(Stream.fromIterable(fixture.partialEvents), Stream.fail(failure)) expect(yield* session.resume(sessionID).pipe(Effect.flip)).toBe(failure) @@ -634,12 +648,11 @@ const verifyPartialFlushOnFailure = (kind: FragmentKind) => const verifyPartialFlushOnInterruption = (kind: FragmentKind) => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service + const session = yield* setup 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* admit(session, prompt) responseStream = Stream.concat( Stream.fromIterable(fixture.partialEvents), Stream.fromEffect(Deferred.succeed(streamed, undefined)).pipe(Stream.flatMap(() => Stream.never)), @@ -667,9 +680,8 @@ const verifyPartialFlushOnInterruption = (kind: FragmentKind) => describe("SessionRunnerLLM", () => { it.effect("advertises and executes a location registered tool", () => Effect.gen(function* () { - yield* setup + const session = yield* setup const registry = yield* ToolRegistry.Service - const session = yield* SessionV2.Service const contexts: Tool.Context[] = [] yield* registry.register({ location_context: Tool.make({ @@ -683,20 +695,8 @@ describe("SessionRunnerLLM", () => { }), }), }) - yield* session.prompt({ - sessionID, - prompt: PromptInput.Prompt.make({ text: "Use application context" }), - resume: false, - }) - responses = [ - [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.toolCall({ id: "call-location", name: "location_context", input: { query: "hello" } }), - LLMEvent.stepFinish({ index: 0, reason: "tool-calls" }), - LLMEvent.finish({ reason: "tool-calls" }), - ], - [], - ] + yield* admit(session, "Use application context") + responses = [reply.tool("call-location", "location_context", { query: "hello" }), []] yield* session.resume(sessionID) @@ -727,13 +727,7 @@ describe("SessionRunnerLLM", () => { it.effect("starts a real runner turn after default prompt recording", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service - requests.length = 0 - responses = undefined - streamGate = undefined - streamStarted = undefined - response = [] + const session = yield* setup const message = yield* session.prompt({ sessionID, @@ -750,16 +744,10 @@ describe("SessionRunnerLLM", () => { it.effect("streams one request with registry definitions from chronological V2 user history", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "First" }), resume: false }) - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Second" }), resume: false }) + const session = yield* setup + yield* admit(session, "First") + yield* admit(session, "Second") - requests.length = 0 - responses = undefined - streamGate = undefined - streamStarted = undefined - response = [] yield* session.resume(sessionID) expect(requests).toHaveLength(1) @@ -775,8 +763,7 @@ describe("SessionRunnerLLM", () => { it.effect("retries the first provider turn after system context becomes available", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service + const session = yield* setup const { db } = yield* Database.Service const messageID = SessionMessage.ID.create() systemUnavailable = true @@ -786,7 +773,6 @@ describe("SessionRunnerLLM", () => { prompt: PromptInput.Prompt.make({ text: "First" }), resume: false, }) - requests.length = 0 const exit = yield* session.resume(sessionID).pipe(Effect.exit) @@ -813,13 +799,10 @@ describe("SessionRunnerLLM", () => { it.effect("interrupts a source Location runner after a Session moves", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service + const session = yield* setup const events = yield* EventV2.Service const { db } = yield* Database.Service - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "First" }), resume: false }) - requests.length = 0 - response = [] + yield* admit(session, "First") yield* session.resume(sessionID) yield* events.publish(SessionEvent.Moved, { @@ -834,7 +817,7 @@ describe("SessionRunnerLLM", () => { .get(), ).toBeUndefined() - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Second" }), resume: false }) + yield* admit(session, "Second") const exit = yield* session.resume(sessionID).pipe(Effect.exit) expect(Exit.isFailure(exit) && Cause.hasInterruptsOnly(exit.cause)).toBe(true) @@ -845,11 +828,9 @@ describe("SessionRunnerLLM", () => { it.effect("copies the context checkpoint to a fork", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service + const session = yield* setup const { db } = yield* Database.Service - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "First" }), resume: false }) - response = [] + yield* admit(session, "First") yield* session.resume(sessionID) const forked = yield* session.fork({ sessionID }) @@ -874,11 +855,9 @@ describe("SessionRunnerLLM", () => { it.effect("heals an undecodable stored applied record by re-announcing context", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service + const session = yield* setup const { db } = yield* Database.Service - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "First" }), resume: false }) - response = [] + yield* admit(session, "First") yield* session.resume(sessionID) yield* db .update(InstructionCheckpointTable) @@ -886,7 +865,7 @@ describe("SessionRunnerLLM", () => { .where(eq(InstructionCheckpointTable.session_id, sessionID)) .run() .pipe(Effect.orDie) - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Second" }), resume: false }) + yield* admit(session, "Second") requests.length = 0 yield* session.resume(sessionID) @@ -908,15 +887,12 @@ describe("SessionRunnerLLM", () => { it.effect("reuses one durable baseline after the context producer changes", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "First" }), resume: false }) + const session = yield* setup + yield* admit(session, "First") - requests.length = 0 - response = [] yield* session.resume(sessionID) systemBaseline = "Changed context" - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Second" }), resume: false }) + yield* admit(session, "Second") yield* session.resume(sessionID) expect(requests.map((request) => request.system.map((part) => part.text))).toEqual([ @@ -942,13 +918,11 @@ describe("SessionRunnerLLM", () => { it.effect("uses the selected model family prompt when the agent does not override it", () => Effect.gen(function* () { - yield* setup + const session = yield* setup currentModel = Model.make({ id: "gpt-5", provider: "openai", route: OpenAIChat.route }) - const session = yield* SessionV2.Service - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "First" }), resume: false }) + yield* admit(session, "First") - requests.length = 0 - response = fragmentFixture("text", "text-provider-prompt", ["Done"]).completeEvents + response = reply.text("Done", "text-provider-prompt") yield* session.resume(sessionID) expect(requests.at(-1)?.system.map((part) => part.text)).toEqual([ @@ -960,7 +934,7 @@ describe("SessionRunnerLLM", () => { it.effect("uses the selected model family prompt when the agent system override is empty", () => Effect.gen(function* () { - yield* setup + const session = yield* setup currentModel = Model.make({ id: "gpt-5", provider: "openai", route: OpenAIChat.route }) const agent = yield* AgentV2.Service yield* agent.transform((editor) => @@ -969,11 +943,9 @@ describe("SessionRunnerLLM", () => { agent.mode = "primary" }), ) - const session = yield* SessionV2.Service - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "First" }), resume: false }) + yield* admit(session, "First") - requests.length = 0 - response = fragmentFixture("text", "text-empty-agent-system", ["Done"]).completeEvents + response = reply.text("Done", "text-empty-agent-system") yield* session.resume(sessionID) expect(requests.at(-1)?.system.map((part) => part.text)).toEqual([ @@ -985,7 +957,7 @@ describe("SessionRunnerLLM", () => { it.effect("includes the effective default agent system before durable context", () => Effect.gen(function* () { - yield* setup + const session = yield* setup const agent = yield* AgentV2.Service yield* agent.transform((editor) => editor.update(AgentV2.ID.make("build"), (agent) => { @@ -993,11 +965,9 @@ describe("SessionRunnerLLM", () => { agent.mode = "primary" }), ) - const session = yield* SessionV2.Service - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "First" }), resume: false }) + yield* admit(session, "First") - requests.length = 0 - response = fragmentFixture("text", "text-build", ["Done"]).completeEvents + response = reply.text("Done", "text-build") yield* session.resume(sessionID) expect(requests.at(-1)?.system.map((part) => part.text)).toEqual(["Build agent instructions", "Initial context"]) @@ -1006,7 +976,7 @@ describe("SessionRunnerLLM", () => { it.effect("uses the configured default agent system for omitted-agent sessions", () => Effect.gen(function* () { - yield* setup + const session = yield* setup const agent = yield* AgentV2.Service yield* agent.transform((editor) => { editor.update(AgentV2.ID.make("build"), (agent) => { @@ -1019,11 +989,9 @@ describe("SessionRunnerLLM", () => { }) editor.default(AgentV2.ID.make("reviewer")) }) - const session = yield* SessionV2.Service - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "First" }), resume: false }) + yield* admit(session, "First") - requests.length = 0 - response = fragmentFixture("text", "text-reviewer", ["Done"]).completeEvents + response = reply.text("Done", "text-reviewer") yield* session.resume(sessionID) expect(requests.at(-1)?.system.map((part) => part.text)).toEqual(["Reviewer instructions", "Initial context"]) @@ -1033,7 +1001,7 @@ describe("SessionRunnerLLM", () => { it.effect("uses only the agent prompt and durable baseline as system parts", () => Effect.gen(function* () { - yield* setup + const session = yield* setup const agent = yield* AgentV2.Service yield* agent.transform((editor) => editor.update(AgentV2.ID.make("build"), (agent) => { @@ -1041,11 +1009,9 @@ describe("SessionRunnerLLM", () => { agent.mode = "primary" }), ) - const session = yield* SessionV2.Service - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "First" }), resume: false }) + yield* admit(session, "First") - requests.length = 0 - response = fragmentFixture("text", "text-no-system", ["Done"]).completeEvents + response = reply.text("Done", "text-no-system") yield* session.resume(sessionID) expect(requests.at(-1)?.system.map((part) => part.text)).toEqual(["Build agent instructions", "Initial context"]) @@ -1054,7 +1020,7 @@ describe("SessionRunnerLLM", () => { it.effect("uses an explicitly selected non-build agent system", () => Effect.gen(function* () { - yield* setup + const session = yield* setup const { db } = yield* Database.Service const agent = yield* AgentV2.Service yield* agent.transform((editor) => @@ -1069,11 +1035,9 @@ describe("SessionRunnerLLM", () => { .where(eq(SessionTable.id, sessionID)) .run() .pipe(Effect.orDie) - const session = yield* SessionV2.Service - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "First" }), resume: false }) + yield* admit(session, "First") - requests.length = 0 - response = fragmentFixture("text", "text-selected", ["Done"]).completeEvents + response = reply.text("Done", "text-selected") yield* session.resume(sessionID) expect(requests.at(-1)?.system.map((part) => part.text)).toEqual(["Reviewer instructions", "Initial context"]) @@ -1083,21 +1047,18 @@ describe("SessionRunnerLLM", () => { it.effect("updates selected-agent skill guidance after an agent switch", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service + const session = yield* setup const events = yield* EventV2.Service skillBaselines.set(AgentV2.ID.make("build"), "Build skills") - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "First" }), resume: false }) + yield* admit(session, "First") - requests.length = 0 - response = [] yield* session.resume(sessionID) skillBaselines.set(AgentV2.ID.make("reviewer"), "Reviewer skills") yield* events.publish(SessionEvent.AgentSelected, { sessionID, agent: "reviewer", }) - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Second" }), resume: false }) + yield* admit(session, "Second") yield* session.resume(sessionID) expect(requests.map((request) => request.system.map((part) => part.text))).toEqual([ @@ -1110,8 +1071,7 @@ describe("SessionRunnerLLM", () => { it.effect("keeps the sampled agent when selection changes during observation", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service + const session = yield* setup const events = yield* EventV2.Service skillBaselines.set(AgentV2.ID.make("build"), "Build skills") skillBaselines.set(AgentV2.ID.make("reviewer"), "Reviewer skills") @@ -1126,10 +1086,8 @@ describe("SessionRunnerLLM", () => { }) .pipe(Effect.asVoid) }) - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "First" }), resume: false }) + yield* admit(session, "First") - requests.length = 0 - response = [] yield* session.resume(sessionID) expect(requests.map((request) => request.system.map((part) => part.text))).toEqual([ @@ -1140,8 +1098,7 @@ describe("SessionRunnerLLM", () => { it.effect("keeps the sampled model when selection changes during model resolution", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service + const session = yield* setup const events = yield* EventV2.Service let switched = false modelResolveHook = Effect.suspend(() => { @@ -1154,10 +1111,8 @@ describe("SessionRunnerLLM", () => { }) .pipe(Effect.asVoid) }) - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "First" }), resume: false }) + yield* admit(session, "First") - requests.length = 0 - response = [] yield* session.resume(sessionID) expect(requests.map((request) => request.model)).toEqual([model]) expect(requests.map((request) => request.system.map((part) => part.text))).toEqual([ @@ -1168,15 +1123,12 @@ describe("SessionRunnerLLM", () => { it.effect("admits removed context as a chronological System message", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "First" }), resume: false }) + const session = yield* setup + yield* admit(session, "First") - requests.length = 0 - response = [] yield* session.resume(sessionID) systemRemoved = true - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Second" }), resume: false }) + yield* admit(session, "Second") yield* session.resume(sessionID) expect(requests[1]?.messages.map((message) => message.role)).toEqual(["user", "system", "user"]) @@ -1189,14 +1141,11 @@ describe("SessionRunnerLLM", () => { it.effect("renders API context entries through the belief lifecycle", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service + const session = yield* setup const contextEntries = yield* InstructionEntry.Service yield* contextEntries.put({ sessionID, key: "deploy-target", value: "production" }) - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "First" }), resume: false }) + yield* admit(session, "First") - requests.length = 0 - response = [] yield* session.resume(sessionID) // String values render verbatim inside the tagged block at baseline. @@ -1207,7 +1156,7 @@ describe("SessionRunnerLLM", () => { // Non-string JSON pretty-prints; the change narrates as a System update. yield* contextEntries.put({ sessionID, key: "deploy-target", value: { region: "us-east-1" } }) - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Second" }), resume: false }) + yield* admit(session, "Second") yield* session.resume(sessionID) expect(requests[1]?.messages.map((message) => message.role)).toEqual(["user", "system", "user"]) @@ -1228,7 +1177,7 @@ describe("SessionRunnerLLM", () => { // Deleting the row announces removal through the stored removal text. yield* contextEntries.remove({ sessionID, key: "deploy-target" }) - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Third" }), resume: false }) + yield* admit(session, "Third") yield* session.resume(sessionID) expect(requests[2]?.messages.map((message) => message.role)).toEqual(["user", "system", "user", "system", "user"]) @@ -1241,23 +1190,20 @@ describe("SessionRunnerLLM", () => { it.effect("keeps the baseline and chronological System updates after a model switch", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service + const session = yield* setup const events = yield* EventV2.Service - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "First" }), resume: false }) + yield* admit(session, "First") - requests.length = 0 - response = [] yield* session.resume(sessionID) systemBaseline = "Changed context" - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Second" }), resume: false }) + yield* admit(session, "Second") yield* session.resume(sessionID) yield* events.publish(SessionEvent.ModelSelected, { sessionID, model: { id: ModelV2.ID.make("replacement"), providerID: ProviderV2.ID.make("fake") }, }) systemBaseline = "Replacement context" - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Third" }), resume: false }) + yield* admit(session, "Third") yield* session.resume(sessionID) expect(requests.map((request) => request.system.map((part) => part.text))).toEqual([ @@ -1277,31 +1223,28 @@ describe("SessionRunnerLLM", () => { ]) yield* replaySessionProjection(sessionID) expect(yield* session.messages({ sessionID })).toHaveLength(6) - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Fourth" }), resume: false }) + yield* admit(session, "Fourth") yield* session.resume(sessionID) }), ) it.effect("preserves the baseline while context is temporarily unavailable", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service + const session = yield* setup const events = yield* EventV2.Service - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "First" }), resume: false }) + yield* admit(session, "First") - requests.length = 0 - response = [] yield* session.resume(sessionID) yield* events.publish(SessionEvent.ModelSelected, { sessionID, model: { id: ModelV2.ID.make("replacement"), providerID: ProviderV2.ID.make("fake") }, }) systemUnavailable = true - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Second" }), resume: false }) + yield* admit(session, "Second") yield* session.resume(sessionID) systemUnavailable = false systemBaseline = "Replacement context" - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Third" }), resume: false }) + yield* admit(session, "Third") yield* session.resume(sessionID) expect(requests.map((request) => request.system.map((part) => part.text))).toEqual([ @@ -1314,13 +1257,10 @@ describe("SessionRunnerLLM", () => { it.effect("rebuilds the baseline directly after completed compaction", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service + const session = yield* setup const events = yield* EventV2.Service - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "First" }), resume: false }) + yield* admit(session, "First") - requests.length = 0 - response = [] yield* session.resume(sessionID) yield* events.publish(SessionEvent.Compaction.Started, { sessionID, @@ -1333,7 +1273,7 @@ describe("SessionRunnerLLM", () => { recent: "", }) systemBaseline = "Replacement context" - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Second" }), resume: false }) + yield* admit(session, "Second") yield* session.resume(sessionID) expect(requests.map((request) => request.system.map((part) => part.text))).toEqual([ @@ -1341,26 +1281,24 @@ describe("SessionRunnerLLM", () => { [defaultSystem, "Replacement context"], ]) yield* replaySessionProjection(sessionID) - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Third" }), resume: false }) + yield* admit(session, "Third") yield* session.resume(sessionID) }), ) it.effect("runs one durable compaction barrier before later steer and queued prompts", () => Effect.gen(function* () { - yield* setup - requests.length = 0 + const session = yield* setup currentModel = recoveryModel - const session = yield* SessionV2.Service streamGate = yield* Deferred.make() streamStarted = yield* Deferred.make() responses = [ - fragmentFixture("text", "text-active", ["Active complete"]).completeEvents, + reply.text("Active complete", "text-active"), [LLMEvent.textDelta({ id: "summary", text: "durable summary" })], - fragmentFixture("text", "text-steer", ["Steer complete"]).completeEvents, - fragmentFixture("text", "text-queue", ["Queue complete"]).completeEvents, + reply.text("Steer complete", "text-steer"), + reply.text("Queue complete", "text-queue"), ] - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Active work" }), resume: false }) + yield* admit(session, "Active work") const active = yield* session.resume(sessionID).pipe(Effect.forkChild) yield* Deferred.await(streamStarted) @@ -1375,11 +1313,7 @@ describe("SessionRunnerLLM", () => { status: "queued", }) - yield* session.prompt({ - sessionID, - prompt: PromptInput.Prompt.make({ text: "Steer after compaction" }), - resume: false, - }) + yield* admit(session, "Steer after compaction") yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Queue after compaction" }), @@ -1406,18 +1340,16 @@ describe("SessionRunnerLLM", () => { it.effect("releases queued prompts when durable compaction fails", () => Effect.gen(function* () { - yield* setup - requests.length = 0 + const session = yield* setup currentModel = recoveryModel - const session = yield* SessionV2.Service streamGate = yield* Deferred.make() streamStarted = yield* Deferred.make() responses = [ - fragmentFixture("text", "text-active-failure", ["Active complete"]).completeEvents, + reply.text("Active complete", "text-active-failure"), [], - fragmentFixture("text", "text-after-failure", ["Continued"]).completeEvents, + reply.text("Continued", "text-after-failure"), ] - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Active work" }), resume: false }) + yield* admit(session, "Active work") const active = yield* session.resume(sessionID).pipe(Effect.forkChild) yield* Deferred.await(streamStarted) @@ -1443,27 +1375,18 @@ describe("SessionRunnerLLM", () => { it.effect("automatically compacts into a completed summary and retained recent turn", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service - response = fragmentFixture("text", "text-first", ["Earlier answer"]).completeEvents - yield* session.prompt({ - sessionID, - prompt: PromptInput.Prompt.make({ text: "Earlier question ".repeat(180) }), - resume: false, - }) + const session = yield* setup + response = reply.text("Earlier answer", "text-first") + yield* admit(session, "Earlier question ".repeat(180)) yield* session.resume(sessionID) currentModel = compactModel requests.length = 0 responses = [ - fragmentFixture("text", "text-summary", ["## Objective\n- Preserve the task"]).completeEvents, - fragmentFixture("text", "text-final", ["Continued"]).completeEvents, + reply.text("## Objective\n- Preserve the task", "text-summary"), + reply.text("Continued", "text-final"), ] - yield* session.prompt({ - sessionID, - prompt: PromptInput.Prompt.make({ text: "Recent exact request ".repeat(180) }), - resume: false, - }) + yield* admit(session, "Recent exact request ".repeat(180)) yield* session.resume(sessionID) expect(requests).toHaveLength(2) @@ -1482,14 +1405,10 @@ describe("SessionRunnerLLM", () => { requests.length = 0 executions.length = 0 responses = [ - fragmentFixture("text", "text-summary-2", ["## Objective\n- Preserve the updated task"]).completeEvents, - fragmentFixture("text", "text-final-2", ["Continued again"]).completeEvents, + reply.text("## Objective\n- Preserve the updated task", "text-summary-2"), + reply.text("Continued again", "text-final-2"), ] - yield* session.prompt({ - sessionID, - prompt: PromptInput.Prompt.make({ text: "Newest exact request ".repeat(180) }), - resume: false, - }) + yield* admit(session, "Newest exact request ".repeat(180)) yield* session.resume(sessionID) expect(requests).toHaveLength(2) @@ -1512,10 +1431,10 @@ describe("SessionRunnerLLM", () => { LLMEvent.stepStart({ index: 0 }), LLMEvent.providerError({ message: "prompt too long", classification: "context-overflow" }), ], - fragmentFixture("text", "text-summary", ["## Objective\n- Recover overflow"]).completeEvents, - fragmentFixture("text", "text-final", ["Recovered"]).completeEvents, + reply.text("## Objective\n- Recover overflow", "text-summary"), + reply.text("Recovered", "text-final"), ] - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Continue" }), resume: false }) + yield* admit(session, "Continue") yield* session.resume(sessionID) expect(requests).toHaveLength(3) @@ -1540,12 +1459,8 @@ describe("SessionRunnerLLM", () => { LLMEvent.stepStart({ index: 0 }), LLMEvent.providerError({ message: "prompt too long", classification: "context-overflow" }), ] - responses = [ - overflow(), - fragmentFixture("text", "text-summary", ["## Objective\n- Recover once"]).completeEvents, - overflow(), - ] - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Continue" }), resume: false }) + responses = [overflow(), reply.text("## Objective\n- Recover once", "text-summary"), overflow()] + yield* admit(session, "Continue") expect((yield* session.resume(sessionID).pipe(Effect.flip)).message).toBe("prompt too long") expect(requests).toHaveLength(3) @@ -1570,10 +1485,10 @@ describe("SessionRunnerLLM", () => { }), ) responses = [ - fragmentFixture("text", "text-summary", ["## Objective\n- Recover raw overflow"]).completeEvents, - fragmentFixture("text", "text-final", ["Recovered"]).completeEvents, + reply.text("## Objective\n- Recover raw overflow", "text-summary"), + reply.text("Recovered", "text-final"), ] - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Continue" }), resume: false }) + yield* admit(session, "Continue") yield* session.resume(sessionID) expect(requests).toHaveLength(3) @@ -1591,7 +1506,7 @@ describe("SessionRunnerLLM", () => { [LLMEvent.providerError({ message: "prompt too long", classification: "context-overflow" })], [LLMEvent.providerError({ message: "summary unavailable" })], ] - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Continue" }), resume: false }) + yield* admit(session, "Continue") expect((yield* session.resume(sessionID).pipe(Effect.flip)).message).toBe("prompt too long") expect(requests).toHaveLength(2) @@ -1609,12 +1524,12 @@ describe("SessionRunnerLLM", () => { const session = yield* setupOverflowRecovery responses = [ [LLMEvent.providerError({ message: "prompt too long", classification: "context-overflow" })], - fragmentFixture("text", "text-summary", ["## Objective\n- Interrupted"]).completeEvents, + reply.text("## Objective\n- Interrupted", "text-summary"), ] const firstGate = yield* Deferred.make() const summaryGate = yield* Deferred.make() streamGate = firstGate - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Continue" }), resume: false }) + yield* admit(session, "Continue") const run = yield* session.resume(sessionID).pipe(Effect.forkChild) while (requests.length < 1) yield* Effect.yieldNow streamGate = summaryGate @@ -1631,16 +1546,13 @@ describe("SessionRunnerLLM", () => { it.effect("rebaselines after compaction from the last-applied belief while unobservable", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service + const session = yield* setup const events = yield* EventV2.Service - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "First" }), resume: false }) + yield* admit(session, "First") - requests.length = 0 - response = [] yield* session.resume(sessionID) systemBaseline = "Changed context" - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Second" }), resume: false }) + yield* admit(session, "Second") yield* session.resume(sessionID) yield* events.publish(SessionEvent.Compaction.Started, { sessionID, @@ -1653,7 +1565,7 @@ describe("SessionRunnerLLM", () => { recent: "", }) systemUnavailable = true - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Third" }), resume: false }) + yield* admit(session, "Third") yield* session.resume(sessionID) // The rebaseline proceeds while the source is unobservable, restating the model's belief. @@ -1664,14 +1576,9 @@ describe("SessionRunnerLLM", () => { it.effect("projects reasoning and tool events without executing or continuing tools", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Use tools" }), resume: false }) + const session = yield* setup + yield* admit(session, "Use tools") - requests.length = 0 - responses = undefined - streamGate = undefined - streamStarted = undefined response = [ LLMEvent.stepStart({ index: 0 }), LLMEvent.reasoningStart({ id: "reasoning-1" }), @@ -1765,31 +1672,10 @@ describe("SessionRunnerLLM", () => { it.effect("continues with reloaded history after durably settling one local tool call", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Echo this" }), resume: false }) + const session = yield* setup + yield* admit(session, "Echo this") - requests.length = 0 - authorizations.length = 0 - executions.length = 0 - streamGate = undefined - streamStarted = undefined - responses = [ - [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.toolCall({ id: "call-echo", name: "echo", input: { text: "hello" } }), - LLMEvent.stepFinish({ index: 0, reason: "tool-calls" }), - LLMEvent.finish({ reason: "tool-calls" }), - ], - [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.textStart({ id: "text-final" }), - LLMEvent.textDelta({ id: "text-final", text: "Done" }), - LLMEvent.textEnd({ id: "text-final" }), - LLMEvent.stepFinish({ index: 0, reason: "stop" }), - LLMEvent.finish({ reason: "stop" }), - ], - ] + responses = [reply.tool("call-echo", "echo", { text: "hello" }), reply.text("Done", "text-final")] yield* session.resume(sessionID) @@ -1831,25 +1717,11 @@ describe("SessionRunnerLLM", () => { it.effect("reloads a model switch before a tool-driven continuation turn", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service + const session = yield* setup const events = yield* EventV2.Service - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Echo this" }), resume: false }) + yield* admit(session, "Echo this") - requests.length = 0 - responses = [ - [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.toolCall({ id: "call-echo", name: "echo", input: { text: "hello" } }), - LLMEvent.stepFinish({ index: 0, reason: "tool-calls" }), - LLMEvent.finish({ reason: "tool-calls" }), - ], - [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.stepFinish({ index: 0, reason: "stop" }), - LLMEvent.finish({ reason: "stop" }), - ], - ] + responses = [reply.tool("call-echo", "echo", { text: "hello" }), reply.stop()] toolExecutionGate = yield* Deferred.make() toolExecutionsStarted = yield* Deferred.make() toolExecutionsReady = 1 @@ -1874,11 +1746,9 @@ describe("SessionRunnerLLM", () => { it.effect("restores durable reasoning provider metadata in a second-turn request", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Think first" }), resume: false }) + const session = yield* setup + yield* admit(session, "Think first") - requests.length = 0 response = [ LLMEvent.stepStart({ index: 0 }), LLMEvent.reasoningStart({ id: "reasoning-anthropic" }), @@ -1923,7 +1793,7 @@ describe("SessionRunnerLLM", () => { }, ]) - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Continue" }), resume: false }) + yield* admit(session, "Continue") response = [] yield* session.resume(sessionID) @@ -1940,11 +1810,9 @@ describe("SessionRunnerLLM", () => { it.effect("replays durable provider-executed tool results inline in a second-turn request", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Search first" }), resume: false }) + const session = yield* setup + yield* admit(session, "Search first") - requests.length = 0 response = [ LLMEvent.stepStart({ index: 0 }), LLMEvent.toolCall({ @@ -1967,7 +1835,7 @@ describe("SessionRunnerLLM", () => { yield* session.resume(sessionID) yield* replaySessionProjection(sessionID) - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Continue" }), resume: false }) + yield* admit(session, "Continue") response = [] yield* session.resume(sessionID) @@ -1995,17 +1863,12 @@ describe("SessionRunnerLLM", () => { it.effect("starts recorded local tools eagerly and awaits settlement before continuing", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Echo five times" }), resume: false }) + const session = yield* setup + yield* admit(session, "Echo five times") - requests.length = 0 - executions.length = 0 toolExecutionGate = yield* Deferred.make() toolExecutionsStarted = yield* Deferred.make() const providerGate = yield* Deferred.make() - response = [] - responses = undefined const initial = Stream.fromIterable([ LLMEvent.stepStart({ index: 0 }), ...Array.from({ length: 5 }, (_, index) => @@ -2016,7 +1879,6 @@ describe("SessionRunnerLLM", () => { LLMEvent.stepFinish({ index: 0, reason: "tool-calls" }), LLMEvent.finish({ reason: "tool-calls" }), ]) - streamGate = undefined responseStream = Stream.concat( initial, Stream.fromEffect(Deferred.await(providerGate)).pipe(Stream.flatMap(() => final)), @@ -2056,25 +1918,12 @@ describe("SessionRunnerLLM", () => { it.effect("settles repeated provider-local tool call IDs against their owning assistant messages", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Echo twice" }), resume: false }) + const session = yield* setup + yield* admit(session, "Echo twice") - requests.length = 0 - executions.length = 0 responses = [ - [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.toolCall({ id: "tool_0", name: "echo", input: { text: "first" } }), - LLMEvent.stepFinish({ index: 0, reason: "tool-calls" }), - LLMEvent.finish({ reason: "tool-calls" }), - ], - [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.toolCall({ id: "tool_0", name: "echo", input: { text: "second" } }), - LLMEvent.stepFinish({ index: 0, reason: "tool-calls" }), - LLMEvent.finish({ reason: "tool-calls" }), - ], + reply.tool("tool_0", "echo", { text: "first" }), + reply.tool("tool_0", "echo", { text: "second" }), [], ] @@ -2144,20 +1993,10 @@ describe("SessionRunnerLLM", () => { it.effect("joins concurrent resume calls into one active provider run", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Run once" }), resume: false }) + const session = yield* setup + yield* admit(session, "Run once") - requests.length = 0 - responses = undefined - response = [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.textStart({ id: "text-once" }), - LLMEvent.textDelta({ id: "text-once", text: "Once" }), - LLMEvent.textEnd({ id: "text-once" }), - LLMEvent.stepFinish({ index: 0, reason: "stop" }), - LLMEvent.finish({ reason: "stop" }), - ] + response = reply.text("Once", "text-once") streamGate = yield* Deferred.make() streamStarted = yield* Deferred.make() @@ -2183,23 +2022,10 @@ describe("SessionRunnerLLM", () => { it.effect("steers an active provider turn with newly recorded prompts", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Start working" }), resume: false }) + const session = yield* setup + yield* admit(session, "Start working") - 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" }), - ], - ] + responses = [reply.stop(), reply.stop()] streamGate = yield* Deferred.make() streamStarted = yield* Deferred.make() @@ -2226,29 +2052,10 @@ describe("SessionRunnerLLM", () => { it.effect("promotes queued input after continuation ends", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Start working" }), resume: false }) + const session = yield* setup + yield* admit(session, "Start working") - requests.length = 0 - responses = [ - [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.toolCall({ id: "call-echo", name: "echo", input: { text: "hello" } }), - LLMEvent.stepFinish({ index: 0, reason: "tool-calls" }), - LLMEvent.finish({ reason: "tool-calls" }), - ], - [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.stepFinish({ index: 0, reason: "stop" }), - LLMEvent.finish({ reason: "stop" }), - ], - [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.stepFinish({ index: 0, reason: "stop" }), - LLMEvent.finish({ reason: "stop" }), - ], - ] + responses = [reply.tool("call-echo", "echo", { text: "hello" }), reply.stop(), reply.stop()] streamGate = yield* Deferred.make() streamStarted = yield* Deferred.make() @@ -2273,24 +2080,11 @@ describe("SessionRunnerLLM", () => { it.effect("preserves durable queued input for a later wake after interruption", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service + const session = yield* setup const { db } = yield* Database.Service - yield* session.prompt({ - sessionID, - prompt: PromptInput.Prompt.make({ text: "Interrupt current work" }), - resume: false, - }) + yield* admit(session, "Interrupt current work") - requests.length = 0 - responses = [ - [], - [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.stepFinish({ index: 0, reason: "stop" }), - LLMEvent.finish({ reason: "stop" }), - ], - ] + responses = [[], reply.stop()] streamGate = yield* Deferred.make() streamStarted = yield* Deferred.make() @@ -2320,24 +2114,11 @@ describe("SessionRunnerLLM", () => { it.effect("preserves durable steering input for a later resume after interruption", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service + const session = yield* setup const { db } = yield* Database.Service - yield* session.prompt({ - sessionID, - prompt: PromptInput.Prompt.make({ text: "Interrupt current work" }), - resume: false, - }) + yield* admit(session, "Interrupt current work") - requests.length = 0 - responses = [ - [], - [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.stepFinish({ index: 0, reason: "stop" }), - LLMEvent.finish({ reason: "stop" }), - ], - ] + responses = [[], reply.stop()] streamGate = yield* Deferred.make() streamStarted = yield* Deferred.make() @@ -2367,28 +2148,10 @@ describe("SessionRunnerLLM", () => { it.effect("promotes queued inputs one at a time in FIFO order", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Start working" }), resume: false }) + const session = yield* setup + yield* admit(session, "Start working") - 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" }), - ], - ] + responses = [reply.stop(), reply.stop(), reply.stop()] streamGate = yield* Deferred.make() streamStarted = yield* Deferred.make() @@ -2410,9 +2173,8 @@ describe("SessionRunnerLLM", () => { it.effect("promotes queued input after steering continuation ends", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Start steering" }), resume: false }) + const session = yield* setup + yield* admit(session, "Start steering") yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Queue for later" }), @@ -2420,19 +2182,7 @@ describe("SessionRunnerLLM", () => { resume: false, }) - requests.length = 0 - responses = [ - [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.stepFinish({ index: 0, reason: "stop" }), - LLMEvent.finish({ reason: "stop" }), - ], - [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.stepFinish({ index: 0, reason: "stop" }), - LLMEvent.finish({ reason: "stop" }), - ], - ] + responses = [reply.stop(), reply.stop()] yield* session.resume(sessionID) @@ -2444,33 +2194,10 @@ describe("SessionRunnerLLM", () => { 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 }) + const session = yield* setup + yield* admit(session, "Start working") - 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" }), - ], - ] + responses = [reply.stop(), reply.stop(), reply.stop(), reply.stop()] const firstGate = yield* Deferred.make() const secondGate = yield* Deferred.make() streamGate = firstGate @@ -2512,23 +2239,10 @@ describe("SessionRunnerLLM", () => { it.effect("coalesces multiple active steering prompts into one continuation turn", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Start working" }), resume: false }) + const session = yield* setup + yield* admit(session, "Start working") - 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" }), - ], - ] + responses = [reply.stop(), reply.stop()] streamGate = yield* Deferred.make() streamStarted = yield* Deferred.make() @@ -2552,13 +2266,9 @@ describe("SessionRunnerLLM", () => { it.effect("runs steering input accepted while the active provider turn fails", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Start working" }), resume: false }) + const session = yield* setup + yield* admit(session, "Start working") - requests.length = 0 - responses = undefined - response = [] streamFailure = invalidRequest() streamGate = yield* Deferred.make() streamStarted = yield* Deferred.make() @@ -2581,14 +2291,9 @@ describe("SessionRunnerLLM", () => { it.effect("durably fails local tools left running by a prior process before continuing", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service + const session = yield* setup const events = yield* EventV2.Service - yield* session.prompt({ - sessionID, - prompt: PromptInput.Prompt.make({ text: "Recover interrupted tool" }), - resume: false, - }) + yield* admit(session, "Recover interrupted tool") yield* SessionInput.promoteSteers((yield* Database.Service).db, events, sessionID) const assistantMessageID = SessionMessage.ID.create() yield* events.publish(SessionEvent.Step.Started, { @@ -2643,14 +2348,9 @@ describe("SessionRunnerLLM", () => { it.effect("durably fails hosted tools left running by a prior process before continuing inline", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service + const session = yield* setup const events = yield* EventV2.Service - yield* session.prompt({ - sessionID, - prompt: PromptInput.Prompt.make({ text: "Recover interrupted hosted tool" }), - resume: false, - }) + yield* admit(session, "Recover interrupted hosted tool") yield* SessionInput.promoteSteers((yield* Database.Service).db, events, sessionID) const assistantMessageID = SessionMessage.ID.create() yield* events.publish(SessionEvent.Step.Started, { @@ -2699,14 +2399,9 @@ describe("SessionRunnerLLM", () => { it.effect("durably fails pending tool input left by a prior process before continuing", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service + const session = yield* setup const events = yield* EventV2.Service - yield* session.prompt({ - sessionID, - prompt: PromptInput.Prompt.make({ text: "Recover interrupted tool input" }), - resume: false, - }) + yield* admit(session, "Recover interrupted tool input") yield* SessionInput.promoteSteers((yield* Database.Service).db, events, sessionID) const assistantMessageID = SessionMessage.ID.create() yield* events.publish(SessionEvent.Step.Started, { @@ -2736,8 +2431,7 @@ describe("SessionRunnerLLM", () => { it.effect("promotes the first queued input when woken while idle", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service + const session = yield* setup yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Wait in queue" }), @@ -2745,7 +2439,6 @@ describe("SessionRunnerLLM", () => { resume: false, }) - requests.length = 0 yield* (yield* SessionExecution.Service).wake(sessionID) yield* Effect.yieldNow @@ -2756,26 +2449,17 @@ describe("SessionRunnerLLM", () => { it.effect("retries inbox input after prompt projection rolls back", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service + const session = yield* setup const events = yield* EventV2.Service const defect = new Error("fail after prompt promotion") let fail = true yield* events.project(SessionEvent.PromptPromoted, () => (fail ? Effect.die(defect) : Effect.void)) - yield* session.prompt({ - sessionID, - prompt: PromptInput.Prompt.make({ text: "Recover promoted input" }), - resume: false, - }) + yield* admit(session, "Recover promoted input") expect(yield* session.resume(sessionID).pipe(Effect.catchDefect(Effect.succeed))).toBe(defect) fail = false requests.length = 0 - response = [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.stepFinish({ index: 0, reason: "stop" }), - LLMEvent.finish({ reason: "stop" }), - ] + response = reply.stop() yield* (yield* SessionExecution.Service).wake(sessionID) while (requests.length === 0) yield* Effect.yieldNow @@ -2786,21 +2470,15 @@ describe("SessionRunnerLLM", () => { it.effect("does not strand a committed promotion when a post-commit listener defects", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service + const session = yield* setup const events = yield* EventV2.Service yield* events.listen((event) => event.type === SessionEvent.PromptPromoted.type ? Effect.die("fail after prompt promotion commits") : Effect.void, ) - yield* session.prompt({ - sessionID, - prompt: PromptInput.Prompt.make({ text: "Run committed promotion" }), - resume: false, - }) + yield* admit(session, "Run committed promotion") - requests.length = 0 yield* session.resume(sessionID) expect(requests).toHaveLength(1) @@ -2810,19 +2488,15 @@ describe("SessionRunnerLLM", () => { it.effect("runs different sessions concurrently", () => Effect.gen(function* () { - yield* setup + const session = yield* setup yield* insertSession(otherSessionID) - const session = yield* SessionV2.Service - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Run first" }), resume: false }) + yield* admit(session, "Run first") yield* session.prompt({ sessionID: otherSessionID, prompt: PromptInput.Prompt.make({ text: "Run second" }), resume: false, }) - requests.length = 0 - responses = undefined - response = [] streamGate = yield* Deferred.make() streamStarted = yield* Deferred.make() @@ -2847,12 +2521,11 @@ describe("SessionRunnerLLM", () => { it.effect("bounds 64-character session prompt cache keys", () => Effect.gen(function* () { - yield* setup + const session = yield* setup const longSessionID = SessionV2.ID.make(`ses_${"a".repeat(64)}`) const otherLongSessionID = SessionV2.ID.make(`ses_${"b".repeat(64)}`) yield* insertSession(longSessionID) yield* insertSession(otherLongSessionID) - const session = yield* SessionV2.Service yield* session.prompt({ sessionID: longSessionID, prompt: PromptInput.Prompt.make({ text: "Run long session" }), @@ -2864,7 +2537,6 @@ describe("SessionRunnerLLM", () => { resume: false, }) - requests.length = 0 yield* session.resume(longSessionID) yield* session.resume(otherLongSessionID) @@ -2877,17 +2549,9 @@ describe("SessionRunnerLLM", () => { it.effect("fans out one failed run and allows a later retry", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service - yield* session.prompt({ - sessionID, - prompt: PromptInput.Prompt.make({ text: "Retry after failure" }), - resume: false, - }) + const session = yield* setup + yield* admit(session, "Retry after failure") - requests.length = 0 - responses = undefined - response = [] streamFailure = invalidRequest() streamGate = yield* Deferred.make() streamStarted = yield* Deferred.make() @@ -2912,30 +2576,10 @@ describe("SessionRunnerLLM", () => { it.effect("durably settles local tool failures before continuing", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Call missing" }), resume: false }) - - requests.length = 0 - responses = [ - [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.toolCall({ id: "call-missing", name: "missing", input: {} }), - LLMEvent.stepFinish({ index: 0, reason: "tool-calls" }), - LLMEvent.finish({ reason: "tool-calls" }), - ], - [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.textStart({ id: "text-after-error" }), - LLMEvent.textDelta({ id: "text-after-error", text: "Recovered" }), - LLMEvent.textEnd({ id: "text-after-error" }), - LLMEvent.stepFinish({ index: 0, reason: "stop" }), - LLMEvent.finish({ reason: "stop" }), - ], - ] - streamGate = undefined - streamStarted = undefined + const session = yield* setup + yield* admit(session, "Call missing") + responses = [reply.tool("call-missing", "missing", {}), reply.text("Recovered", "text-after-error")] yield* session.resume(sessionID) expect(requests).toHaveLength(2) @@ -2958,27 +2602,10 @@ describe("SessionRunnerLLM", () => { it.effect("returns unexpected local tool defects to the model and continues", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Call defect" }), resume: false }) + const session = yield* setup + yield* admit(session, "Call defect") - requests.length = 0 - responses = [ - [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.toolCall({ id: "call-defect", name: "defect", input: {} }), - LLMEvent.stepFinish({ index: 0, reason: "tool-calls" }), - LLMEvent.finish({ reason: "tool-calls" }), - ], - [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.textStart({ id: "text-after-defect" }), - LLMEvent.textDelta({ id: "text-after-defect", text: "Recovered" }), - LLMEvent.textEnd({ id: "text-after-defect" }), - LLMEvent.stepFinish({ index: 0, reason: "stop" }), - LLMEvent.finish({ reason: "stop" }), - ], - ] + responses = [reply.tool("call-defect", "defect", {}), reply.text("Recovered", "text-after-defect")] yield* session.resume(sessionID) @@ -3014,8 +2641,7 @@ describe("SessionRunnerLLM", () => { it.effect("returns policy-blocked tools to the model and continues", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service + const session = yield* setup const registry = yield* ToolRegistry.Service yield* registry.register({ blocked: Tool.make({ @@ -3028,22 +2654,9 @@ describe("SessionRunnerLLM", () => { ), }), }) - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Call blocked" }), resume: false }) + yield* admit(session, "Call blocked") - requests.length = 0 - responses = [ - [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.toolCall({ id: "call-blocked", name: "blocked", input: {} }), - LLMEvent.stepFinish({ index: 0, reason: "tool-calls" }), - LLMEvent.finish({ reason: "tool-calls" }), - ], - [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.stepFinish({ index: 0, reason: "stop" }), - LLMEvent.finish({ reason: "stop" }), - ], - ] + responses = [reply.tool("call-blocked", "blocked", {}), reply.stop()] yield* session.resume(sessionID) @@ -3063,8 +2676,7 @@ describe("SessionRunnerLLM", () => { it.effect("interrupts runner continuation when permission approval is declined", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service + const session = yield* setup const registry = yield* ToolRegistry.Service yield* registry.register({ declined: Tool.make({ @@ -3074,15 +2686,9 @@ describe("SessionRunnerLLM", () => { execute: () => Effect.die(new PermissionV2.DeclinedError()), }), }) - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Call declined" }), resume: false }) + yield* admit(session, "Call declined") - requests.length = 0 - response = [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.toolCall({ id: "call-declined", name: "declined", input: {} }), - LLMEvent.stepFinish({ index: 0, reason: "tool-calls" }), - LLMEvent.finish({ reason: "tool-calls" }), - ] + response = reply.tool("call-declined", "declined", {}) const exit = yield* session.resume(sessionID).pipe(Effect.exit) @@ -3107,8 +2713,7 @@ describe("SessionRunnerLLM", () => { it.effect("returns permission corrections to the model and continues", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service + const session = yield* setup const registry = yield* ToolRegistry.Service yield* registry.register({ corrected: Tool.make({ @@ -3121,22 +2726,9 @@ describe("SessionRunnerLLM", () => { ), }), }) - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Call corrected" }), resume: false }) + yield* admit(session, "Call corrected") - requests.length = 0 - responses = [ - [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.toolCall({ id: "call-corrected", name: "corrected", input: {} }), - LLMEvent.stepFinish({ index: 0, reason: "tool-calls" }), - LLMEvent.finish({ reason: "tool-calls" }), - ], - [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.stepFinish({ index: 0, reason: "stop" }), - LLMEvent.finish({ reason: "stop" }), - ], - ] + responses = [reply.tool("call-corrected", "corrected", {}), reply.stop()] yield* session.resume(sessionID) @@ -3156,20 +2748,10 @@ describe("SessionRunnerLLM", () => { it.effect("fails the drain when tool output persistence fails", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Call storefail" }), resume: false }) + const session = yield* setup + yield* admit(session, "Call storefail") - requests.length = 0 - responses = [ - [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.toolCall({ id: "call-storefail", name: "storefail", input: {} }), - LLMEvent.stepFinish({ index: 0, reason: "tool-calls" }), - LLMEvent.finish({ reason: "tool-calls" }), - ], - [], - ] + responses = [reply.tool("call-storefail", "storefail", {}), []] const exit = yield* session.resume(sessionID).pipe(Effect.exit) @@ -3201,8 +2783,7 @@ describe("SessionRunnerLLM", () => { it.effect("preserves permission rejection and stops before continuation", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service + const session = yield* setup const registry = yield* ToolRegistry.Service yield* registry.register({ permissionfail: Tool.make({ @@ -3220,19 +2801,9 @@ describe("SessionRunnerLLM", () => { }), }), }) - yield* session.prompt({ - sessionID, - prompt: PromptInput.Prompt.make({ text: "Reject permission" }), - resume: false, - }) - requests.length = 0 + yield* admit(session, "Reject permission") responses = [ - [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.toolCall({ id: "call-permission", name: "permissionfail", input: {} }), - LLMEvent.stepFinish({ index: 0, reason: "tool-calls" }), - LLMEvent.finish({ reason: "tool-calls" }), - ], + reply.tool("call-permission", "permissionfail", {}), [LLMEvent.stepStart({ index: 0 }), LLMEvent.stepFinish({ index: 0, reason: "stop" })], ] @@ -3270,8 +2841,7 @@ describe("SessionRunnerLLM", () => { it.effect("interrupts runner continuation when a question is cancelled", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service + const session = yield* setup const registry = yield* ToolRegistry.Service yield* registry.register({ question: Tool.make({ @@ -3281,18 +2851,9 @@ describe("SessionRunnerLLM", () => { execute: () => Effect.die(new QuestionTool.CancelledError()), }), }) - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Ask then stop" }), resume: false }) + yield* admit(session, "Ask then stop") - requests.length = 0 - responses = [ - [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.toolCall({ id: "call-question", name: "question", input: {} }), - LLMEvent.stepFinish({ index: 0, reason: "tool-calls" }), - LLMEvent.finish({ reason: "tool-calls" }), - ], - [], - ] + responses = [reply.tool("call-question", "question", {}), []] const run = yield* session.resume(sessionID).pipe(Effect.exit, Effect.forkChild) const exit = yield* Fiber.join(run) @@ -3318,13 +2879,8 @@ describe("SessionRunnerLLM", () => { it.effect("awaits started local tools before surfacing provider stream failure", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service - yield* session.prompt({ - sessionID, - prompt: PromptInput.Prompt.make({ text: "Settle before failing" }), - resume: false, - }) + const session = yield* setup + yield* admit(session, "Settle before failing") const failure = providerUnavailable() toolExecutionGate = yield* Deferred.make() responseStream = Stream.concat( @@ -3364,14 +2920,8 @@ describe("SessionRunnerLLM", () => { it.effect("durably fails blocked local tools when a provider turn is interrupted", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service - yield* session.prompt({ - sessionID, - prompt: PromptInput.Prompt.make({ text: "Interrupt blocked tool" }), - resume: false, - }) - executions.length = 0 + const session = yield* setup + yield* admit(session, "Interrupt blocked tool") toolExecutionGate = yield* Deferred.make() responseStream = Stream.concat( Stream.fromIterable([ @@ -3426,15 +2976,8 @@ describe("SessionRunnerLLM", () => { it.effect("interrupts a blocked provider turn without local tool execution", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service - yield* session.prompt({ - sessionID, - prompt: PromptInput.Prompt.make({ text: "Interrupt provider" }), - resume: false, - }) - requests.length = 0 - response = [] + const session = yield* setup + yield* admit(session, "Interrupt provider") streamGate = yield* Deferred.make() streamStarted = yield* Deferred.make() @@ -3458,23 +3001,12 @@ describe("SessionRunnerLLM", () => { it.effect("durably fails blocked local tools when interrupted while awaiting settlement", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service - yield* session.prompt({ - sessionID, - prompt: PromptInput.Prompt.make({ text: "Interrupt tool settlement" }), - resume: false, - }) - executions.length = 0 + const session = yield* setup + yield* admit(session, "Interrupt tool settlement") toolExecutionGate = yield* Deferred.make() toolExecutionsStarted = yield* Deferred.make() toolExecutionsReady = 1 - response = [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.toolCall({ id: "call-await-interrupt", name: "echo", input: { text: "blocked" } }), - LLMEvent.stepFinish({ index: 0, reason: "tool-calls" }), - LLMEvent.finish({ reason: "tool-calls" }), - ] + response = reply.tool("call-await-interrupt", "echo", { text: "blocked" }) const runner = yield* SessionRunner.Service const run = yield* runner.drain({ sessionID, force: true }).pipe(Effect.forkChild) @@ -3506,35 +3038,18 @@ describe("SessionRunnerLLM", () => { it.effect("forces a text response on an agent's configured final step", () => Effect.gen(function* () { - yield* setup + const session = yield* setup const agents = yield* AgentV2.Service yield* agents.transform((editor) => editor.update(AgentV2.ID.make("build"), (agent) => { agent.steps = 2 }), ) - const session = yield* SessionV2.Service - yield* session.prompt({ - sessionID, - prompt: PromptInput.Prompt.make({ text: "Finish at the limit" }), - resume: false, - }) + yield* admit(session, "Finish at the limit") - requests.length = 0 - executions.length = 0 responses = [ - [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.toolCall({ id: "call-terminal", name: "echo", input: { text: "done" } }), - LLMEvent.stepFinish({ index: 0, reason: "tool-calls" }), - LLMEvent.finish({ reason: "tool-calls" }), - ], - [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.toolCall({ id: "call-forbidden", name: "echo", input: { text: "forbidden" } }), - LLMEvent.stepFinish({ index: 0, reason: "tool-calls" }), - LLMEvent.finish({ reason: "tool-calls" }), - ], + reply.tool("call-terminal", "echo", { text: "done" }), + reply.tool("call-forbidden", "echo", { text: "forbidden" }), ] yield* session.resume(sessionID) @@ -3558,36 +3073,19 @@ describe("SessionRunnerLLM", () => { it.effect("resets the configured step allowance when steering input promotes", () => Effect.gen(function* () { - yield* setup + const session = yield* setup const agents = yield* AgentV2.Service yield* agents.transform((editor) => editor.update(AgentV2.ID.make("build"), (agent) => { agent.steps = 2 }), ) - const session = yield* SessionV2.Service - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Start work" }), resume: false }) + yield* admit(session, "Start work") - requests.length = 0 - executions.length = 0 responses = [ - [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.toolCall({ id: "call-before-steer", name: "echo", input: { text: "before" } }), - LLMEvent.stepFinish({ index: 0, reason: "tool-calls" }), - LLMEvent.finish({ reason: "tool-calls" }), - ], - [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.toolCall({ id: "call-after-steer", name: "echo", input: { text: "after" } }), - LLMEvent.stepFinish({ index: 0, reason: "tool-calls" }), - LLMEvent.finish({ reason: "tool-calls" }), - ], - [ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.stepFinish({ index: 0, reason: "stop" }), - LLMEvent.finish({ reason: "stop" }), - ], + reply.tool("call-before-steer", "echo", { text: "before" }), + reply.tool("call-after-steer", "echo", { text: "after" }), + reply.stop(), ] streamGate = yield* Deferred.make() streamStarted = yield* Deferred.make() @@ -3610,14 +3108,9 @@ describe("SessionRunnerLLM", () => { it.effect("projects provider errors as terminal assistant step failures", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Fail durably" }), resume: false }) + const session = yield* setup + yield* admit(session, "Fail durably") - requests.length = 0 - responses = undefined - streamGate = undefined - streamStarted = undefined response = [LLMEvent.stepStart({ index: 0 }), LLMEvent.providerError({ message: "Provider unavailable" })] expect((yield* session.resume(sessionID).pipe(Effect.flip)).message).toBe("Provider unavailable") @@ -3632,11 +3125,9 @@ describe("SessionRunnerLLM", () => { it.effect("projects provider errors emitted before assistant step start", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Fail before step" }), resume: false }) + const session = yield* setup + yield* admit(session, "Fail before step") - requests.length = 0 response = [LLMEvent.providerError({ message: "Provider unavailable" })] expect((yield* session.resume(sessionID).pipe(Effect.flip)).message).toBe("Provider unavailable") @@ -3651,9 +3142,8 @@ describe("SessionRunnerLLM", () => { it.effect("projects content-filter finishes as visible terminal failures", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Blocked response" }), resume: false }) + const session = yield* setup + yield* admit(session, "Blocked response") response = [ LLMEvent.stepStart({ index: 0 }), LLMEvent.textStart({ id: "partial" }), @@ -3688,13 +3178,8 @@ describe("SessionRunnerLLM", () => { it.effect("settles a local tool before one content-filter step failure", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service - yield* session.prompt({ - sessionID, - prompt: PromptInput.Prompt.make({ text: "Tool before blocked response" }), - resume: false, - }) + const session = yield* setup + yield* admit(session, "Tool before blocked response") toolExecutionGate = yield* Deferred.make() toolExecutionsStarted = yield* Deferred.make() toolExecutionsReady = 1 @@ -3728,15 +3213,9 @@ describe("SessionRunnerLLM", () => { it.effect("does not recover context overflow after durable assistant output", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service - yield* session.prompt({ - sessionID, - prompt: PromptInput.Prompt.make({ text: "Fail after output" }), - resume: false, - }) + const session = yield* setup + yield* admit(session, "Fail after output") - requests.length = 0 response = [ LLMEvent.stepStart({ index: 0 }), LLMEvent.textStart({ id: "text-partial" }), @@ -3761,13 +3240,8 @@ describe("SessionRunnerLLM", () => { it.effect("projects raw provider stream failures as terminal assistant step failures", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service - yield* session.prompt({ - sessionID, - prompt: PromptInput.Prompt.make({ text: "Fail raw stream durably" }), - resume: false, - }) + const session = yield* setup + yield* admit(session, "Fail raw stream durably") const failure = invalidRequest() responseStream = Stream.fail(failure) @@ -3782,12 +3256,10 @@ describe("SessionRunnerLLM", () => { it.effect("retries eligible pre-output failures after exponential backoff", () => 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 + const session = yield* setup + yield* admit(session, "Retry transport") responseStream = Stream.fail(providerUnavailable()) - response = fragmentFixture("text", "retry-success", ["Recovered"]).completeEvents + response = reply.text("Recovered", "retry-success") const run = yield* session.resume(sessionID).pipe(Effect.forkChild) while (requests.length < 1) yield* Effect.yieldNow @@ -3811,12 +3283,10 @@ describe("SessionRunnerLLM", () => { it.effect("uses a larger provider retry-after delay", () => 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 + const session = yield* setup + yield* admit(session, "Retry rate limit") responseStream = Stream.fail(rateLimited(5_000)) - response = fragmentFixture("text", "retry-after-success", ["Recovered"]).completeEvents + response = reply.text("Recovered", "retry-after-success") const run = yield* session.resume(sessionID).pipe(Effect.forkChild) while (requests.length < 1) yield* Effect.yieldNow @@ -3830,10 +3300,8 @@ describe("SessionRunnerLLM", () => { it.effect("stops after five total retry attempts", () => 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 + const session = yield* setup + yield* admit(session, "Exhaust retries") streamFailure = providerUnavailable() const run = yield* session.resume(sessionID).pipe(Effect.forkChild) @@ -3866,20 +3334,14 @@ describe("SessionRunnerLLM", () => { it.effect("counts retry attempts against the agent step allowance", () => Effect.gen(function* () { - yield* setup + const session = yield* setup const agents = yield* AgentV2.Service yield* agents.transform((editor) => editor.update(AgentV2.ID.make("build"), (agent) => { agent.steps = 2 }), ) - const session = yield* SessionV2.Service - yield* session.prompt({ - sessionID, - prompt: PromptInput.Prompt.make({ text: "Bound retries by steps" }), - resume: false, - }) - requests.length = 0 + yield* admit(session, "Bound retries by steps") const failure = providerUnavailable() responseStream = Stream.fail(failure) streamFailure = failure @@ -3899,10 +3361,8 @@ describe("SessionRunnerLLM", () => { it.effect("does not retry non-eligible provider failures", () => 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 session = yield* setup + yield* admit(session, "Do not retry") const failure = invalidRequest() streamFailure = failure @@ -3914,16 +3374,9 @@ describe("SessionRunnerLLM", () => { it.effect("does not continue automatically after a provider error follows a local tool call", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service - yield* session.prompt({ - sessionID, - prompt: PromptInput.Prompt.make({ text: "Do not continue failed provider" }), - resume: false, - }) + const session = yield* setup + yield* admit(session, "Do not continue failed provider") - requests.length = 0 - const executionCount = executions.length toolExecutionGate = yield* Deferred.make() toolExecutionsStarted = yield* Deferred.make() toolExecutionsReady = 1 @@ -3941,7 +3394,7 @@ describe("SessionRunnerLLM", () => { toolExecutionsStarted = undefined expect(requests).toHaveLength(1) - expect(executions.slice(executionCount)).toEqual(["settled"]) + expect(executions).toEqual(["settled"]) const context = yield* session.context(sessionID) const assistant = requireAssistant(context) expect((yield* recordedStepSettlementEvents(sessionID, assistant.id)).map((event) => event.type)).toEqual([ @@ -3955,23 +3408,12 @@ describe("SessionRunnerLLM", () => { it.effect("durably fails a hosted tool when its provider errors before returning a result", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service - yield* session.prompt({ - sessionID, - prompt: PromptInput.Prompt.make({ text: "Fail hosted tool durably" }), - resume: false, - }) + const session = yield* setup + yield* admit(session, "Fail hosted tool durably") - requests.length = 0 response = [ LLMEvent.stepStart({ index: 0 }), - LLMEvent.toolCall({ - id: "call-hosted-provider-error", - name: "web_search", - input: { query: "effect" }, - providerExecuted: true, - }), + hostedCall("call-hosted-provider-error", "effect"), LLMEvent.providerError({ message: "Provider unavailable" }), ] @@ -3998,13 +3440,8 @@ describe("SessionRunnerLLM", () => { it.effect("preserves a tool defect before provider failure settlement", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service - yield* session.prompt({ - sessionID, - prompt: PromptInput.Prompt.make({ text: "Defect while provider fails" }), - resume: false, - }) + const session = yield* setup + yield* admit(session, "Defect while provider fails") response = [ LLMEvent.stepStart({ index: 0 }), LLMEvent.toolCall({ id: "call-defect-provider-error", name: "defect", input: {} }), @@ -4028,22 +3465,9 @@ describe("SessionRunnerLLM", () => { it.effect("durably fails a hosted tool left unresolved at normal provider EOF", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service - yield* session.prompt({ - sessionID, - 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, - }), - ] + const session = yield* setup + yield* admit(session, "Fail hosted tool at EOF") + response = [LLMEvent.stepStart({ index: 0 }), hostedCall("call-hosted-eof", "effect")] expect((yield* session.resume(sessionID).pipe(Effect.flip)).message).toBe("Provider did not return a tool result") const assistant = requireAssistant(yield* session.context(sessionID)) @@ -4073,21 +3497,11 @@ describe("SessionRunnerLLM", () => { it.effect("fails an unresolved hosted tool before one clean step end", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service - yield* session.prompt({ - sessionID, - prompt: PromptInput.Prompt.make({ text: "Settle hosted tool before ending" }), - resume: false, - }) + const session = yield* setup + yield* admit(session, "Settle hosted tool before ending") response = [ LLMEvent.stepStart({ index: 0 }), - LLMEvent.toolCall({ - id: "call-hosted-clean-end", - name: "web_search", - input: { query: "effect" }, - providerExecuted: true, - }), + hostedCall("call-hosted-clean-end", "effect"), LLMEvent.stepFinish({ index: 0, reason: "stop" }), LLMEvent.finish({ reason: "stop" }), ] @@ -4110,13 +3524,8 @@ describe("SessionRunnerLLM", () => { it.effect("settles unresolved local and hosted tools before one raw provider failure", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service - yield* session.prompt({ - sessionID, - prompt: PromptInput.Prompt.make({ text: "Fail unresolved tools" }), - resume: false, - }) + const session = yield* setup + yield* admit(session, "Fail unresolved tools") const failure = invalidRequest() const providerFailed = yield* Deferred.make() toolExecutionGate = yield* Deferred.make() @@ -4124,12 +3533,7 @@ describe("SessionRunnerLLM", () => { 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, - }), + hostedCall("call-hosted-raw-failure-pair", "effect"), ]), Stream.fromEffect(Deferred.succeed(providerFailed, undefined)).pipe(Stream.flatMap(() => Stream.fail(failure))), ) @@ -4158,25 +3562,11 @@ describe("SessionRunnerLLM", () => { it.effect("durably fails a hosted tool left unresolved by a raw provider stream failure", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service - yield* session.prompt({ - sessionID, - prompt: PromptInput.Prompt.make({ text: "Fail hosted tool on raw failure" }), - resume: false, - }) - requests.length = 0 + const session = yield* setup + yield* admit(session, "Fail hosted tool on raw failure") const failure = providerUnavailable() responseStream = Stream.concat( - Stream.fromIterable([ - LLMEvent.stepStart({ index: 0 }), - LLMEvent.toolCall({ - id: "call-hosted-raw-failure", - name: "web_search", - input: { query: "effect" }, - providerExecuted: true, - }), - ]), + Stream.fromIterable([LLMEvent.stepStart({ index: 0 }), hostedCall("call-hosted-raw-failure", "effect")]), Stream.fail(failure), ) @@ -4208,13 +3598,9 @@ describe("SessionRunnerLLM", () => { it.effect("rejects a second text start before the open fragment ends", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Two blocks" }), resume: false }) + const session = yield* setup + yield* admit(session, "Two blocks") - responses = undefined - streamGate = undefined - streamStarted = undefined response = [ LLMEvent.stepStart({ index: 0 }), LLMEvent.textStart({ id: "text-1" }), @@ -4230,13 +3616,9 @@ describe("SessionRunnerLLM", () => { it.effect("projects sequential text fragments as separate content parts", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service - yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Two blocks" }), resume: false }) + const session = yield* setup + yield* admit(session, "Two blocks") - responses = undefined - streamGate = undefined - streamStarted = undefined response = [ LLMEvent.stepStart({ index: 0 }), LLMEvent.textStart({ id: "text-1" }), @@ -4278,11 +3660,7 @@ describe("SessionRunnerLLM", () => { it.effect("rejects duplicate streamed text starts", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service - responses = undefined - streamGate = undefined - streamStarted = undefined + const session = yield* setup response = [LLMEvent.textStart({ id: "text-1" }), LLMEvent.textStart({ id: "text-1" })] const defect = yield* session.resume(sessionID).pipe(Effect.catchDefect(Effect.succeed)) @@ -4294,23 +3672,15 @@ describe("SessionRunnerLLM", () => { it.effect("transitions streamed raw tool input to parsed called input", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service - yield* session.prompt({ - sessionID, - prompt: PromptInput.Prompt.make({ text: "Call provider tool" }), - resume: false, - }) + const session = yield* setup + yield* admit(session, "Call provider tool") - 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 }), + hostedCall("call-parsed", "hello"), LLMEvent.stepFinish({ index: 0, reason: "stop" }), LLMEvent.finish({ reason: "stop" }), ] @@ -4329,11 +3699,7 @@ describe("SessionRunnerLLM", () => { it.effect("rejects malformed streamed tool input ordering", () => Effect.gen(function* () { - yield* setup - const session = yield* SessionV2.Service - responses = undefined - streamGate = undefined - streamStarted = undefined + const session = yield* setup response = [LLMEvent.toolInputDelta({ id: "call-1", name: "read", text: "{}" })] const defect = yield* session.resume(sessionID).pipe(Effect.catchDefect(Effect.succeed))