Compare commits
10 commits
dev
...
refactor/c
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
7ad5c27e7c | ||
|
|
576a541f67 | ||
|
|
ad639f89fb | ||
|
|
cacb08c1be | ||
|
|
b795f415af | ||
|
|
652ec048e7 | ||
|
|
1cb6f7b091 | ||
|
|
5462e08ac1 | ||
|
|
ad04846593 | ||
|
|
4f6a2c5b69 |
2 changed files with 340 additions and 281 deletions
|
|
@ -1,37 +1,12 @@
|
||||||
import {
|
import { DateTime, Effect, Layer } from "effect"
|
||||||
LLM,
|
|
||||||
LLMClient,
|
|
||||||
LLMError,
|
|
||||||
LLMEvent,
|
|
||||||
SystemPart,
|
|
||||||
isContextOverflowFailure,
|
|
||||||
type ProviderErrorEvent,
|
|
||||||
} from "@opencode-ai/llm"
|
|
||||||
import { Cause, DateTime, Effect, FiberSet, Layer, Option, Schema, Semaphore, Stream } from "effect"
|
|
||||||
import { AgentV2 } from "../../agent"
|
|
||||||
import { Config } from "../../config"
|
|
||||||
import { Database } from "../../database/database"
|
import { Database } from "../../database/database"
|
||||||
import { EventV2 } from "../../event"
|
import { EventV2 } from "../../event"
|
||||||
import { Location } from "../../location"
|
|
||||||
import { ModelV2 } from "../../model"
|
|
||||||
import { ProviderV2 } from "../../provider"
|
|
||||||
import { QuestionV2 } from "../../question"
|
|
||||||
import { SystemContext } from "../../system-context/index"
|
|
||||||
import { SystemContextRegistry } from "../../system-context/registry"
|
|
||||||
import { SkillGuidance } from "../../skill/guidance"
|
|
||||||
import { ToolRegistry } from "../../tool/registry"
|
|
||||||
import { ToolOutputStore } from "../../tool-output-store"
|
|
||||||
import { SessionContextEpoch } from "../context-epoch"
|
|
||||||
import { SessionCompaction } from "../compaction"
|
|
||||||
import { SessionEvent } from "../event"
|
import { SessionEvent } from "../event"
|
||||||
import { SessionHistory } from "../history"
|
|
||||||
import { SessionInput } from "../input"
|
import { SessionInput } from "../input"
|
||||||
import { SessionSchema } from "../schema"
|
import { SessionSchema } from "../schema"
|
||||||
import { SessionStore } from "../store"
|
import { SessionStore } from "../store"
|
||||||
import { type RunError, Service, StepLimitExceededError } from "./index"
|
import { Service, StepLimitExceededError } from "./index"
|
||||||
import { SessionRunnerModel } from "./model"
|
import { RunTurn } from "./run-turn"
|
||||||
import { createLLMEventPublisher } from "./publish-llm-event"
|
|
||||||
import { toLLMMessages } from "./to-llm-message"
|
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Runs one durable coding-agent Session until it settles.
|
* Runs one durable coding-agent Session until it settles.
|
||||||
|
|
@ -75,8 +50,9 @@ import { toLLMMessages } from "./to-llm-message"
|
||||||
* - [ ] Coalesce streamed deltas and add covering projected-history indexes.
|
* - [ ] Coalesce streamed deltas and add covering projected-history indexes.
|
||||||
* - [ ] Update title, summaries, compaction state, and cleanup in bounded background work.
|
* - [ ] Update title, summaries, compaction state, and cleanup in bounded background work.
|
||||||
*
|
*
|
||||||
* Use `llm.stream(request)` for each provider turn. Keep tool execution and continuation here.
|
* `RunTurn` owns provider-turn preparation, streaming, tool settlement, and continuation signals.
|
||||||
* Durable activity recovery remains a separate future slice with an explicit retry policy.
|
* This module owns durable activity scheduling and bounded continuation. Durable activity recovery
|
||||||
|
* remains a separate future slice with an explicit retry policy.
|
||||||
*
|
*
|
||||||
* The current slice loads V2 history, translates it, resolves a model through a core service, and persists one
|
* The current slice loads V2 history, translates it, resolves a model through a core service, and persists one
|
||||||
* provider turn. Registry definitions are advertised, local tool calls are settled durably, and a
|
* provider turn. Registry definitions are advertised, local tool calls are settled durably, and a
|
||||||
|
|
@ -90,22 +66,9 @@ export const layer = Layer.effect(
|
||||||
Service,
|
Service,
|
||||||
Effect.gen(function* () {
|
Effect.gen(function* () {
|
||||||
const events = yield* EventV2.Service
|
const events = yield* EventV2.Service
|
||||||
const llm = yield* LLMClient.Service
|
|
||||||
const agents = yield* AgentV2.Service
|
|
||||||
const tools = yield* ToolRegistry.Service
|
|
||||||
const models = yield* SessionRunnerModel.Service
|
|
||||||
const store = yield* SessionStore.Service
|
const store = yield* SessionStore.Service
|
||||||
const location = yield* Location.Service
|
|
||||||
const systemContext = yield* SystemContextRegistry.Service
|
|
||||||
const skillGuidance = yield* SkillGuidance.Service
|
|
||||||
const config = yield* Config.Service
|
|
||||||
const db = (yield* Database.Service).db
|
const db = (yield* Database.Service).db
|
||||||
const compaction = SessionCompaction.make({ events, llm, config: yield* config.entries() })
|
const runTurn = yield* RunTurn.make
|
||||||
const getSession = Effect.fn("SessionRunner.getSession")(function* (sessionID: SessionSchema.ID) {
|
|
||||||
const session = yield* store.get(sessionID)
|
|
||||||
if (!session) return yield* Effect.die(`Session not found: ${sessionID}`)
|
|
||||||
return session
|
|
||||||
})
|
|
||||||
|
|
||||||
const getContext = Effect.fn("SessionRunner.getContext")(function* (sessionID: SessionSchema.ID) {
|
const getContext = Effect.fn("SessionRunner.getContext")(function* (sessionID: SessionSchema.ID) {
|
||||||
return yield* store.context(sessionID)
|
return yield* store.context(sessionID)
|
||||||
|
|
@ -132,242 +95,6 @@ export const layer = Layer.effect(
|
||||||
}
|
}
|
||||||
})
|
})
|
||||||
|
|
||||||
const awaitToolFibers = (fibers: FiberSet.FiberSet<void, ToolOutputStore.Error>) =>
|
|
||||||
Effect.raceFirst(FiberSet.join(fibers), FiberSet.awaitEmpty(fibers))
|
|
||||||
|
|
||||||
// Match V1: dismissing a question halts the loop instead of becoming model-facing tool output.
|
|
||||||
const isQuestionRejected = (cause: Cause.Cause<unknown>) =>
|
|
||||||
cause.reasons.some((reason) => Cause.isDieReason(reason) && reason.defect instanceof QuestionV2.RejectedError)
|
|
||||||
|
|
||||||
type TurnTransition =
|
|
||||||
// Request preparation observed a concurrent Session change and must restart from durable state.
|
|
||||||
| { readonly _tag: "RebuildPreparedTurn"; readonly promotion?: SessionInput.Delivery }
|
|
||||||
// Overflow compaction completed; rebuild once through the path without overflow recovery.
|
|
||||||
| { readonly _tag: "ContinueAfterOverflowCompaction" }
|
|
||||||
|
|
||||||
class TurnTransitionError extends Error {
|
|
||||||
constructor(readonly transition: TurnTransition) {
|
|
||||||
super()
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
const rebuildPreparedTurn = (promotion?: SessionInput.Delivery) =>
|
|
||||||
new TurnTransitionError({ _tag: "RebuildPreparedTurn", promotion })
|
|
||||||
const continueAfterOverflowCompaction = new TurnTransitionError({
|
|
||||||
_tag: "ContinueAfterOverflowCompaction",
|
|
||||||
})
|
|
||||||
|
|
||||||
const retryAgentMismatch = (promotion: SessionInput.Delivery | undefined) =>
|
|
||||||
Effect.catchDefect((defect) =>
|
|
||||||
defect instanceof SessionContextEpoch.AgentMismatch
|
|
||||||
? Effect.die(rebuildPreparedTurn(promotion))
|
|
||||||
: Effect.die(defect),
|
|
||||||
)
|
|
||||||
|
|
||||||
const sameModel = Schema.toEquivalence(Schema.UndefinedOr(ModelV2.Ref))
|
|
||||||
const loadSystemContext = (agent: AgentV2.Selection) =>
|
|
||||||
Effect.all([systemContext.load(), skillGuidance.load(agent)], { concurrency: "unbounded" }).pipe(
|
|
||||||
Effect.map(SystemContext.combine),
|
|
||||||
)
|
|
||||||
|
|
||||||
const runTurnAttempt = Effect.fn("SessionRunner.runTurn")(function* (
|
|
||||||
sessionID: SessionSchema.ID,
|
|
||||||
promotion: SessionInput.Delivery | undefined,
|
|
||||||
recoverOverflow?: typeof compaction.compactAfterOverflow,
|
|
||||||
) {
|
|
||||||
const session = yield* getSession(sessionID)
|
|
||||||
if (session.location.directory !== location.directory || session.location.workspaceID !== location.workspaceID)
|
|
||||||
return yield* Effect.interrupt
|
|
||||||
const agent = yield* agents.select(session.agent)
|
|
||||||
const initialized = yield* SessionContextEpoch.initialize(
|
|
||||||
db,
|
|
||||||
loadSystemContext(agent),
|
|
||||||
session.id,
|
|
||||||
session.location,
|
|
||||||
agent.id,
|
|
||||||
).pipe(retryAgentMismatch(promotion))
|
|
||||||
const toolFibers = yield* FiberSet.make<void, ToolOutputStore.Error>()
|
|
||||||
let needsContinuation = false
|
|
||||||
if (promotion) {
|
|
||||||
const cutoff = yield* SessionInput.latestSeq(db, session.id)
|
|
||||||
if (promotion === "steer") yield* SessionInput.promoteSteers(db, events, session.id, cutoff)
|
|
||||||
if (promotion === "queue") {
|
|
||||||
yield* SessionInput.promoteNextQueued(db, events, session.id)
|
|
||||||
yield* SessionInput.promoteSteers(db, events, session.id, cutoff)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
const system =
|
|
||||||
initialized ??
|
|
||||||
(yield* SessionContextEpoch.prepare(
|
|
||||||
db,
|
|
||||||
events,
|
|
||||||
loadSystemContext(agent),
|
|
||||||
session.id,
|
|
||||||
session.location,
|
|
||||||
agent.id,
|
|
||||||
).pipe(retryAgentMismatch(undefined)))
|
|
||||||
const current = yield* getSession(sessionID)
|
|
||||||
if ((yield* agents.select(current.agent)).id !== agent.id || !sameModel(current.model, session.model))
|
|
||||||
return yield* Effect.die(rebuildPreparedTurn())
|
|
||||||
const model = yield* models.resolve(session)
|
|
||||||
const entries = yield* SessionHistory.entriesForRunner(db, session.id, system.baselineSeq)
|
|
||||||
const context = entries.map((entry) => entry.message)
|
|
||||||
const toolMaterialization = yield* tools.materialize(agent.info?.permissions)
|
|
||||||
const promptCacheKey = /^ses_[0-9a-f]{64}$/.test(session.id) ? session.id.slice(4) : session.id
|
|
||||||
const request = LLM.request({
|
|
||||||
model,
|
|
||||||
providerOptions: { openai: { promptCacheKey } },
|
|
||||||
system: [agent.info?.system, system.baseline]
|
|
||||||
.filter((part): part is string => part !== undefined && part.length > 0)
|
|
||||||
.map(SystemPart.make),
|
|
||||||
messages: toLLMMessages(context, model),
|
|
||||||
tools: toolMaterialization.definitions,
|
|
||||||
})
|
|
||||||
if (yield* compaction.compactIfNeeded({ sessionID: session.id, entries, model, request }))
|
|
||||||
return yield* Effect.die(rebuildPreparedTurn())
|
|
||||||
const publisher = createLLMEventPublisher(events, {
|
|
||||||
sessionID: session.id,
|
|
||||||
agent: agent.id,
|
|
||||||
model: {
|
|
||||||
id: ModelV2.ID.make(model.id),
|
|
||||||
providerID: ProviderV2.ID.make(model.provider),
|
|
||||||
...(session.model?.variant === undefined ? {} : { variant: session.model.variant }),
|
|
||||||
},
|
|
||||||
})
|
|
||||||
const withPublication = Semaphore.makeUnsafe(1).withPermit
|
|
||||||
const publish = (event: LLMEvent, outputPaths: ReadonlyArray<string> = []) =>
|
|
||||||
withPublication(publisher.publish(event, outputPaths))
|
|
||||||
let overflowFailure: ProviderErrorEvent | undefined
|
|
||||||
if (!(yield* SessionContextEpoch.current(db, session.id, agent.id, system.revision)))
|
|
||||||
return yield* Effect.die(rebuildPreparedTurn())
|
|
||||||
const providerStream = llm.stream(request).pipe(
|
|
||||||
Stream.runForEach((event) =>
|
|
||||||
Effect.gen(function* () {
|
|
||||||
if (overflowFailure || publisher.hasProviderError()) return
|
|
||||||
if (LLMEvent.is.providerError(event)) {
|
|
||||||
if (isContextOverflowFailure(event) && !publisher.hasAssistantStarted()) {
|
|
||||||
overflowFailure = event
|
|
||||||
return
|
|
||||||
}
|
|
||||||
}
|
|
||||||
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(
|
|
||||||
toolMaterialization.settle({
|
|
||||||
sessionID: session.id,
|
|
||||||
agent: 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())),
|
|
||||||
)
|
|
||||||
|
|
||||||
return yield* Effect.uninterruptibleMask((restore) =>
|
|
||||||
Effect.gen(function* () {
|
|
||||||
const stream = yield* restore(providerStream).pipe(Effect.exit)
|
|
||||||
const failure =
|
|
||||||
stream._tag === "Failure" ? Option.getOrUndefined(Cause.findErrorOption(stream.cause)) : undefined
|
|
||||||
if (
|
|
||||||
recoverOverflow &&
|
|
||||||
!publisher.hasAssistantStarted() &&
|
|
||||||
isContextOverflowFailure(overflowFailure ?? failure) &&
|
|
||||||
(yield* restore(recoverOverflow({ sessionID: session.id, entries, model, request })))
|
|
||||||
)
|
|
||||||
return yield* Effect.die(continueAfterOverflowCompaction)
|
|
||||||
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* withPublication(
|
|
||||||
events.publish(SessionEvent.Step.Failed, {
|
|
||||||
sessionID: session.id,
|
|
||||||
timestamp: yield* DateTime.now,
|
|
||||||
assistantMessageID: yield* publisher.startAssistant(),
|
|
||||||
error: { type: "unknown", message: llmFailure.reason.message },
|
|
||||||
}),
|
|
||||||
)
|
|
||||||
}
|
|
||||||
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(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(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* withPublication(publisher.failUnsettledTools(`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))
|
|
||||||
if (stream._tag === "Failure") return yield* Effect.failCause(stream.cause)
|
|
||||||
if (settled._tag === "Failure") return yield* Effect.failCause(settled.cause)
|
|
||||||
return !publisher.hasProviderError() && needsContinuation
|
|
||||||
}),
|
|
||||||
)
|
|
||||||
}, Effect.scoped)
|
|
||||||
type RunTurn = (
|
|
||||||
sessionID: SessionSchema.ID,
|
|
||||||
promotion: SessionInput.Delivery | undefined,
|
|
||||||
) => Effect.Effect<boolean, RunError>
|
|
||||||
|
|
||||||
const runAfterOverflowCompaction: RunTurn = Effect.fnUntraced(function* (sessionID, promotion) {
|
|
||||||
return yield* runTurnAttempt(sessionID, promotion).pipe(
|
|
||||||
Effect.catchDefect(
|
|
||||||
Effect.fnUntraced(function* (defect) {
|
|
||||||
if (!(defect instanceof TurnTransitionError)) return yield* Effect.die(defect)
|
|
||||||
if (defect.transition._tag === "ContinueAfterOverflowCompaction")
|
|
||||||
return yield* Effect.die("Post-compaction provider attempt cannot recover another overflow")
|
|
||||||
yield* Effect.yieldNow
|
|
||||||
return yield* runAfterOverflowCompaction(sessionID, defect.transition.promotion)
|
|
||||||
}),
|
|
||||||
),
|
|
||||||
)
|
|
||||||
})
|
|
||||||
|
|
||||||
const runTurn: RunTurn = Effect.fnUntraced(function* (sessionID, promotion) {
|
|
||||||
return yield* runTurnAttempt(sessionID, promotion, compaction.compactAfterOverflow).pipe(
|
|
||||||
Effect.catchDefect(
|
|
||||||
Effect.fnUntraced(function* (defect) {
|
|
||||||
if (!(defect instanceof TurnTransitionError)) return yield* Effect.die(defect)
|
|
||||||
yield* Effect.yieldNow
|
|
||||||
if (defect.transition._tag === "ContinueAfterOverflowCompaction")
|
|
||||||
return yield* runAfterOverflowCompaction(sessionID, undefined)
|
|
||||||
return yield* runTurn(sessionID, defect.transition.promotion)
|
|
||||||
}),
|
|
||||||
),
|
|
||||||
)
|
|
||||||
})
|
|
||||||
|
|
||||||
const run = Effect.fn("SessionRunner.run")(function* (input: {
|
const run = Effect.fn("SessionRunner.run")(function* (input: {
|
||||||
readonly sessionID: SessionSchema.ID
|
readonly sessionID: SessionSchema.ID
|
||||||
readonly force?: boolean
|
readonly force?: boolean
|
||||||
|
|
@ -381,7 +108,7 @@ export const layer = Layer.effect(
|
||||||
while (openActivity) {
|
while (openActivity) {
|
||||||
let needsContinuation = true
|
let needsContinuation = true
|
||||||
for (let step = 0; step < MAX_STEPS; step++) {
|
for (let step = 0; step < MAX_STEPS; step++) {
|
||||||
needsContinuation = yield* runTurn(input.sessionID, promotion)
|
needsContinuation = yield* runTurn({ sessionID: input.sessionID, delivery: promotion })
|
||||||
promotion = "steer"
|
promotion = "steer"
|
||||||
if (!needsContinuation) needsContinuation = yield* SessionInput.hasPending(db, input.sessionID, "steer")
|
if (!needsContinuation) needsContinuation = yield* SessionInput.hasPending(db, input.sessionID, "steer")
|
||||||
if (!needsContinuation) break
|
if (!needsContinuation) break
|
||||||
|
|
|
||||||
332
packages/core/src/session/runner/run-turn.ts
Normal file
332
packages/core/src/session/runner/run-turn.ts
Normal file
|
|
@ -0,0 +1,332 @@
|
||||||
|
export * as RunTurn from "./run-turn"
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Sends the next request to the model and finishes every tool call it starts.
|
||||||
|
*
|
||||||
|
* Before sending, it makes admitted input visible, loads the latest Session
|
||||||
|
* history and instructions, and compacts oversized history. If the model rejects
|
||||||
|
* the request for being too large before producing output, it may compact and
|
||||||
|
* try once more. Returns `true` when tool results require another model request.
|
||||||
|
*/
|
||||||
|
|
||||||
|
import {
|
||||||
|
LLM,
|
||||||
|
LLMClient,
|
||||||
|
LLMError,
|
||||||
|
LLMEvent,
|
||||||
|
SystemPart,
|
||||||
|
isContextOverflowFailure,
|
||||||
|
type ProviderErrorEvent,
|
||||||
|
} from "@opencode-ai/llm"
|
||||||
|
import { Cause, DateTime, Effect, Exit, FiberSet, Option, Schema, Semaphore, Stream } from "effect"
|
||||||
|
import { AgentV2 } from "../../agent"
|
||||||
|
import { Config } from "../../config"
|
||||||
|
import { Database } from "../../database/database"
|
||||||
|
import { EventV2 } from "../../event"
|
||||||
|
import { Location } from "../../location"
|
||||||
|
import { ModelV2 } from "../../model"
|
||||||
|
import { ProviderV2 } from "../../provider"
|
||||||
|
import { QuestionV2 } from "../../question"
|
||||||
|
import { SkillGuidance } from "../../skill/guidance"
|
||||||
|
import { SystemContext } from "../../system-context/index"
|
||||||
|
import { SystemContextRegistry } from "../../system-context/registry"
|
||||||
|
import { ToolOutputStore } from "../../tool-output-store"
|
||||||
|
import { ToolRegistry } from "../../tool/registry"
|
||||||
|
import { SessionCompaction } from "../compaction"
|
||||||
|
import { SessionContextEpoch } from "../context-epoch"
|
||||||
|
import { SessionEvent } from "../event"
|
||||||
|
import { SessionHistory } from "../history"
|
||||||
|
import { SessionInput } from "../input"
|
||||||
|
import { SessionSchema } from "../schema"
|
||||||
|
import { SessionStore } from "../store"
|
||||||
|
import type { RunError } from "./index"
|
||||||
|
import { SessionRunnerModel } from "./model"
|
||||||
|
import { createLLMEventPublisher } from "./publish-llm-event"
|
||||||
|
import { toLLMMessages } from "./to-llm-message"
|
||||||
|
|
||||||
|
export interface Input {
|
||||||
|
readonly sessionID: SessionSchema.ID
|
||||||
|
readonly delivery?: SessionInput.Delivery
|
||||||
|
}
|
||||||
|
|
||||||
|
const AttemptResult = Schema.TaggedUnion({
|
||||||
|
Complete: { needsContinuation: Schema.Boolean },
|
||||||
|
CompactedOverflow: {},
|
||||||
|
})
|
||||||
|
|
||||||
|
const isQuestionRejected = (cause: Cause.Cause<unknown>) =>
|
||||||
|
cause.reasons.some((reason) => Cause.isDieReason(reason) && reason.defect instanceof QuestionV2.RejectedError)
|
||||||
|
|
||||||
|
const awaitTools = (fibers: FiberSet.FiberSet<void, ToolOutputStore.Error>) =>
|
||||||
|
Effect.raceFirst(FiberSet.join(fibers), FiberSet.awaitEmpty(fibers))
|
||||||
|
|
||||||
|
export const make = Effect.gen(function* () {
|
||||||
|
const events = yield* EventV2.Service
|
||||||
|
const llm = yield* LLMClient.Service
|
||||||
|
const agents = yield* AgentV2.Service
|
||||||
|
const tools = yield* ToolRegistry.Service
|
||||||
|
const models = yield* SessionRunnerModel.Service
|
||||||
|
const store = yield* SessionStore.Service
|
||||||
|
const location = yield* Location.Service
|
||||||
|
const systemContext = yield* SystemContextRegistry.Service
|
||||||
|
const skillGuidance = yield* SkillGuidance.Service
|
||||||
|
const config = yield* Config.Service
|
||||||
|
const db = (yield* Database.Service).db
|
||||||
|
const compaction = SessionCompaction.make({ events, llm, config: yield* config.entries() })
|
||||||
|
|
||||||
|
const getSession = Effect.fn("SessionRunner.getSession")(function* (sessionID: SessionSchema.ID) {
|
||||||
|
const session = yield* store.get(sessionID)
|
||||||
|
if (!session) return yield* Effect.die(`Session not found: ${sessionID}`)
|
||||||
|
return session
|
||||||
|
})
|
||||||
|
const stale = Symbol("stale turn preparation")
|
||||||
|
const retryAgentMismatch = <A, E, R>(effect: Effect.Effect<A, E, R>) =>
|
||||||
|
effect.pipe(
|
||||||
|
Effect.catchDefect((defect) =>
|
||||||
|
defect instanceof SessionContextEpoch.AgentMismatch ? Effect.succeed(stale) : Effect.die(defect),
|
||||||
|
),
|
||||||
|
)
|
||||||
|
const loadSystemContext = (agent: AgentV2.Selection) =>
|
||||||
|
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.
|
||||||
|
*
|
||||||
|
* Initial instructions must be available before admitted input becomes visible.
|
||||||
|
* 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 buildRequest = Effect.fn("SessionRunner.buildRequest")(function* (
|
||||||
|
sessionID: SessionSchema.ID,
|
||||||
|
delivery: SessionInput.Delivery | undefined,
|
||||||
|
) {
|
||||||
|
let pendingDelivery = delivery
|
||||||
|
while (true) {
|
||||||
|
const session = yield* getSession(sessionID)
|
||||||
|
if (session.location.directory !== location.directory || session.location.workspaceID !== location.workspaceID)
|
||||||
|
return yield* Effect.interrupt
|
||||||
|
const agent = yield* agents.select(session.agent)
|
||||||
|
const initialized = yield* retryAgentMismatch(
|
||||||
|
SessionContextEpoch.initialize(db, loadSystemContext(agent), session.id, session.location, agent.id),
|
||||||
|
)
|
||||||
|
if (initialized === stale) continue
|
||||||
|
if (pendingDelivery) {
|
||||||
|
yield* promoteDelivery(session.id, pendingDelivery)
|
||||||
|
pendingDelivery = undefined
|
||||||
|
}
|
||||||
|
const prepared =
|
||||||
|
initialized ??
|
||||||
|
(yield* retryAgentMismatch(
|
||||||
|
SessionContextEpoch.prepare(db, events, loadSystemContext(agent), session.id, session.location, agent.id),
|
||||||
|
))
|
||||||
|
if (prepared === stale) continue
|
||||||
|
const system = prepared
|
||||||
|
const model = yield* models.resolve(session)
|
||||||
|
const entries = yield* SessionHistory.entriesForRunner(db, session.id, system.baselineSeq)
|
||||||
|
const toolMaterialization = yield* tools.materialize(agent.info?.permissions)
|
||||||
|
const request = LLM.request({
|
||||||
|
model,
|
||||||
|
providerOptions: {
|
||||||
|
openai: { promptCacheKey: /^ses_[0-9a-f]{64}$/.test(session.id) ? session.id.slice(4) : session.id },
|
||||||
|
},
|
||||||
|
system: [agent.info?.system, system.baseline]
|
||||||
|
.filter((part): part is string => part !== undefined && part.length > 0)
|
||||||
|
.map(SystemPart.make),
|
||||||
|
messages: toLLMMessages(
|
||||||
|
entries.map((entry) => entry.message),
|
||||||
|
model,
|
||||||
|
),
|
||||||
|
tools: toolMaterialization.definitions,
|
||||||
|
})
|
||||||
|
if (yield* compaction.compactIfNeeded({ sessionID: session.id, entries, model, request })) {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
if (!(yield* SessionContextEpoch.current(db, session.id, agent.id, system.revision))) {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
return { session, agent, model, entries, request, toolMaterialization }
|
||||||
|
}
|
||||||
|
})
|
||||||
|
|
||||||
|
type RequestSnapshot = Effect.Success<ReturnType<typeof buildRequest>>
|
||||||
|
|
||||||
|
/**
|
||||||
|
* 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 streamAndSettle = Effect.fn("SessionRunner.streamAndSettle")(function* (
|
||||||
|
prepared: RequestSnapshot,
|
||||||
|
canRecoverOverflow: boolean,
|
||||||
|
) {
|
||||||
|
const publisher = createLLMEventPublisher(events, {
|
||||||
|
sessionID: prepared.session.id,
|
||||||
|
agent: prepared.agent.id,
|
||||||
|
model: {
|
||||||
|
id: ModelV2.ID.make(prepared.model.id),
|
||||||
|
providerID: ProviderV2.ID.make(prepared.model.provider),
|
||||||
|
...(prepared.session.model?.variant === undefined ? {} : { variant: prepared.session.model.variant }),
|
||||||
|
},
|
||||||
|
})
|
||||||
|
const withPublication = Semaphore.makeUnsafe(1).withPermit
|
||||||
|
const publish = (event: LLMEvent, outputPaths: ReadonlyArray<string> = []) =>
|
||||||
|
withPublication(publisher.publish(event, outputPaths))
|
||||||
|
const failUnsettled = (message: string, providerExecuted = false) =>
|
||||||
|
withPublication(publisher.failUnsettledTools(message, providerExecuted))
|
||||||
|
let needsContinuation = false
|
||||||
|
let overflowFailure: ProviderErrorEvent | undefined
|
||||||
|
const toolEffect = Effect.fnUntraced(function* (event: Extract<LLMEvent, { readonly type: "tool-call" }>) {
|
||||||
|
needsContinuation = true
|
||||||
|
const assistantMessageID = yield* publisher.assistantMessageID(event.id)
|
||||||
|
return yield* 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 ?? [],
|
||||||
|
),
|
||||||
|
),
|
||||||
|
)
|
||||||
|
})
|
||||||
|
const handleProviderEvent = Effect.fnUntraced(function* (
|
||||||
|
event: LLMEvent,
|
||||||
|
tools: FiberSet.FiberSet<void, ToolOutputStore.Error>,
|
||||||
|
) {
|
||||||
|
if (overflowFailure || publisher.hasProviderError()) return
|
||||||
|
if (LLMEvent.is.providerError(event) && isContextOverflowFailure(event) && !publisher.hasAssistantStarted()) {
|
||||||
|
overflowFailure = event
|
||||||
|
return
|
||||||
|
}
|
||||||
|
yield* publish(event)
|
||||||
|
if (event.type !== "tool-call" || event.providerExecuted) return
|
||||||
|
yield* toolEffect(event).pipe(FiberSet.run(tools))
|
||||||
|
})
|
||||||
|
const providerStream = (tools: FiberSet.FiberSet<void, ToolOutputStore.Error>) =>
|
||||||
|
llm.stream(prepared.request).pipe(
|
||||||
|
Stream.runForEach((event) => handleProviderEvent(event, tools)),
|
||||||
|
Effect.ensuring(withPublication(publisher.flush())),
|
||||||
|
)
|
||||||
|
const projectProviderFailure = Effect.fnUntraced(function* (failure: unknown) {
|
||||||
|
if (overflowFailure) yield* publish(overflowFailure)
|
||||||
|
if (!(failure instanceof LLMError) || publisher.hasProviderError()) return
|
||||||
|
yield* failUnsettled("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* publisher.startAssistant(),
|
||||||
|
error: { type: "unknown", message: failure.reason.message },
|
||||||
|
}),
|
||||||
|
)
|
||||||
|
})
|
||||||
|
const finishToolsSuccessfully = Effect.fnUntraced(function* <A, E>(stream: Exit.Exit<A, E>) {
|
||||||
|
if (Exit.hasInterrupts(stream) || publisher.hasProviderError()) yield* failUnsettled("Tool execution interrupted")
|
||||||
|
if (Exit.isSuccess(stream) && !publisher.hasProviderError())
|
||||||
|
yield* failUnsettled("Provider did not return a tool result", true)
|
||||||
|
if (Exit.isFailure(stream)) return yield* Effect.failCause(stream.cause)
|
||||||
|
})
|
||||||
|
const finishToolsAfterFailure = Effect.fnUntraced(function* <A, E>(
|
||||||
|
stream: Exit.Exit<A, E>,
|
||||||
|
tools: FiberSet.FiberSet<void, ToolOutputStore.Error>,
|
||||||
|
cause: Cause.Cause<ToolOutputStore.Error>,
|
||||||
|
) {
|
||||||
|
if (isQuestionRejected(cause)) {
|
||||||
|
yield* FiberSet.clear(tools)
|
||||||
|
yield* failUnsettled("Tool execution interrupted")
|
||||||
|
return yield* Effect.interrupt
|
||||||
|
}
|
||||||
|
const interrupted = Cause.hasInterrupts(cause)
|
||||||
|
if (interrupted) yield* FiberSet.clear(tools)
|
||||||
|
if (Exit.hasInterrupts(stream) || interrupted || publisher.hasProviderError())
|
||||||
|
yield* failUnsettled("Tool execution interrupted")
|
||||||
|
if (!interrupted) {
|
||||||
|
const failure = Cause.squash(cause)
|
||||||
|
yield* failUnsettled(`Tool execution failed: ${failure instanceof Error ? failure.message : String(failure)}`)
|
||||||
|
}
|
||||||
|
if (Exit.isFailure(stream)) return yield* Effect.failCause(stream.cause)
|
||||||
|
return yield* Effect.failCause(cause)
|
||||||
|
})
|
||||||
|
const finishTools = Effect.fnUntraced(function* <A, E>(
|
||||||
|
stream: Exit.Exit<A, E>,
|
||||||
|
tools: FiberSet.FiberSet<void, ToolOutputStore.Error>,
|
||||||
|
wait: Effect.Effect<void, ToolOutputStore.Error>,
|
||||||
|
) {
|
||||||
|
if (Exit.hasInterrupts(stream)) yield* FiberSet.clear(tools)
|
||||||
|
return yield* wait.pipe(
|
||||||
|
Effect.matchCauseEffect({
|
||||||
|
onFailure: (cause) => finishToolsAfterFailure(stream, tools, cause),
|
||||||
|
onSuccess: () => finishToolsSuccessfully(stream),
|
||||||
|
}),
|
||||||
|
)
|
||||||
|
})
|
||||||
|
const settleProviderTurn = Effect.fnUntraced(function* () {
|
||||||
|
const tools = yield* FiberSet.make<void, ToolOutputStore.Error>()
|
||||||
|
return yield* Effect.uninterruptibleMask((restore) =>
|
||||||
|
Effect.gen(function* () {
|
||||||
|
const stream = yield* restore(providerStream(tools)).pipe(Effect.exit)
|
||||||
|
const failure =
|
||||||
|
stream._tag === "Failure" ? Option.getOrUndefined(Cause.findErrorOption(stream.cause)) : undefined
|
||||||
|
const recovered =
|
||||||
|
canRecoverOverflow &&
|
||||||
|
!publisher.hasAssistantStarted() &&
|
||||||
|
isContextOverflowFailure(overflowFailure ?? failure) &&
|
||||||
|
(yield* restore(
|
||||||
|
compaction.compactAfterOverflow({
|
||||||
|
sessionID: prepared.session.id,
|
||||||
|
entries: prepared.entries,
|
||||||
|
model: prepared.model,
|
||||||
|
request: prepared.request,
|
||||||
|
}),
|
||||||
|
))
|
||||||
|
if (recovered) return AttemptResult.cases.CompactedOverflow.make({})
|
||||||
|
|
||||||
|
yield* projectProviderFailure(failure)
|
||||||
|
yield* finishTools(stream, tools, restore(awaitTools(tools)))
|
||||||
|
return AttemptResult.cases.Complete.make({
|
||||||
|
needsContinuation: !publisher.hasProviderError() && needsContinuation,
|
||||||
|
})
|
||||||
|
}),
|
||||||
|
)
|
||||||
|
})
|
||||||
|
return yield* settleProviderTurn()
|
||||||
|
}, Effect.scoped)
|
||||||
|
|
||||||
|
const run = Effect.fn("SessionRunner.runTurn")(function* (input: Input): Effect.fn.Return<boolean, RunError> {
|
||||||
|
let pendingDelivery = input.delivery
|
||||||
|
let canRecoverOverflow = true
|
||||||
|
while (true) {
|
||||||
|
const request = yield* buildRequest(input.sessionID, pendingDelivery)
|
||||||
|
pendingDelivery = undefined
|
||||||
|
const result = yield* streamAndSettle(request, canRecoverOverflow)
|
||||||
|
const next = AttemptResult.match(result, {
|
||||||
|
Complete: (completed) => completed.needsContinuation,
|
||||||
|
CompactedOverflow: () => undefined,
|
||||||
|
})
|
||||||
|
if (next !== undefined) return next
|
||||||
|
canRecoverOverflow = false
|
||||||
|
}
|
||||||
|
})
|
||||||
|
|
||||||
|
return run
|
||||||
|
})
|
||||||
Loading…
Add table
Add a link
Reference in a new issue