diff --git a/packages/core/src/session/runner/llm.ts b/packages/core/src/session/runner/llm.ts index c895558e02..704bf2a4bd 100644 --- a/packages/core/src/session/runner/llm.ts +++ b/packages/core/src/session/runner/llm.ts @@ -108,7 +108,7 @@ export const layer = Layer.effect( while (openActivity) { let needsContinuation = true for (let step = 0; step < MAX_STEPS; step++) { - needsContinuation = yield* runTurn(input.sessionID, promotion) + needsContinuation = yield* runTurn({ sessionID: input.sessionID, delivery: promotion }) promotion = "steer" if (!needsContinuation) needsContinuation = yield* SessionInput.hasPending(db, input.sessionID, "steer") if (!needsContinuation) break diff --git a/packages/core/src/session/runner/run-turn.ts b/packages/core/src/session/runner/run-turn.ts index cf9f5b8b05..9790a90764 100644 --- a/packages/core/src/session/runner/run-turn.ts +++ b/packages/core/src/session/runner/run-turn.ts @@ -44,10 +44,12 @@ import { SessionRunnerModel } from "./model" import { createLLMEventPublisher } from "./publish-llm-event" import { toLLMMessages } from "./to-llm-message" -export type Run = ( - sessionID: SessionSchema.ID, - promotion: SessionInput.Delivery | undefined, -) => Effect.Effect +export interface Input { + readonly sessionID: SessionSchema.ID + readonly delivery?: SessionInput.Delivery +} + +export type Run = (input: Input) => Effect.Effect const AttemptResult = Schema.TaggedUnion({ Complete: { needsContinuation: Schema.Boolean }, @@ -96,11 +98,11 @@ export const make = Effect.gen(function* () { * Once input is promoted, retries load it from history instead of promoting * again. This matters for queued input because promotion opens the next item. */ - const prepareTurn = Effect.fn("SessionRunner.prepareTurn")(function* ( + const buildRequest = Effect.fn("SessionRunner.buildRequest")(function* ( sessionID: SessionSchema.ID, - promotion: SessionInput.Delivery | undefined, + delivery: SessionInput.Delivery | undefined, ) { - let pendingPromotion = promotion + let pendingDelivery = delivery while (true) { const session = yield* getSession(sessionID) if (session.location.directory !== location.directory || session.location.workspaceID !== location.workspaceID) @@ -110,14 +112,14 @@ export const make = Effect.gen(function* () { SessionContextEpoch.initialize(db, loadSystemContext(agent), session.id, session.location, agent.id), ) if (initialized === stale) continue - if (pendingPromotion) { + if (pendingDelivery) { const cutoff = yield* SessionInput.latestSeq(db, session.id) - if (pendingPromotion === "steer") yield* SessionInput.promoteSteers(db, events, session.id, cutoff) - if (pendingPromotion === "queue") { + if (pendingDelivery === "steer") yield* SessionInput.promoteSteers(db, events, session.id, cutoff) + if (pendingDelivery === "queue") { yield* SessionInput.promoteNextQueued(db, events, session.id) yield* SessionInput.promoteSteers(db, events, session.id, cutoff) } - pendingPromotion = undefined + pendingDelivery = undefined } const prepared = initialized ?? @@ -153,13 +155,20 @@ export const make = Effect.gen(function* () { } }) - type RequestSnapshot = Effect.Success> + type RequestSnapshot = Effect.Success> /** - * Provider events and tool results can arrive concurrently. They share one - * permit so their durable Session events are written in order. + * Reads one model response and finishes every tool it starts. + * + * Provider events and tool results share one permit so durable events stay in + * order. Tool calls are recorded before their side effects begin. A pre-output + * overflow is held back so successful compaction does not leave a terminal + * error in Session history. */ - const startTurn = Effect.fnUntraced(function* (prepared: RequestSnapshot) { + const streamAndSettle = Effect.fn("SessionRunner.streamAndSettle")(function* ( + prepared: RequestSnapshot, + recoverOverflow?: typeof compaction.compactAfterOverflow, + ) { const publisher = createLLMEventPublisher(events, { sessionID: prepared.session.id, agent: prepared.agent.id, @@ -170,41 +179,23 @@ export const make = Effect.gen(function* () { }, }) const withPublication = Semaphore.makeUnsafe(1).withPermit - return { - publisher, - withPublication, - publish: (event: LLMEvent, outputPaths: ReadonlyArray = []) => - withPublication(publisher.publish(event, outputPaths)), - toolFibers: yield* FiberSet.make(), - needsContinuation: false, - overflowFailure: undefined as ProviderErrorEvent | undefined, - } - }) - - type ActiveTurn = Effect.Success> - - /** - * Reads one model response. A tool call is recorded before its side effect - * starts. An overflow error is held back briefly so successful compaction does - * not leave a terminal error in Session history. - */ - const consumeProvider = (prepared: RequestSnapshot, runtime: ActiveTurn) => - llm.stream(prepared.request).pipe( + const publish = (event: LLMEvent, outputPaths: ReadonlyArray = []) => + withPublication(publisher.publish(event, outputPaths)) + const toolFibers = yield* FiberSet.make() + let needsContinuation = false + let overflowFailure: ProviderErrorEvent | undefined + const providerStream = llm.stream(prepared.request).pipe( Stream.runForEach((event) => Effect.gen(function* () { - if (runtime.overflowFailure || runtime.publisher.hasProviderError()) return - if ( - LLMEvent.is.providerError(event) && - isContextOverflowFailure(event) && - !runtime.publisher.hasAssistantStarted() - ) { - runtime.overflowFailure = event + if (overflowFailure || publisher.hasProviderError()) return + if (LLMEvent.is.providerError(event) && isContextOverflowFailure(event) && !publisher.hasAssistantStarted()) { + overflowFailure = event return } - yield* runtime.publish(event) + yield* publish(event) if (event.type !== "tool-call" || event.providerExecuted) return - runtime.needsContinuation = true - const assistantMessageID = yield* runtime.publisher.assistantMessageID(event.id) + needsContinuation = true + const assistantMessageID = yield* publisher.assistantMessageID(event.id) yield* Effect.uninterruptibleMask((restore) => restore( prepared.toolMaterialization.settle({ @@ -215,7 +206,7 @@ export const make = Effect.gen(function* () { }), ).pipe( Effect.flatMap((settlement) => - runtime.publish( + publish( LLMEvent.toolResult({ id: event.id, name: event.name, @@ -226,32 +217,23 @@ export const make = Effect.gen(function* () { ), ), ), - ).pipe(FiberSet.run(runtime.toolFibers)) + ).pipe(FiberSet.run(toolFibers)) }), ), - Effect.ensuring(runtime.withPublication(runtime.publisher.flush())), + Effect.ensuring(withPublication(publisher.flush())), ) - /** - * The model response and tools remain interruptible. The short handoff after - * the response ends is protected so no started tool is forgotten before cleanup. - */ - const runAttempt = Effect.fn("SessionRunner.runTurn")(function* ( - sessionID: SessionSchema.ID, - promotion: SessionInput.Delivery | undefined, - recoverOverflow?: typeof compaction.compactAfterOverflow, - ) { - const prepared = yield* prepareTurn(sessionID, promotion) - const runtime = yield* startTurn(prepared) + // Keep cleanup protected after the response ends so no started tool is + // forgotten, while the response stream and tool work remain interruptible. return yield* Effect.uninterruptibleMask((restore) => Effect.gen(function* () { - const stream = yield* restore(consumeProvider(prepared, runtime)).pipe(Effect.exit) + const stream = yield* restore(providerStream).pipe(Effect.exit) const failure = stream._tag === "Failure" ? Option.getOrUndefined(Cause.findErrorOption(stream.cause)) : undefined if ( recoverOverflow && - !runtime.publisher.hasAssistantStarted() && - isContextOverflowFailure(runtime.overflowFailure ?? failure) && + !publisher.hasAssistantStarted() && + isContextOverflowFailure(overflowFailure ?? failure) && (yield* restore( recoverOverflow({ sessionID: prepared.session.id, @@ -262,63 +244,65 @@ export const make = Effect.gen(function* () { )) ) return AttemptResult.cases.CompactedOverflow.make({}) - if (runtime.overflowFailure) yield* runtime.publish(runtime.overflowFailure) + if (overflowFailure) yield* publish(overflowFailure) const llmFailure = failure instanceof LLMError ? failure : undefined - if (llmFailure && !runtime.publisher.hasProviderError()) { - yield* runtime.withPublication( - runtime.publisher.failUnsettledTools("Provider did not return a tool result", true), - ) - yield* runtime.withPublication( + if (llmFailure && !publisher.hasProviderError()) { + yield* withPublication(publisher.failUnsettledTools("Provider did not return a tool result", true)) + yield* withPublication( events.publish(SessionEvent.Step.Failed, { sessionID: prepared.session.id, timestamp: yield* DateTime.now, - assistantMessageID: yield* runtime.publisher.startAssistant(), + assistantMessageID: yield* publisher.startAssistant(), error: { type: "unknown", message: llmFailure.reason.message }, }), ) } - if (stream._tag === "Failure" && Cause.hasInterrupts(stream.cause)) yield* FiberSet.clear(runtime.toolFibers) - const settled = yield* restore(awaitToolFibers(runtime.toolFibers)).pipe(Effect.exit) + if (stream._tag === "Failure" && Cause.hasInterrupts(stream.cause)) yield* FiberSet.clear(toolFibers) + const settled = yield* restore(awaitToolFibers(toolFibers)).pipe(Effect.exit) if (settled._tag === "Failure" && isQuestionRejected(settled.cause)) { - yield* FiberSet.clear(runtime.toolFibers) - yield* runtime.withPublication(runtime.publisher.failUnsettledTools("Tool execution interrupted")) + yield* FiberSet.clear(toolFibers) + yield* withPublication(publisher.failUnsettledTools("Tool execution interrupted")) return yield* Effect.interrupt } if ( (stream._tag === "Failure" && Cause.hasInterrupts(stream.cause)) || (settled._tag === "Failure" && Cause.hasInterrupts(settled.cause)) ) { - yield* FiberSet.clear(runtime.toolFibers) - yield* runtime.withPublication(runtime.publisher.failUnsettledTools("Tool execution interrupted")) + yield* FiberSet.clear(toolFibers) + yield* withPublication(publisher.failUnsettledTools("Tool execution interrupted")) } if (settled._tag === "Failure" && !Cause.hasInterrupts(settled.cause)) { const failure = Cause.squash(settled.cause) const message = failure instanceof Error ? failure.message : String(failure) - yield* runtime.withPublication(runtime.publisher.failUnsettledTools(`Tool execution failed: ${message}`)) + yield* withPublication(publisher.failUnsettledTools(`Tool execution failed: ${message}`)) } - if (runtime.publisher.hasProviderError()) - yield* runtime.withPublication(runtime.publisher.failUnsettledTools("Tool execution interrupted")) - if (stream._tag === "Success" && !runtime.publisher.hasProviderError()) - yield* runtime.withPublication( - runtime.publisher.failUnsettledTools("Provider did not return a tool result", true), - ) + if (publisher.hasProviderError()) + yield* withPublication(publisher.failUnsettledTools("Tool execution interrupted")) + if (stream._tag === "Success" && !publisher.hasProviderError()) + 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 AttemptResult.cases.Complete.make({ - needsContinuation: !runtime.publisher.hasProviderError() && runtime.needsContinuation, + needsContinuation: !publisher.hasProviderError() && needsContinuation, }) }), ) }, Effect.scoped) - const run: Run = Effect.fnUntraced(function* (sessionID, promotion) { - const first = yield* runAttempt(sessionID, promotion, compaction.compactAfterOverflow) - if (first._tag === "Complete") return first.needsContinuation - const second = yield* runAttempt(sessionID, undefined) - return AttemptResult.match(second, { - Complete: (result) => result.needsContinuation, - CompactedOverflow: () => false, - }) + const run: Run = Effect.fn("SessionRunner.runTurn")(function* (input) { + let pendingDelivery = input.delivery + let canRecoverOverflow = true + while (true) { + const request = yield* buildRequest(input.sessionID, pendingDelivery) + pendingDelivery = undefined + const result = yield* streamAndSettle(request, canRecoverOverflow ? compaction.compactAfterOverflow : undefined) + const next = AttemptResult.match(result, { + Complete: (completed) => completed.needsContinuation, + CompactedOverflow: () => undefined, + }) + if (next !== undefined) return next + canRecoverOverflow = false + } }) return run