refactor(core): clarify provider turn helpers

This commit is contained in:
Kit Langton 2026-06-06 22:27:19 -04:00
commit b795f415af

View file

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