diff --git a/packages/core/src/session/runner/run-turn.ts b/packages/core/src/session/runner/run-turn.ts index 9790a90764b..cbcda8ffd4c 100644 --- a/packages/core/src/session/runner/run-turn.ts +++ b/packages/core/src/session/runner/run-turn.ts @@ -49,8 +49,6 @@ export interface Input { readonly delivery?: SessionInput.Delivery } -export type Run = (input: Input) => Effect.Effect - const AttemptResult = Schema.TaggedUnion({ Complete: { needsContinuation: Schema.Boolean }, CompactedOverflow: {}, @@ -90,6 +88,11 @@ export const make = Effect.gen(function* () { Effect.all([systemContext.load(), skillGuidance.load(agent)], { concurrency: "unbounded" }).pipe( Effect.map(SystemContext.combine), ) + const promoteDelivery = Effect.fnUntraced(function* (sessionID: SessionSchema.ID, delivery: SessionInput.Delivery) { + const cutoff = yield* SessionInput.latestSeq(db, sessionID) + if (delivery === "queue") yield* SessionInput.promoteNextQueued(db, events, sessionID) + yield* SessionInput.promoteSteers(db, events, sessionID, cutoff) + }) /** * Builds the next model request from durable Session state. @@ -113,12 +116,7 @@ export const make = Effect.gen(function* () { ) if (initialized === stale) continue if (pendingDelivery) { - const cutoff = yield* SessionInput.latestSeq(db, session.id) - 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) - } + yield* promoteDelivery(session.id, pendingDelivery) pendingDelivery = undefined } const prepared = @@ -167,7 +165,7 @@ export const make = Effect.gen(function* () { */ const streamAndSettle = Effect.fn("SessionRunner.streamAndSettle")(function* ( prepared: RequestSnapshot, - recoverOverflow?: typeof compaction.compactAfterOverflow, + canRecoverOverflow: boolean, ) { const publisher = createLLMEventPublisher(events, { sessionID: prepared.session.id, @@ -181,9 +179,37 @@ export const make = Effect.gen(function* () { const withPublication = Semaphore.makeUnsafe(1).withPermit const publish = (event: LLMEvent, outputPaths: ReadonlyArray = []) => withPublication(publisher.publish(event, outputPaths)) + const failUnsettled = (message: string, providerExecuted = false) => + withPublication(publisher.failUnsettledTools(message, providerExecuted)) const toolFibers = yield* FiberSet.make() let needsContinuation = false let overflowFailure: ProviderErrorEvent | undefined + const startTool = Effect.fnUntraced(function* (event: Extract) { + needsContinuation = true + const assistantMessageID = yield* publisher.assistantMessageID(event.id) + yield* Effect.uninterruptibleMask((restore) => + restore( + prepared.toolMaterialization.settle({ + sessionID: prepared.session.id, + agent: prepared.agent.id, + assistantMessageID, + call: event, + }), + ).pipe( + Effect.flatMap((settlement) => + publish( + LLMEvent.toolResult({ + id: event.id, + name: event.name, + result: settlement.result, + output: settlement.output, + }), + settlement.outputPaths ?? [], + ), + ), + ), + ).pipe(FiberSet.run(toolFibers)) + }) const providerStream = llm.stream(prepared.request).pipe( Stream.runForEach((event) => Effect.gen(function* () { @@ -194,30 +220,7 @@ export const make = Effect.gen(function* () { } yield* publish(event) if (event.type !== "tool-call" || event.providerExecuted) return - needsContinuation = true - const assistantMessageID = yield* publisher.assistantMessageID(event.id) - yield* Effect.uninterruptibleMask((restore) => - restore( - prepared.toolMaterialization.settle({ - sessionID: prepared.session.id, - agent: prepared.agent.id, - assistantMessageID, - call: event, - }), - ).pipe( - Effect.flatMap((settlement) => - publish( - LLMEvent.toolResult({ - id: event.id, - name: event.name, - result: settlement.result, - output: settlement.output, - }), - settlement.outputPaths ?? [], - ), - ), - ), - ).pipe(FiberSet.run(toolFibers)) + yield* startTool(event) }), ), Effect.ensuring(withPublication(publisher.flush())), @@ -231,11 +234,11 @@ export const make = Effect.gen(function* () { const failure = stream._tag === "Failure" ? Option.getOrUndefined(Cause.findErrorOption(stream.cause)) : undefined if ( - recoverOverflow && + canRecoverOverflow && !publisher.hasAssistantStarted() && isContextOverflowFailure(overflowFailure ?? failure) && (yield* restore( - recoverOverflow({ + compaction.compactAfterOverflow({ sessionID: prepared.session.id, entries: prepared.entries, model: prepared.model, @@ -247,7 +250,7 @@ export const make = Effect.gen(function* () { if (overflowFailure) yield* publish(overflowFailure) const llmFailure = failure instanceof LLMError ? failure : undefined if (llmFailure && !publisher.hasProviderError()) { - yield* withPublication(publisher.failUnsettledTools("Provider did not return a tool result", true)) + yield* failUnsettled("Provider did not return a tool result", true) yield* withPublication( events.publish(SessionEvent.Step.Failed, { sessionID: prepared.session.id, @@ -257,29 +260,25 @@ export const make = Effect.gen(function* () { }), ) } - if (stream._tag === "Failure" && Cause.hasInterrupts(stream.cause)) yield* FiberSet.clear(toolFibers) + const streamInterrupted = stream._tag === "Failure" && Cause.hasInterrupts(stream.cause) + if (streamInterrupted) yield* FiberSet.clear(toolFibers) const settled = yield* restore(awaitToolFibers(toolFibers)).pipe(Effect.exit) if (settled._tag === "Failure" && isQuestionRejected(settled.cause)) { yield* FiberSet.clear(toolFibers) - yield* withPublication(publisher.failUnsettledTools("Tool execution interrupted")) + yield* failUnsettled("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(toolFibers) - yield* withPublication(publisher.failUnsettledTools("Tool execution interrupted")) - } - if (settled._tag === "Failure" && !Cause.hasInterrupts(settled.cause)) { + const toolInterrupted = settled._tag === "Failure" && Cause.hasInterrupts(settled.cause) + if (toolInterrupted) yield* FiberSet.clear(toolFibers) + if (streamInterrupted || toolInterrupted || publisher.hasProviderError()) + yield* failUnsettled("Tool execution interrupted") + if (settled._tag === "Failure" && !toolInterrupted) { const failure = Cause.squash(settled.cause) const message = failure instanceof Error ? failure.message : String(failure) - yield* withPublication(publisher.failUnsettledTools(`Tool execution failed: ${message}`)) + yield* failUnsettled(`Tool execution failed: ${message}`) } - 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)) + yield* failUnsettled("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({ @@ -289,13 +288,13 @@ export const make = Effect.gen(function* () { ) }, Effect.scoped) - const run: Run = Effect.fn("SessionRunner.runTurn")(function* (input) { + const run = Effect.fn("SessionRunner.runTurn")(function* (input: Input): Effect.fn.Return { 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 result = yield* streamAndSettle(request, canRecoverOverflow) const next = AttemptResult.match(result, { Complete: (completed) => completed.needsContinuation, CompactedOverflow: () => undefined,