test(core): preserve runner execution boundaries
This commit is contained in:
parent
68b23a80eb
commit
8046b1066f
3 changed files with 444 additions and 393 deletions
|
|
@ -14,9 +14,9 @@ class RunnerEndedError extends Error {
|
|||
}
|
||||
}
|
||||
|
||||
class RunnerScenarioUsedError extends Error {
|
||||
class RunnerScenarioActiveError extends Error {
|
||||
constructor() {
|
||||
super("RunnerScenario.run may only be called once")
|
||||
super("RunnerScenario.run is already active")
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -48,29 +48,41 @@ export interface RunnerLLMCall {
|
|||
}
|
||||
}
|
||||
|
||||
interface CallState {
|
||||
readonly generation: number
|
||||
readonly status: "pending" | "responded" | "closed"
|
||||
}
|
||||
|
||||
interface State {
|
||||
readonly started: boolean
|
||||
readonly accepting: boolean
|
||||
readonly ended: boolean
|
||||
readonly active: number | undefined
|
||||
readonly nextGeneration: number
|
||||
readonly nextID: number
|
||||
readonly calls: ReadonlyMap<number, "pending" | "responded" | "closed">
|
||||
readonly calls: ReadonlyMap<number, CallState>
|
||||
readonly requests: ReadonlyArray<LLMRequest>
|
||||
}
|
||||
|
||||
type Registration =
|
||||
| { readonly _tag: "Registered"; readonly id: number }
|
||||
| { readonly _tag: "Registered"; readonly id: number; readonly generation: number }
|
||||
| { readonly _tag: "Unexpected"; readonly error: Error }
|
||||
|
||||
interface PendingCall {
|
||||
readonly _tag: "Call"
|
||||
readonly generation: number
|
||||
readonly id: number
|
||||
readonly call: RunnerLLMCall
|
||||
}
|
||||
|
||||
interface RunnerEnd {
|
||||
readonly _tag: "End"
|
||||
readonly generation: number
|
||||
}
|
||||
|
||||
interface RunnerLLMInternal {
|
||||
readonly llm: RunnerLLM
|
||||
readonly begin: Effect.Effect<void, RunnerScenarioUsedError>
|
||||
readonly finishInteraction: Effect.Effect<void, Error>
|
||||
readonly end: Effect.Effect<void>
|
||||
readonly begin: Effect.Effect<number, RunnerScenarioActiveError>
|
||||
readonly finishInteraction: (generation: number) => Effect.Effect<void, Error>
|
||||
readonly end: (generation: number) => Effect.Effect<void>
|
||||
}
|
||||
|
||||
export interface RunnerLLM {
|
||||
|
|
@ -115,11 +127,11 @@ const toolCall = (
|
|||
}
|
||||
|
||||
const makeRunnerLLM = Effect.gen(function* () {
|
||||
const calls = yield* Queue.unbounded<PendingCall, Cause.Done>()
|
||||
const calls = yield* Queue.unbounded<PendingCall | RunnerEnd>()
|
||||
const state = yield* Ref.make<State>({
|
||||
started: false,
|
||||
accepting: false,
|
||||
ended: false,
|
||||
active: undefined,
|
||||
nextGeneration: 1,
|
||||
nextID: 1,
|
||||
calls: new Map(),
|
||||
requests: [],
|
||||
|
|
@ -133,8 +145,9 @@ const makeRunnerLLM = Effect.gen(function* () {
|
|||
Effect.uninterruptible(
|
||||
Effect.gen(function* () {
|
||||
const pending = yield* Ref.modify(state, (current) => {
|
||||
if (current.calls.get(id) !== "pending") return [false, current]
|
||||
return [true, { ...current, calls: new Map(current.calls).set(id, "responded") }]
|
||||
const call = current.calls.get(id)
|
||||
if (call?.status !== "pending") return [false, current]
|
||||
return [true, { ...current, calls: new Map(current.calls).set(id, { ...call, status: "responded" }) }]
|
||||
})
|
||||
if (!pending) return yield* Effect.fail(new RunnerLLMResponseClosedError(id))
|
||||
yield* Deferred.succeed(response, stream)
|
||||
|
|
@ -147,7 +160,7 @@ const makeRunnerLLM = Effect.gen(function* () {
|
|||
Effect.gen(function* () {
|
||||
const response = yield* Deferred.make<Stream.Stream<LLMEvent, LLMError>>()
|
||||
const registered = yield* Ref.modify<State, Registration>(state, (current) => {
|
||||
if (!current.accepting) {
|
||||
if (!current.accepting || current.active === undefined) {
|
||||
const error = new RunnerLLMRequestError(current.nextID)
|
||||
return [
|
||||
{ _tag: "Unexpected" as const, error },
|
||||
|
|
@ -155,11 +168,14 @@ const makeRunnerLLM = Effect.gen(function* () {
|
|||
]
|
||||
}
|
||||
return [
|
||||
{ _tag: "Registered" as const, id: current.nextID },
|
||||
{ _tag: "Registered" as const, id: current.nextID, generation: current.active },
|
||||
{
|
||||
...current,
|
||||
nextID: current.nextID + 1,
|
||||
calls: new Map(current.calls).set(current.nextID, "pending"),
|
||||
calls: new Map(current.calls).set(current.nextID, {
|
||||
generation: current.active,
|
||||
status: "pending",
|
||||
}),
|
||||
requests: [...current.requests, request],
|
||||
},
|
||||
]
|
||||
|
|
@ -167,6 +183,8 @@ const makeRunnerLLM = Effect.gen(function* () {
|
|||
if (registered._tag === "Unexpected") return yield* Effect.fail(registered.error)
|
||||
const reply = (stream: Stream.Stream<LLMEvent, LLMError>) => respond(registered.id, response, stream)
|
||||
yield* Queue.offer(calls, {
|
||||
_tag: "Call",
|
||||
generation: registered.generation,
|
||||
id: registered.id,
|
||||
call: {
|
||||
request,
|
||||
|
|
@ -182,11 +200,11 @@ const makeRunnerLLM = Effect.gen(function* () {
|
|||
})
|
||||
return yield* restore(Deferred.await(response)).pipe(
|
||||
Effect.onInterrupt(() =>
|
||||
Ref.update(state, (current) =>
|
||||
current.calls.get(registered.id) === "pending"
|
||||
? { ...current, calls: new Map(current.calls).set(registered.id, "closed") }
|
||||
: current,
|
||||
),
|
||||
Ref.update(state, (current) => {
|
||||
const call = current.calls.get(registered.id)
|
||||
if (call?.status !== "pending") return current
|
||||
return { ...current, calls: new Map(current.calls).set(registered.id, { ...call, status: "closed" }) }
|
||||
}),
|
||||
),
|
||||
)
|
||||
}),
|
||||
|
|
@ -195,11 +213,12 @@ const makeRunnerLLM = Effect.gen(function* () {
|
|||
|
||||
const next = (): Effect.Effect<RunnerLLMCall, Error> =>
|
||||
Effect.gen(function* () {
|
||||
if ((yield* Ref.get(state)).ended) return yield* Effect.fail(new RunnerEndedError())
|
||||
const pending = yield* Queue.take(calls).pipe(Effect.mapError(() => new RunnerEndedError()))
|
||||
const current = yield* Ref.get(state)
|
||||
if (current.ended) return yield* Effect.fail(new RunnerEndedError())
|
||||
if (current.calls.get(pending.id) === "pending") return pending.call
|
||||
const generation = (yield* Ref.get(state)).active
|
||||
if (generation === undefined) return yield* Effect.fail(new RunnerEndedError())
|
||||
const pending = yield* Queue.take(calls)
|
||||
if (pending.generation !== generation) return yield* next()
|
||||
if (pending._tag === "End") return yield* Effect.fail(new RunnerEndedError())
|
||||
if ((yield* Ref.get(state)).calls.get(pending.id)?.status === "pending") return pending.call
|
||||
return yield* next()
|
||||
})
|
||||
|
||||
|
|
@ -216,22 +235,44 @@ const makeRunnerLLM = Effect.gen(function* () {
|
|||
next,
|
||||
requests: Ref.get(state).pipe(Effect.map((state) => state.requests)),
|
||||
},
|
||||
begin: Ref.modify<State, Effect.Effect<void, RunnerScenarioUsedError>>(state, (current) =>
|
||||
current.started
|
||||
? [Effect.fail(new RunnerScenarioUsedError()), current]
|
||||
: [Effect.void, { ...current, started: true, accepting: true }],
|
||||
begin: Ref.modify<State, Effect.Effect<number, RunnerScenarioActiveError>>(state, (current) =>
|
||||
current.active !== undefined
|
||||
? [Effect.fail(new RunnerScenarioActiveError()), current]
|
||||
: [
|
||||
Effect.succeed(current.nextGeneration),
|
||||
{
|
||||
...current,
|
||||
accepting: true,
|
||||
active: current.nextGeneration,
|
||||
nextGeneration: current.nextGeneration + 1,
|
||||
},
|
||||
],
|
||||
).pipe(Effect.flatten),
|
||||
finishInteraction: Effect.gen(function* () {
|
||||
const unanswered = yield* Ref.modify(state, (current) => [
|
||||
[...current.calls].find(([, status]) => status === "pending")?.[0],
|
||||
{ ...current, accepting: false },
|
||||
])
|
||||
if (unanswered !== undefined) return yield* Effect.fail(new RunnerLLMRequestError(unanswered))
|
||||
}),
|
||||
end: Effect.gen(function* () {
|
||||
yield* Ref.update(state, (current) => ({ ...current, accepting: false, ended: true }))
|
||||
yield* Queue.end(calls)
|
||||
}),
|
||||
finishInteraction: (generation) =>
|
||||
Effect.gen(function* () {
|
||||
const unanswered = yield* Ref.modify(state, (current) => [
|
||||
[...current.calls].find(([, call]) => call.generation === generation && call.status === "pending")?.[0],
|
||||
current.active === generation ? { ...current, accepting: false } : current,
|
||||
])
|
||||
if (unanswered !== undefined) return yield* Effect.fail(new RunnerLLMRequestError(unanswered))
|
||||
}),
|
||||
end: (generation) =>
|
||||
Effect.gen(function* () {
|
||||
yield* Ref.update(state, (current) => ({
|
||||
...current,
|
||||
accepting: current.active === generation ? false : current.accepting,
|
||||
active: current.active === generation ? undefined : current.active,
|
||||
calls: new Map(
|
||||
[...current.calls].map(([id, call]) => [
|
||||
id,
|
||||
call.generation === generation && call.status === "pending"
|
||||
? { ...call, status: "closed" as const }
|
||||
: call,
|
||||
]),
|
||||
),
|
||||
}))
|
||||
yield* Queue.offer(calls, { _tag: "End", generation })
|
||||
}),
|
||||
} satisfies RunnerLLMInternal
|
||||
})
|
||||
|
||||
|
|
@ -253,13 +294,13 @@ export namespace RunnerScenario {
|
|||
): Effect.Effect<A, RunError | E | Error, RunRequirements | R> =>
|
||||
Effect.uninterruptibleMask((restore) =>
|
||||
Effect.gen(function* () {
|
||||
yield* internal.begin
|
||||
const runner = yield* start(internal.llm).pipe(Effect.ensuring(internal.end), Effect.forkChild)
|
||||
const generation = yield* internal.begin
|
||||
const runner = yield* start(internal.llm).pipe(Effect.ensuring(internal.end(generation)), Effect.forkChild)
|
||||
return yield* restore(
|
||||
Effect.gen(function* () {
|
||||
const interactionExit = yield* Effect.gen(interaction).pipe(Effect.exit)
|
||||
if (Exit.isSuccess(interactionExit)) {
|
||||
yield* internal.finishInteraction
|
||||
yield* internal.finishInteraction(generation)
|
||||
yield* Fiber.join(runner)
|
||||
return interactionExit.value
|
||||
}
|
||||
|
|
@ -267,7 +308,7 @@ export namespace RunnerScenario {
|
|||
if (Option.isSome(error) && error.value instanceof RunnerEndedError) yield* Fiber.join(runner)
|
||||
return yield* Effect.failCause(interactionExit.cause)
|
||||
}),
|
||||
).pipe(Effect.ensuring(Fiber.interrupt(runner).pipe(Effect.andThen(internal.end))))
|
||||
).pipe(Effect.ensuring(Fiber.interrupt(runner).pipe(Effect.andThen(internal.end(generation)))))
|
||||
}),
|
||||
)
|
||||
return {
|
||||
|
|
|
|||
|
|
@ -224,11 +224,13 @@ describe("RunnerScenario", () => {
|
|||
it.effect("interrupts a runner that hangs after its final response", () =>
|
||||
Effect.gen(function* () {
|
||||
const started = yield* Deferred.make<void>()
|
||||
const interrupted = yield* Deferred.make<void>()
|
||||
const scenario = yield* RunnerScenario.make((llm) =>
|
||||
LLM.stream(request("Hang")).pipe(
|
||||
Stream.runDrain,
|
||||
Effect.andThen(Deferred.succeed(started, undefined)),
|
||||
Effect.andThen(Effect.never),
|
||||
Effect.onInterrupt(() => Deferred.succeed(interrupted, undefined)),
|
||||
Effect.provide(llm.layer),
|
||||
),
|
||||
)
|
||||
|
|
@ -242,11 +244,12 @@ describe("RunnerScenario", () => {
|
|||
|
||||
yield* Fiber.interrupt(run)
|
||||
|
||||
expect(Deferred.isDoneUnsafe(interrupted)).toBe(true)
|
||||
expect(run.pollUnsafe()).toBeDefined()
|
||||
}),
|
||||
)
|
||||
|
||||
it.effect("rejects running one scenario twice", () =>
|
||||
it.effect("runs one scenario sequentially while retaining request history", () =>
|
||||
Effect.gen(function* () {
|
||||
const scenario = yield* RunnerScenario.make((llm) =>
|
||||
LLM.stream(request("Once")).pipe(Stream.runDrain, Effect.provide(llm.layer)),
|
||||
|
|
@ -254,13 +257,35 @@ describe("RunnerScenario", () => {
|
|||
yield* scenario.run(function* () {
|
||||
yield* (yield* scenario.llm.next()).respond.stop()
|
||||
})
|
||||
yield* scenario.run(function* () {
|
||||
yield* (yield* scenario.llm.next()).respond.stop()
|
||||
})
|
||||
|
||||
expect(yield* scenario.llm.requests).toHaveLength(2)
|
||||
}),
|
||||
)
|
||||
|
||||
it.effect("rejects concurrent runs of one scenario", () =>
|
||||
Effect.gen(function* () {
|
||||
const scenario = yield* RunnerScenario.make((llm) =>
|
||||
LLM.stream(request("Once")).pipe(Stream.runDrain, Effect.provide(llm.layer)),
|
||||
)
|
||||
const active = yield* scenario
|
||||
.run(function* () {
|
||||
yield* scenario.llm.next()
|
||||
yield* Effect.never
|
||||
})
|
||||
.pipe(Effect.forkChild)
|
||||
yield* Effect.yieldNow
|
||||
|
||||
const failure = yield* scenario.run(function* (): Effect.gen.Return<void> {}).pipe(Effect.exit)
|
||||
expect(failure._tag).toBe("Failure")
|
||||
if (failure._tag !== "Failure") throw new Error("Expected second run to fail")
|
||||
if (failure._tag !== "Failure") throw new Error("Expected concurrent run to fail")
|
||||
const error = failure.cause.reasons.find((reason) => reason._tag === "Fail")
|
||||
expect(error?._tag).toBe("Fail")
|
||||
if (error?._tag !== "Fail") throw new Error("Expected a typed failure")
|
||||
expect(error.error).toMatchObject({ message: "RunnerScenario.run may only be called once" })
|
||||
expect(error.error).toMatchObject({ message: "RunnerScenario.run is already active" })
|
||||
yield* Fiber.interrupt(active)
|
||||
}),
|
||||
)
|
||||
})
|
||||
|
|
|
|||
|
|
@ -389,16 +389,21 @@ const rateLimited = (retryAfterMs?: number) =>
|
|||
reason: new RateLimitReason({ message: "Rate limited", retryAfterMs }),
|
||||
})
|
||||
|
||||
const setupOverflowRecovery = Effect.gen(function* () {
|
||||
yield* setup
|
||||
const session = yield* SessionV2.Service
|
||||
yield* session.prompt({
|
||||
sessionID,
|
||||
prompt: PromptInput.Prompt.make({ text: "Earlier question ".repeat(700) }),
|
||||
resume: false,
|
||||
const setupOverflowRecovery = (scenario: SessionScenario) =>
|
||||
Effect.gen(function* () {
|
||||
yield* setup
|
||||
const session = yield* SessionV2.Service
|
||||
yield* session.prompt({
|
||||
sessionID,
|
||||
prompt: PromptInput.Prompt.make({ text: "Earlier question ".repeat(700) }),
|
||||
resume: false,
|
||||
})
|
||||
yield* scenario.run(function* () {
|
||||
yield* (yield* scenario.llm.next()).respond.text("Earlier answer", { id: "text-earlier" })
|
||||
})
|
||||
currentModel = recoveryModel
|
||||
return session
|
||||
})
|
||||
return session
|
||||
})
|
||||
|
||||
const messageTexts = (request: LLMRequest, role: "user" | "system") =>
|
||||
request.messages.flatMap((message) =>
|
||||
|
|
@ -663,7 +668,11 @@ describe("SessionRunnerLLM", () => {
|
|||
resume: false,
|
||||
})
|
||||
|
||||
const exit = yield* session.resume(sessionID).pipe(Effect.exit)
|
||||
const exit = yield* scenario
|
||||
.run(function* () {
|
||||
yield* scenario.llm.next()
|
||||
})
|
||||
.pipe(Effect.exit)
|
||||
|
||||
expect(Exit.isFailure(exit)).toBe(true)
|
||||
if (Exit.isFailure(exit)) expect(Cause.squash(exit.cause)).toBeInstanceOf(Instructions.InitializationBlocked)
|
||||
|
|
@ -697,26 +706,29 @@ describe("SessionRunnerLLM", () => {
|
|||
const { db } = yield* Database.Service
|
||||
yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "First" }), resume: false })
|
||||
|
||||
const exit = yield* scenario.run(function* () {
|
||||
yield* scenario.run(function* () {
|
||||
yield* (yield* scenario.llm.next()).respond.events()
|
||||
yield* session.wait(sessionID)
|
||||
|
||||
yield* events.publish(SessionEvent.Moved, {
|
||||
sessionID,
|
||||
location: Location.Ref.make({ directory: AbsolutePath.make("/moved") }),
|
||||
})
|
||||
expect(
|
||||
yield* db
|
||||
.select()
|
||||
.from(InstructionCheckpointTable)
|
||||
.where(eq(InstructionCheckpointTable.session_id, sessionID))
|
||||
.get(),
|
||||
).toBeUndefined()
|
||||
|
||||
yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Second" }), resume: false })
|
||||
return yield* session.resume(sessionID).pipe(Effect.exit)
|
||||
})
|
||||
|
||||
yield* events.publish(SessionEvent.Moved, {
|
||||
sessionID,
|
||||
location: Location.Ref.make({ directory: AbsolutePath.make("/moved") }),
|
||||
})
|
||||
expect(
|
||||
yield* db
|
||||
.select()
|
||||
.from(InstructionCheckpointTable)
|
||||
.where(eq(InstructionCheckpointTable.session_id, sessionID))
|
||||
.get(),
|
||||
).toBeUndefined()
|
||||
|
||||
yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Second" }), resume: false })
|
||||
const exit = yield* scenario
|
||||
.run(function* () {
|
||||
yield* scenario.llm.next()
|
||||
})
|
||||
.pipe(Effect.exit)
|
||||
|
||||
expect(Exit.isFailure(exit) && Cause.hasInterruptsOnly(exit.cause)).toBe(true)
|
||||
expect(yield* scenario.llm.requests).toHaveLength(1)
|
||||
expect(yield* SessionInput.hasPending(db, sessionID, "steer")).toBe(true)
|
||||
|
|
@ -760,22 +772,23 @@ describe("SessionRunnerLLM", () => {
|
|||
const { db } = yield* Database.Service
|
||||
yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "First" }), resume: false })
|
||||
yield* scenario.run(function* () {
|
||||
const first = yield* scenario.llm.next()
|
||||
yield* db
|
||||
.update(InstructionCheckpointTable)
|
||||
.set({ snapshot: { invalid: { value: "bad" } } })
|
||||
.where(eq(InstructionCheckpointTable.session_id, sessionID))
|
||||
.run()
|
||||
.pipe(Effect.orDie)
|
||||
yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Second" }), resume: false })
|
||||
yield* first.respond.events()
|
||||
yield* (yield* scenario.llm.next()).respond.events()
|
||||
})
|
||||
yield* db
|
||||
.update(InstructionCheckpointTable)
|
||||
.set({ snapshot: { invalid: { value: "bad" } } })
|
||||
.where(eq(InstructionCheckpointTable.session_id, sessionID))
|
||||
.run()
|
||||
.pipe(Effect.orDie)
|
||||
yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Second" }), resume: false })
|
||||
|
||||
const second = yield* scenario.llm.next()
|
||||
yield* scenario.run(function* () {
|
||||
const call = yield* scenario.llm.next()
|
||||
// Comparison state was lost, so every source re-announces as new.
|
||||
expect(second.request.system.map((part) => part.text)).toEqual([defaultSystem, "Initial context"])
|
||||
expect(second.request.messages.map((message) => message.role)).toEqual(["user", "system", "user"])
|
||||
expect(second.request.messages.at(1)?.content).toEqual([{ type: "text", text: "Initial context" }])
|
||||
yield* second.respond.events()
|
||||
expect(call.request.system.map((part) => part.text)).toEqual([defaultSystem, "Initial context"])
|
||||
expect(call.request.messages.map((message) => message.role)).toEqual(["user", "system", "user"])
|
||||
expect(call.request.messages.at(1)?.content).toEqual([{ type: "text", text: "Initial context" }])
|
||||
yield* call.respond.events()
|
||||
})
|
||||
|
||||
expect(yield* scenario.llm.requests).toHaveLength(2)
|
||||
|
|
@ -796,17 +809,19 @@ describe("SessionRunnerLLM", () => {
|
|||
yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "First" }), resume: false })
|
||||
|
||||
yield* scenario.run(function* () {
|
||||
const first = yield* scenario.llm.next()
|
||||
expect(first.request.system.map((part) => part.text)).toEqual([defaultSystem, "Initial context"])
|
||||
systemBaseline = "Changed context"
|
||||
yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Second" }), resume: false })
|
||||
yield* first.respond.events()
|
||||
const call = yield* scenario.llm.next()
|
||||
expect(call.request.system.map((part) => part.text)).toEqual([defaultSystem, "Initial context"])
|
||||
yield* call.respond.events()
|
||||
})
|
||||
systemBaseline = "Changed context"
|
||||
yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Second" }), resume: false })
|
||||
|
||||
const second = yield* scenario.llm.next()
|
||||
expect(second.request.system.map((part) => part.text)).toEqual([defaultSystem, "Initial context"])
|
||||
expect(second.request.messages.map((message) => message.role)).toEqual(["user", "system", "user"])
|
||||
expect(second.request.messages.at(1)?.content).toEqual([{ type: "text", text: "Changed context" }])
|
||||
yield* second.respond.events()
|
||||
yield* scenario.run(function* () {
|
||||
const call = yield* scenario.llm.next()
|
||||
expect(call.request.system.map((part) => part.text)).toEqual([defaultSystem, "Initial context"])
|
||||
expect(call.request.messages.map((message) => message.role)).toEqual(["user", "system", "user"])
|
||||
expect(call.request.messages.at(1)?.content).toEqual([{ type: "text", text: "Changed context" }])
|
||||
yield* call.respond.events()
|
||||
})
|
||||
|
||||
expect(yield* session.messages({ sessionID })).toHaveLength(3)
|
||||
|
|
@ -976,26 +991,22 @@ describe("SessionRunnerLLM", () => {
|
|||
yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "First" }), resume: false })
|
||||
|
||||
yield* scenario.run(function* () {
|
||||
const first = yield* scenario.llm.next()
|
||||
expect(first.request.system.map((part) => part.text)).toEqual([
|
||||
defaultSystem,
|
||||
"Initial context\n\nBuild skills",
|
||||
])
|
||||
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* first.respond.events()
|
||||
const call = yield* scenario.llm.next()
|
||||
expect(call.request.system.map((part) => part.text)).toEqual([defaultSystem, "Initial context\n\nBuild skills"])
|
||||
yield* call.respond.events()
|
||||
})
|
||||
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 })
|
||||
|
||||
const second = yield* scenario.llm.next()
|
||||
expect(second.request.system.map((part) => part.text)).toEqual([
|
||||
defaultSystem,
|
||||
"Initial context\n\nBuild skills",
|
||||
])
|
||||
expect(systemTexts(second.request)).toContainEqual(expect.stringContaining("Reviewer skills"))
|
||||
yield* second.respond.events()
|
||||
yield* scenario.run(function* () {
|
||||
const call = yield* scenario.llm.next()
|
||||
expect(call.request.system.map((part) => part.text)).toEqual([defaultSystem, "Initial context\n\nBuild skills"])
|
||||
expect(systemTexts(call.request)).toContainEqual(expect.stringContaining("Reviewer skills"))
|
||||
yield* call.respond.events()
|
||||
})
|
||||
}),
|
||||
)
|
||||
|
|
@ -1062,17 +1073,18 @@ describe("SessionRunnerLLM", () => {
|
|||
yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "First" }), resume: false })
|
||||
|
||||
yield* scenario.run(function* () {
|
||||
const first = yield* scenario.llm.next()
|
||||
systemRemoved = true
|
||||
yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Second" }), resume: false })
|
||||
yield* first.respond.events()
|
||||
yield* (yield* scenario.llm.next()).respond.events()
|
||||
})
|
||||
systemRemoved = true
|
||||
yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Second" }), resume: false })
|
||||
|
||||
const second = yield* scenario.llm.next()
|
||||
expect(second.request.messages.map((message) => message.role)).toEqual(["user", "system", "user"])
|
||||
expect(second.request.messages.at(1)?.content).toEqual([
|
||||
yield* scenario.run(function* () {
|
||||
const call = yield* scenario.llm.next()
|
||||
expect(call.request.messages.map((message) => message.role)).toEqual(["user", "system", "user"])
|
||||
expect(call.request.messages.at(1)?.content).toEqual([
|
||||
{ type: "text", text: "System context source removed: test/context" },
|
||||
])
|
||||
yield* second.respond.events()
|
||||
yield* call.respond.events()
|
||||
})
|
||||
|
||||
expect(yield* session.messages({ sessionID })).toHaveLength(3)
|
||||
|
|
@ -1088,21 +1100,23 @@ describe("SessionRunnerLLM", () => {
|
|||
yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "First" }), resume: false })
|
||||
|
||||
yield* scenario.run(function* () {
|
||||
const first = yield* scenario.llm.next()
|
||||
const call = yield* scenario.llm.next()
|
||||
// String values render verbatim inside the tagged block at baseline.
|
||||
expect(first.request.system.map((part) => part.text)).toEqual([
|
||||
expect(call.request.system.map((part) => part.text)).toEqual([
|
||||
defaultSystem,
|
||||
["Initial context", "", '<context key="deploy-target">', "production", "</context>"].join("\n"),
|
||||
])
|
||||
yield* call.respond.events()
|
||||
})
|
||||
|
||||
// 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* first.respond.events()
|
||||
// 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 })
|
||||
|
||||
const second = yield* scenario.llm.next()
|
||||
expect(second.request.messages.map((message) => message.role)).toEqual(["user", "system", "user"])
|
||||
expect(second.request.messages.at(1)?.content).toEqual([
|
||||
yield* scenario.run(function* () {
|
||||
const call = yield* scenario.llm.next()
|
||||
expect(call.request.messages.map((message) => message.role)).toEqual(["user", "system", "user"])
|
||||
expect(call.request.messages.at(1)?.content).toEqual([
|
||||
{
|
||||
type: "text",
|
||||
text: [
|
||||
|
|
@ -1115,27 +1129,27 @@ describe("SessionRunnerLLM", () => {
|
|||
].join("\n"),
|
||||
},
|
||||
])
|
||||
expect(yield* contextEntries.list(sessionID)).toEqual([
|
||||
{ key: "deploy-target", value: { region: "us-east-1" } },
|
||||
])
|
||||
yield* call.respond.events()
|
||||
})
|
||||
expect(yield* contextEntries.list(sessionID)).toEqual([{ key: "deploy-target", value: { region: "us-east-1" } }])
|
||||
|
||||
// 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* second.respond.events()
|
||||
// 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 })
|
||||
|
||||
const third = yield* scenario.llm.next()
|
||||
expect(third.request.messages.map((message) => message.role)).toEqual([
|
||||
yield* scenario.run(function* () {
|
||||
const call = yield* scenario.llm.next()
|
||||
expect(call.request.messages.map((message) => message.role)).toEqual([
|
||||
"user",
|
||||
"system",
|
||||
"user",
|
||||
"system",
|
||||
"user",
|
||||
])
|
||||
expect(third.request.messages.at(-2)?.content).toEqual([
|
||||
expect(call.request.messages.at(-2)?.content).toEqual([
|
||||
{ type: "text", text: 'The context under "deploy-target" no longer applies. Disregard it.' },
|
||||
])
|
||||
yield* third.respond.events()
|
||||
yield* call.respond.events()
|
||||
})
|
||||
|
||||
expect(yield* contextEntries.list(sessionID)).toEqual([])
|
||||
|
|
@ -1150,39 +1164,45 @@ describe("SessionRunnerLLM", () => {
|
|||
yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "First" }), resume: false })
|
||||
|
||||
yield* scenario.run(function* () {
|
||||
const first = yield* scenario.llm.next()
|
||||
expect(first.request.system.map((part) => part.text)).toEqual([defaultSystem, "Initial context"])
|
||||
systemBaseline = "Changed context"
|
||||
yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Second" }), resume: false })
|
||||
yield* first.respond.events()
|
||||
const call = yield* scenario.llm.next()
|
||||
expect(call.request.system.map((part) => part.text)).toEqual([defaultSystem, "Initial context"])
|
||||
yield* call.respond.events()
|
||||
})
|
||||
systemBaseline = "Changed context"
|
||||
yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Second" }), resume: false })
|
||||
|
||||
const second = yield* scenario.llm.next()
|
||||
expect(second.request.system.map((part) => part.text)).toEqual([defaultSystem, "Initial context"])
|
||||
expect(second.request.messages.map((message) => message.role)).toEqual(["user", "system", "user"])
|
||||
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* second.respond.events()
|
||||
yield* scenario.run(function* () {
|
||||
const call = yield* scenario.llm.next()
|
||||
expect(call.request.system.map((part) => part.text)).toEqual([defaultSystem, "Initial context"])
|
||||
expect(call.request.messages.map((message) => message.role)).toEqual(["user", "system", "user"])
|
||||
yield* call.respond.events()
|
||||
})
|
||||
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 })
|
||||
|
||||
const third = yield* scenario.llm.next()
|
||||
expect(third.request.system.map((part) => part.text)).toEqual([defaultSystem, "Initial context"])
|
||||
expect(third.request.messages.filter((message) => message.role === "system")).toHaveLength(2)
|
||||
expect((yield* session.context(sessionID)).map((message) => message.type)).toEqual([
|
||||
"user",
|
||||
"system",
|
||||
"user",
|
||||
"model-switched",
|
||||
"system",
|
||||
"user",
|
||||
])
|
||||
yield* replaySessionProjection(sessionID)
|
||||
expect(yield* session.messages({ sessionID })).toHaveLength(6)
|
||||
yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Fourth" }), resume: false })
|
||||
yield* third.respond.events()
|
||||
yield* scenario.run(function* () {
|
||||
const call = yield* scenario.llm.next()
|
||||
expect(call.request.system.map((part) => part.text)).toEqual([defaultSystem, "Initial context"])
|
||||
expect(call.request.messages.filter((message) => message.role === "system")).toHaveLength(2)
|
||||
yield* call.respond.events()
|
||||
})
|
||||
expect((yield* session.context(sessionID)).map((message) => message.type)).toEqual([
|
||||
"user",
|
||||
"system",
|
||||
"user",
|
||||
"model-switched",
|
||||
"system",
|
||||
"user",
|
||||
])
|
||||
yield* replaySessionProjection(sessionID)
|
||||
expect(yield* session.messages({ sessionID })).toHaveLength(6)
|
||||
yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Fourth" }), resume: false })
|
||||
|
||||
yield* scenario.run(function* () {
|
||||
yield* (yield* scenario.llm.next()).respond.events()
|
||||
})
|
||||
}),
|
||||
|
|
@ -1196,26 +1216,30 @@ describe("SessionRunnerLLM", () => {
|
|||
yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "First" }), resume: false })
|
||||
|
||||
yield* scenario.run(function* () {
|
||||
const first = yield* scenario.llm.next()
|
||||
expect(first.request.system.map((part) => part.text)).toEqual([defaultSystem, "Initial context"])
|
||||
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* first.respond.events()
|
||||
const call = yield* scenario.llm.next()
|
||||
expect(call.request.system.map((part) => part.text)).toEqual([defaultSystem, "Initial context"])
|
||||
yield* call.respond.events()
|
||||
})
|
||||
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 })
|
||||
|
||||
const second = yield* scenario.llm.next()
|
||||
expect(second.request.system.map((part) => part.text)).toEqual([defaultSystem, "Initial context"])
|
||||
systemUnavailable = false
|
||||
systemBaseline = "Replacement context"
|
||||
yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Third" }), resume: false })
|
||||
yield* second.respond.events()
|
||||
yield* scenario.run(function* () {
|
||||
const call = yield* scenario.llm.next()
|
||||
expect(call.request.system.map((part) => part.text)).toEqual([defaultSystem, "Initial context"])
|
||||
yield* call.respond.events()
|
||||
})
|
||||
systemUnavailable = false
|
||||
systemBaseline = "Replacement context"
|
||||
yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Third" }), resume: false })
|
||||
|
||||
const third = yield* scenario.llm.next()
|
||||
expect(third.request.system.map((part) => part.text)).toEqual([defaultSystem, "Initial context"])
|
||||
yield* third.respond.events()
|
||||
yield* scenario.run(function* () {
|
||||
const call = yield* scenario.llm.next()
|
||||
expect(call.request.system.map((part) => part.text)).toEqual([defaultSystem, "Initial context"])
|
||||
yield* call.respond.events()
|
||||
})
|
||||
}),
|
||||
)
|
||||
|
|
@ -1228,28 +1252,32 @@ describe("SessionRunnerLLM", () => {
|
|||
yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "First" }), resume: false })
|
||||
|
||||
yield* scenario.run(function* () {
|
||||
const first = yield* scenario.llm.next()
|
||||
expect(first.request.system.map((part) => part.text)).toEqual([defaultSystem, "Initial context"])
|
||||
yield* events.publish(SessionEvent.Compaction.Started, {
|
||||
sessionID,
|
||||
reason: "manual",
|
||||
})
|
||||
yield* events.publish(SessionEvent.Compaction.Ended, {
|
||||
sessionID,
|
||||
reason: "manual",
|
||||
text: "summary",
|
||||
recent: "",
|
||||
})
|
||||
systemBaseline = "Replacement context"
|
||||
yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Second" }), resume: false })
|
||||
yield* first.respond.events()
|
||||
const call = yield* scenario.llm.next()
|
||||
expect(call.request.system.map((part) => part.text)).toEqual([defaultSystem, "Initial context"])
|
||||
yield* call.respond.events()
|
||||
})
|
||||
yield* events.publish(SessionEvent.Compaction.Started, {
|
||||
sessionID,
|
||||
reason: "manual",
|
||||
})
|
||||
yield* events.publish(SessionEvent.Compaction.Ended, {
|
||||
sessionID,
|
||||
reason: "manual",
|
||||
text: "summary",
|
||||
recent: "",
|
||||
})
|
||||
systemBaseline = "Replacement context"
|
||||
yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Second" }), resume: false })
|
||||
|
||||
const second = yield* scenario.llm.next()
|
||||
expect(second.request.system.map((part) => part.text)).toEqual([defaultSystem, "Replacement context"])
|
||||
yield* replaySessionProjection(sessionID)
|
||||
yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Third" }), resume: false })
|
||||
yield* second.respond.events()
|
||||
yield* scenario.run(function* () {
|
||||
const call = yield* scenario.llm.next()
|
||||
expect(call.request.system.map((part) => part.text)).toEqual([defaultSystem, "Replacement context"])
|
||||
yield* call.respond.events()
|
||||
})
|
||||
yield* replaySessionProjection(sessionID)
|
||||
yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Third" }), resume: false })
|
||||
|
||||
yield* scenario.run(function* () {
|
||||
yield* (yield* scenario.llm.next()).respond.events()
|
||||
})
|
||||
}),
|
||||
|
|
@ -1358,15 +1386,16 @@ describe("SessionRunnerLLM", () => {
|
|||
})
|
||||
|
||||
yield* scenario.run(function* () {
|
||||
const earlier = yield* scenario.llm.next()
|
||||
currentModel = compactModel
|
||||
yield* session.prompt({
|
||||
sessionID,
|
||||
prompt: PromptInput.Prompt.make({ text: "Recent exact request ".repeat(180) }),
|
||||
resume: false,
|
||||
})
|
||||
yield* earlier.respond.text("Earlier answer", { id: "text-first" })
|
||||
yield* (yield* scenario.llm.next()).respond.text("Earlier answer", { id: "text-first" })
|
||||
})
|
||||
currentModel = compactModel
|
||||
yield* session.prompt({
|
||||
sessionID,
|
||||
prompt: PromptInput.Prompt.make({ text: "Recent exact request ".repeat(180) }),
|
||||
resume: false,
|
||||
})
|
||||
|
||||
yield* scenario.run(function* () {
|
||||
const summary = yield* scenario.llm.next()
|
||||
expect(userTexts(summary.request)[0]).toContain("## Objective")
|
||||
yield* summary.respond.text("## Objective\n- Preserve the task", { id: "text-summary" })
|
||||
|
|
@ -1375,13 +1404,15 @@ describe("SessionRunnerLLM", () => {
|
|||
expect(userTexts(continued.request)).toHaveLength(1)
|
||||
expect(userTexts(continued.request)[0]).toContain("<summary>\n## Objective\n- Preserve the task\n</summary>")
|
||||
expect(userTexts(continued.request)[0]).toContain(`[User]: ${"Recent exact request ".repeat(180)}`)
|
||||
yield* session.prompt({
|
||||
sessionID,
|
||||
prompt: PromptInput.Prompt.make({ text: "Newest exact request ".repeat(180) }),
|
||||
resume: false,
|
||||
})
|
||||
yield* continued.respond.text("Continued", { id: "text-final" })
|
||||
})
|
||||
yield* session.prompt({
|
||||
sessionID,
|
||||
prompt: PromptInput.Prompt.make({ text: "Newest exact request ".repeat(180) }),
|
||||
resume: false,
|
||||
})
|
||||
|
||||
yield* scenario.run(function* () {
|
||||
const updatedSummary = yield* scenario.llm.next()
|
||||
expect(userTexts(updatedSummary.request)[0]).toContain(
|
||||
"<previous-summary>\n## Objective\n- Preserve the task\n</previous-summary>",
|
||||
|
|
@ -1402,13 +1433,9 @@ describe("SessionRunnerLLM", () => {
|
|||
|
||||
scenarioIt("forces one compaction and retries after provider context overflow", (scenario) =>
|
||||
Effect.gen(function* () {
|
||||
const session = yield* setupOverflowRecovery
|
||||
const session = yield* setupOverflowRecovery(scenario)
|
||||
yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Continue" }), resume: false })
|
||||
yield* scenario.run(function* () {
|
||||
const earlier = yield* scenario.llm.next()
|
||||
currentModel = recoveryModel
|
||||
yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Continue" }), resume: false })
|
||||
yield* earlier.respond.text("Earlier answer", { id: "text-earlier" })
|
||||
|
||||
yield* (yield* scenario.llm.next()).respond.events(
|
||||
LLMEvent.stepStart({ index: 0 }),
|
||||
LLMEvent.providerError({ message: "prompt too long", classification: "context-overflow" }),
|
||||
|
|
@ -1438,7 +1465,8 @@ describe("SessionRunnerLLM", () => {
|
|||
|
||||
scenarioIt("persists a second context overflow after one recovery", (scenario) =>
|
||||
Effect.gen(function* () {
|
||||
const session = yield* setupOverflowRecovery
|
||||
const session = yield* setupOverflowRecovery(scenario)
|
||||
yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Continue" }), resume: false })
|
||||
const overflow = () => [
|
||||
LLMEvent.stepStart({ index: 0 }),
|
||||
LLMEvent.providerError({ message: "prompt too long", classification: "context-overflow" }),
|
||||
|
|
@ -1446,15 +1474,6 @@ describe("SessionRunnerLLM", () => {
|
|||
expect(
|
||||
(yield* scenario
|
||||
.run(function* () {
|
||||
const earlier = yield* scenario.llm.next()
|
||||
currentModel = recoveryModel
|
||||
yield* session.prompt({
|
||||
sessionID,
|
||||
prompt: PromptInput.Prompt.make({ text: "Continue" }),
|
||||
resume: false,
|
||||
})
|
||||
yield* earlier.respond.text("Earlier answer", { id: "text-earlier" })
|
||||
|
||||
yield* (yield* scenario.llm.next()).respond.events(...overflow())
|
||||
yield* (yield* scenario.llm.next()).respond.text("## Objective\n- Recover once", {
|
||||
id: "text-summary",
|
||||
|
|
@ -1474,13 +1493,9 @@ describe("SessionRunnerLLM", () => {
|
|||
|
||||
scenarioIt("recovers once from a raw context overflow failure", (scenario) =>
|
||||
Effect.gen(function* () {
|
||||
const session = yield* setupOverflowRecovery
|
||||
const session = yield* setupOverflowRecovery(scenario)
|
||||
yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Continue" }), resume: false })
|
||||
yield* scenario.run(function* () {
|
||||
const earlier = yield* scenario.llm.next()
|
||||
currentModel = recoveryModel
|
||||
yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Continue" }), resume: false })
|
||||
yield* earlier.respond.text("Earlier answer", { id: "text-earlier" })
|
||||
|
||||
yield* (yield* scenario.llm.next()).respond.fail(
|
||||
new LLMError({
|
||||
module: "test",
|
||||
|
|
@ -1507,19 +1522,11 @@ describe("SessionRunnerLLM", () => {
|
|||
|
||||
scenarioIt("publishes the original overflow when recovery summarization fails", (scenario) =>
|
||||
Effect.gen(function* () {
|
||||
const session = yield* setupOverflowRecovery
|
||||
const session = yield* setupOverflowRecovery(scenario)
|
||||
yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Continue" }), resume: false })
|
||||
expect(
|
||||
(yield* scenario
|
||||
.run(function* () {
|
||||
const earlier = yield* scenario.llm.next()
|
||||
currentModel = recoveryModel
|
||||
yield* session.prompt({
|
||||
sessionID,
|
||||
prompt: PromptInput.Prompt.make({ text: "Continue" }),
|
||||
resume: false,
|
||||
})
|
||||
yield* earlier.respond.text("Earlier answer", { id: "text-earlier" })
|
||||
|
||||
yield* (yield* scenario.llm.next()).respond.events(
|
||||
LLMEvent.providerError({ message: "prompt too long", classification: "context-overflow" }),
|
||||
)
|
||||
|
|
@ -1542,18 +1549,10 @@ describe("SessionRunnerLLM", () => {
|
|||
|
||||
scenarioIt("interrupts overflow recovery while the summary provider is running", (scenario) =>
|
||||
Effect.gen(function* () {
|
||||
const session = yield* setupOverflowRecovery
|
||||
const session = yield* setupOverflowRecovery(scenario)
|
||||
yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Continue" }), resume: false })
|
||||
const exit = yield* scenario
|
||||
.run(function* () {
|
||||
const earlier = yield* scenario.llm.next()
|
||||
currentModel = recoveryModel
|
||||
yield* session.prompt({
|
||||
sessionID,
|
||||
prompt: PromptInput.Prompt.make({ text: "Continue" }),
|
||||
resume: false,
|
||||
})
|
||||
yield* earlier.respond.text("Earlier answer", { id: "text-earlier" })
|
||||
|
||||
yield* (yield* scenario.llm.next()).respond.events(
|
||||
LLMEvent.providerError({ message: "prompt too long", classification: "context-overflow" }),
|
||||
)
|
||||
|
|
@ -1576,31 +1575,33 @@ describe("SessionRunnerLLM", () => {
|
|||
yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "First" }), resume: false })
|
||||
|
||||
yield* scenario.run(function* () {
|
||||
const first = yield* scenario.llm.next()
|
||||
systemBaseline = "Changed context"
|
||||
yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Second" }), resume: false })
|
||||
yield* first.respond.events()
|
||||
yield* (yield* scenario.llm.next()).respond.events()
|
||||
})
|
||||
systemBaseline = "Changed context"
|
||||
yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Second" }), resume: false })
|
||||
|
||||
const second = yield* scenario.llm.next()
|
||||
yield* events.publish(SessionEvent.Compaction.Started, {
|
||||
sessionID,
|
||||
reason: "manual",
|
||||
})
|
||||
yield* events.publish(SessionEvent.Compaction.Ended, {
|
||||
sessionID,
|
||||
reason: "manual",
|
||||
text: "summary",
|
||||
recent: "",
|
||||
})
|
||||
systemUnavailable = true
|
||||
yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Third" }), resume: false })
|
||||
yield* second.respond.events()
|
||||
yield* scenario.run(function* () {
|
||||
yield* (yield* scenario.llm.next()).respond.events()
|
||||
})
|
||||
yield* events.publish(SessionEvent.Compaction.Started, {
|
||||
sessionID,
|
||||
reason: "manual",
|
||||
})
|
||||
yield* events.publish(SessionEvent.Compaction.Ended, {
|
||||
sessionID,
|
||||
reason: "manual",
|
||||
text: "summary",
|
||||
recent: "",
|
||||
})
|
||||
systemUnavailable = true
|
||||
yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Third" }), resume: false })
|
||||
|
||||
const third = yield* scenario.llm.next()
|
||||
yield* scenario.run(function* () {
|
||||
const call = yield* scenario.llm.next()
|
||||
// The rebaseline proceeds while the source is unobservable, restating the model's belief.
|
||||
expect(third.request.system.map((part) => part.text)).toEqual([defaultSystem, "Changed context"])
|
||||
expect(systemTexts(third.request)).not.toContain("Changed context")
|
||||
yield* third.respond.events()
|
||||
expect(call.request.system.map((part) => part.text)).toEqual([defaultSystem, "Changed context"])
|
||||
expect(systemTexts(call.request)).not.toContain("Changed context")
|
||||
yield* call.respond.events()
|
||||
})
|
||||
}),
|
||||
)
|
||||
|
|
@ -1805,69 +1806,56 @@ describe("SessionRunnerLLM", () => {
|
|||
const session = yield* SessionV2.Service
|
||||
yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Think first" }), resume: false })
|
||||
|
||||
const streamed = yield* Deferred.make<void>()
|
||||
const release = yield* Deferred.make<void>()
|
||||
yield* scenario.run(function* () {
|
||||
const first = yield* scenario.llm.next()
|
||||
yield* first.respond.stream(
|
||||
Stream.concat(
|
||||
Stream.fromIterable([
|
||||
LLMEvent.stepStart({ index: 0 }),
|
||||
LLMEvent.reasoningStart({ id: "reasoning-anthropic" }),
|
||||
LLMEvent.reasoningDelta({ id: "reasoning-anthropic", text: "Signed thought" }),
|
||||
LLMEvent.reasoningEnd({
|
||||
id: "reasoning-anthropic",
|
||||
providerMetadata: { fake: { signature: "sig_1" }, anthropic: { ignored: true } },
|
||||
}),
|
||||
LLMEvent.reasoningStart({
|
||||
id: "reasoning-openai",
|
||||
providerMetadata: {
|
||||
fake: { itemId: "rs_1", reasoningEncryptedContent: null },
|
||||
openai: { ignored: true },
|
||||
},
|
||||
}),
|
||||
LLMEvent.reasoningDelta({ id: "reasoning-openai", text: "Encrypted thought" }),
|
||||
LLMEvent.reasoningEnd({
|
||||
id: "reasoning-openai",
|
||||
providerMetadata: {
|
||||
fake: { itemId: "rs_1", reasoningEncryptedContent: "encrypted-state" },
|
||||
openai: { ignored: true },
|
||||
},
|
||||
}),
|
||||
LLMEvent.stepFinish({ index: 0, reason: "stop" }),
|
||||
LLMEvent.finish({ reason: "stop" }),
|
||||
]),
|
||||
Stream.unwrap(
|
||||
Deferred.succeed(streamed, undefined).pipe(
|
||||
Effect.andThen(Deferred.await(release)),
|
||||
Effect.as(Stream.empty),
|
||||
),
|
||||
),
|
||||
),
|
||||
yield* (yield* scenario.llm.next()).respond.events(
|
||||
LLMEvent.stepStart({ index: 0 }),
|
||||
LLMEvent.reasoningStart({ id: "reasoning-anthropic" }),
|
||||
LLMEvent.reasoningDelta({ id: "reasoning-anthropic", text: "Signed thought" }),
|
||||
LLMEvent.reasoningEnd({
|
||||
id: "reasoning-anthropic",
|
||||
providerMetadata: { fake: { signature: "sig_1" }, anthropic: { ignored: true } },
|
||||
}),
|
||||
LLMEvent.reasoningStart({
|
||||
id: "reasoning-openai",
|
||||
providerMetadata: {
|
||||
fake: { itemId: "rs_1", reasoningEncryptedContent: null },
|
||||
openai: { ignored: true },
|
||||
},
|
||||
}),
|
||||
LLMEvent.reasoningDelta({ id: "reasoning-openai", text: "Encrypted thought" }),
|
||||
LLMEvent.reasoningEnd({
|
||||
id: "reasoning-openai",
|
||||
providerMetadata: {
|
||||
fake: { itemId: "rs_1", reasoningEncryptedContent: "encrypted-state" },
|
||||
openai: { ignored: true },
|
||||
},
|
||||
}),
|
||||
LLMEvent.stepFinish({ index: 0, reason: "stop" }),
|
||||
LLMEvent.finish({ reason: "stop" }),
|
||||
)
|
||||
yield* Deferred.await(streamed)
|
||||
yield* replaySessionProjection(sessionID)
|
||||
})
|
||||
yield* replaySessionProjection(sessionID)
|
||||
|
||||
expect(yield* session.context(sessionID)).toMatchObject([
|
||||
{ type: "user", text: "Think first" },
|
||||
{
|
||||
type: "assistant",
|
||||
content: [
|
||||
{ type: "reasoning", text: "Signed thought", state: { signature: "sig_1" } },
|
||||
{
|
||||
type: "reasoning",
|
||||
text: "Encrypted thought",
|
||||
state: { itemId: "rs_1", reasoningEncryptedContent: "encrypted-state" },
|
||||
},
|
||||
],
|
||||
},
|
||||
])
|
||||
expect(yield* session.context(sessionID)).toMatchObject([
|
||||
{ type: "user", text: "Think first" },
|
||||
{
|
||||
type: "assistant",
|
||||
content: [
|
||||
{ type: "reasoning", text: "Signed thought", state: { signature: "sig_1" } },
|
||||
{
|
||||
type: "reasoning",
|
||||
text: "Encrypted thought",
|
||||
state: { itemId: "rs_1", reasoningEncryptedContent: "encrypted-state" },
|
||||
},
|
||||
],
|
||||
},
|
||||
])
|
||||
|
||||
yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Continue" }), resume: false })
|
||||
yield* Deferred.succeed(release, undefined)
|
||||
yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Continue" }), resume: false })
|
||||
|
||||
const second = yield* scenario.llm.next()
|
||||
expect(second.request.messages[1]?.content).toEqual([
|
||||
yield* scenario.run(function* () {
|
||||
const call = yield* scenario.llm.next()
|
||||
expect(call.request.messages[1]?.content).toEqual([
|
||||
{ type: "reasoning", text: "Signed thought", providerMetadata: { fake: { signature: "sig_1" } } },
|
||||
{
|
||||
type: "reasoning",
|
||||
|
|
@ -1875,8 +1863,10 @@ describe("SessionRunnerLLM", () => {
|
|||
providerMetadata: { fake: { itemId: "rs_1", reasoningEncryptedContent: "encrypted-state" } },
|
||||
},
|
||||
])
|
||||
yield* second.respond.events()
|
||||
yield* call.respond.events()
|
||||
})
|
||||
|
||||
expect(yield* scenario.llm.requests).toHaveLength(2)
|
||||
}),
|
||||
)
|
||||
|
||||
|
|
@ -1886,48 +1876,35 @@ describe("SessionRunnerLLM", () => {
|
|||
const session = yield* SessionV2.Service
|
||||
yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Search first" }), resume: false })
|
||||
|
||||
const streamed = yield* Deferred.make<void>()
|
||||
const release = yield* Deferred.make<void>()
|
||||
yield* scenario.run(function* () {
|
||||
const first = yield* scenario.llm.next()
|
||||
yield* first.respond.stream(
|
||||
Stream.concat(
|
||||
Stream.fromIterable([
|
||||
LLMEvent.stepStart({ index: 0 }),
|
||||
LLMEvent.toolCall({
|
||||
id: "hosted-search",
|
||||
name: "web_search",
|
||||
input: { query: "Effect" },
|
||||
providerExecuted: true,
|
||||
providerMetadata: { fake: { itemId: "hosted-search" }, openai: { ignored: true } },
|
||||
}),
|
||||
LLMEvent.toolResult({
|
||||
id: "hosted-search",
|
||||
name: "web_search",
|
||||
result: { type: "json", value: [{ title: "Effect" }] },
|
||||
providerExecuted: true,
|
||||
providerMetadata: { fake: { blockType: "web_search_tool_result" }, anthropic: { ignored: true } },
|
||||
}),
|
||||
LLMEvent.stepFinish({ index: 0, reason: "stop" }),
|
||||
LLMEvent.finish({ reason: "stop" }),
|
||||
]),
|
||||
Stream.unwrap(
|
||||
Deferred.succeed(streamed, undefined).pipe(
|
||||
Effect.andThen(Deferred.await(release)),
|
||||
Effect.as(Stream.empty),
|
||||
),
|
||||
),
|
||||
),
|
||||
yield* (yield* scenario.llm.next()).respond.events(
|
||||
LLMEvent.stepStart({ index: 0 }),
|
||||
LLMEvent.toolCall({
|
||||
id: "hosted-search",
|
||||
name: "web_search",
|
||||
input: { query: "Effect" },
|
||||
providerExecuted: true,
|
||||
providerMetadata: { fake: { itemId: "hosted-search" }, openai: { ignored: true } },
|
||||
}),
|
||||
LLMEvent.toolResult({
|
||||
id: "hosted-search",
|
||||
name: "web_search",
|
||||
result: { type: "json", value: [{ title: "Effect" }] },
|
||||
providerExecuted: true,
|
||||
providerMetadata: { fake: { blockType: "web_search_tool_result" }, anthropic: { ignored: true } },
|
||||
}),
|
||||
LLMEvent.stepFinish({ index: 0, reason: "stop" }),
|
||||
LLMEvent.finish({ reason: "stop" }),
|
||||
)
|
||||
yield* Deferred.await(streamed)
|
||||
yield* replaySessionProjection(sessionID)
|
||||
})
|
||||
yield* replaySessionProjection(sessionID)
|
||||
|
||||
yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Continue" }), resume: false })
|
||||
yield* Deferred.succeed(release, undefined)
|
||||
yield* session.prompt({ sessionID, prompt: PromptInput.Prompt.make({ text: "Continue" }), resume: false })
|
||||
|
||||
const second = yield* scenario.llm.next()
|
||||
expect(second.request.messages.map((message) => message.role)).toEqual(["user", "assistant", "user"])
|
||||
expect(second.request.messages[1]?.content).toMatchObject([
|
||||
yield* scenario.run(function* () {
|
||||
const call = yield* scenario.llm.next()
|
||||
expect(call.request.messages.map((message) => message.role)).toEqual(["user", "assistant", "user"])
|
||||
expect(call.request.messages[1]?.content).toMatchObject([
|
||||
{
|
||||
type: "tool-call",
|
||||
id: "hosted-search",
|
||||
|
|
@ -1945,8 +1922,10 @@ describe("SessionRunnerLLM", () => {
|
|||
providerMetadata: { fake: { blockType: "web_search_tool_result" } },
|
||||
},
|
||||
])
|
||||
yield* second.respond.events()
|
||||
yield* call.respond.events()
|
||||
})
|
||||
|
||||
expect(yield* scenario.llm.requests).toHaveLength(2)
|
||||
}),
|
||||
)
|
||||
|
||||
|
|
@ -2767,7 +2746,6 @@ describe("SessionRunnerLLM", () => {
|
|||
Exit.Exit<void, SessionRunner.RunError | SessionV2.NotFoundError>,
|
||||
]
|
||||
>()
|
||||
const retry = yield* Deferred.make<void>()
|
||||
const scenario = yield* RunnerScenario.make(() =>
|
||||
SessionV2.Service.use((session) =>
|
||||
Effect.gen(function* () {
|
||||
|
|
@ -2775,8 +2753,6 @@ describe("SessionRunnerLLM", () => {
|
|||
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)
|
||||
}),
|
||||
),
|
||||
)
|
||||
|
|
@ -2798,8 +2774,9 @@ describe("SessionRunnerLLM", () => {
|
|||
yield* first.respond.fail(failure)
|
||||
const [firstExit, secondExit] = yield* Deferred.await(joined)
|
||||
expect(secondExit).toEqual(firstExit)
|
||||
})
|
||||
|
||||
yield* Deferred.succeed(retry, undefined)
|
||||
yield* scenario.run(function* () {
|
||||
yield* (yield* scenario.llm.next()).respond.events()
|
||||
})
|
||||
|
||||
|
|
@ -3254,6 +3231,14 @@ describe("SessionRunnerLLM", () => {
|
|||
{ type: "user", text: "Interrupt blocked tool" },
|
||||
{ type: "assistant", content: [{ type: "tool", id: "call-before-interrupt", state: { status: "error" } }] },
|
||||
])
|
||||
|
||||
yield* scenario.run(function* () {
|
||||
const call = yield* scenario.llm.next()
|
||||
expect(call.request.messages.map((message) => message.role)).toEqual(["user", "assistant", "tool"])
|
||||
yield* call.respond.events()
|
||||
})
|
||||
|
||||
expect(yield* scenario.llm.requests).toHaveLength(2)
|
||||
}),
|
||||
)
|
||||
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue