fix(core): preserve queue after provider failure
This commit is contained in:
parent
dc468bdcfd
commit
24c672e89e
5 changed files with 51 additions and 3 deletions
|
|
@ -308,7 +308,10 @@ export const layer = Layer.effect(
|
|||
yield* withPublication(publisher.failUnsettledTools("Provider did not return a tool result", true))
|
||||
if (stream._tag === "Failure") return yield* Effect.failCause(stream.cause)
|
||||
if (settled._tag === "Failure") return yield* Effect.failCause(settled.cause)
|
||||
return { needsContinuation: !publisher.hasProviderError() && needsContinuation, step: currentStep }
|
||||
return {
|
||||
outcome: publisher.hasProviderError() ? "failed" : needsContinuation ? "continue" : "complete",
|
||||
step: currentStep,
|
||||
} as const
|
||||
}),
|
||||
)
|
||||
}, Effect.scoped)
|
||||
|
|
@ -316,7 +319,10 @@ export const layer = Layer.effect(
|
|||
sessionID: SessionSchema.ID,
|
||||
promotion: SessionInput.Delivery | undefined,
|
||||
step: number,
|
||||
) => Effect.Effect<{ readonly needsContinuation: boolean; readonly step: number }, RunError>
|
||||
) => Effect.Effect<
|
||||
{ readonly outcome: "continue" | "complete" | "failed"; readonly step: number },
|
||||
RunError
|
||||
>
|
||||
|
||||
const runAfterOverflowCompaction: RunTurn = Effect.fnUntraced(function* (sessionID, promotion, step) {
|
||||
return yield* runTurnAttempt(sessionID, promotion, step).pipe(
|
||||
|
|
@ -361,7 +367,8 @@ export const layer = Layer.effect(
|
|||
let step = 1
|
||||
while (needsContinuation) {
|
||||
const result = yield* runTurn(input.sessionID, promotion, step)
|
||||
needsContinuation = result.needsContinuation
|
||||
if (result.outcome === "failed") return
|
||||
needsContinuation = result.outcome === "continue"
|
||||
step = result.step + 1
|
||||
promotion = "steer"
|
||||
if (!needsContinuation) needsContinuation = yield* SessionInput.hasPending(db, input.sessionID, "steer")
|
||||
|
|
|
|||
|
|
@ -2927,6 +2927,43 @@ describe("SessionRunnerLLM", () => {
|
|||
}),
|
||||
)
|
||||
|
||||
it.effect("leaves queued input pending after a terminal provider error", () =>
|
||||
Effect.gen(function* () {
|
||||
yield* setup
|
||||
const session = yield* SessionV2.Service
|
||||
const { db } = yield* Database.Service
|
||||
yield* session.prompt({ sessionID, prompt: new Prompt({ text: "Fail first" }), resume: false })
|
||||
yield* session.prompt({
|
||||
sessionID,
|
||||
prompt: new Prompt({ text: "Run after explicit resume" }),
|
||||
delivery: "queue",
|
||||
resume: false,
|
||||
})
|
||||
|
||||
requests.length = 0
|
||||
responses = [
|
||||
[LLMEvent.stepStart({ index: 0 }), LLMEvent.providerError({ message: "Provider unavailable" })],
|
||||
[
|
||||
LLMEvent.stepStart({ index: 0 }),
|
||||
LLMEvent.stepFinish({ index: 0, reason: "stop" }),
|
||||
LLMEvent.finish({ reason: "stop" }),
|
||||
],
|
||||
]
|
||||
|
||||
yield* session.resume(sessionID)
|
||||
|
||||
expect(requests).toHaveLength(1)
|
||||
expect(yield* SessionInput.hasPending(db, sessionID, "queue")).toBe(true)
|
||||
expect(userTexts(requests[0]!)).toEqual(["Fail first"])
|
||||
|
||||
yield* session.resume(sessionID)
|
||||
|
||||
expect(requests).toHaveLength(2)
|
||||
expect(yield* SessionInput.hasPending(db, sessionID, "queue")).toBe(false)
|
||||
expect(userTexts(requests[1]!)).toEqual(["Fail first", "Run after explicit resume"])
|
||||
}),
|
||||
)
|
||||
|
||||
it.effect("projects provider errors emitted before assistant step start", () =>
|
||||
Effect.gen(function* () {
|
||||
yield* setup
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue