Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
35 changes: 34 additions & 1 deletion packages/core/src/session/runner/llm.ts
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
export * as SessionRunnerLLM from "./llm.js"

import { Message } from "@opencode/ai"
import { and, desc, eq, sql } from "drizzle-orm"
import { and, desc, eq, inArray, sql } from "drizzle-orm"
import { Cause, Effect, Exit, FiberMap, Layer } from "effect"
import { Database } from "../../database/database.js"
import { Bus } from "../../bus.js"
Expand Down Expand Up @@ -35,6 +35,9 @@ import { MAX_STEPS_PROMPT } from "./max-steps.js"
const CONTINUE_AFTER_INCOMPLETE_STREAM =
"The previous response was interrupted. Continue from where you left off without repeating completed content."

const CONTINUE_AFTER_OUTPUT_LIMIT =
"Your last response hit the output token limit (stop reason: length). Do not apologize, recap, or repeat yourself. Break the remaining work into smaller pieces."

const layer = Layer.effect(
Service,
Effect.gen(function* () {
Expand Down Expand Up @@ -284,6 +287,36 @@ const layer = Layer.effect(
yield* bus.publish(SessionEvent.Synthetic, { sessionID, text: CONTINUE_AFTER_INCOMPLETE_STREAM })
assistantMessageID = SessionMessage.ID.create()
}),
OutputLimit: Effect.fn("SessionRunner.continueOutputLimit")(function* (outcome) {
// The current response is included. User input or any other finish breaks the streak.
const rows = yield* db
.select()
.from(SessionMessageTable)
.where(
and(
eq(SessionMessageTable.session_id, sessionID),
inArray(SessionMessageTable.type, ["user", "assistant"]),
),
)
.orderBy(desc(SessionMessageTable.seq))
.limit(3)
.all()
.pipe(Effect.orDie)
const recent = yield* Effect.forEach(rows, SessionHistory.decodeMessageRow)
const exhausted =
recent.length === 3 &&
recent.every((message) => message.type === "assistant" && message.finish === "length")
if (exhausted)
return yield* new StepFailedError({
error: {
type: "output-limit",
message: "Response still truncated after two output token limit continuations",
},
})
if (outcome.needsContinuation) return true
yield* bus.publish(SessionEvent.Synthetic, { sessionID, text: CONTINUE_AFTER_OUTPUT_LIMIT })
assistantMessageID = SessionMessage.ID.create()
}),
Compacted: Effect.fnUntraced(function* () {
recoverOverflow = false
assistantMessageID = SessionMessage.ID.create()
Expand Down
81 changes: 66 additions & 15 deletions packages/core/src/session/runner/publish-llm-event.ts
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,32 @@ type Input = {
const asRecord = (value: unknown): Record<string, unknown> =>
typeof value === "object" && value !== null && !Array.isArray(value) ? (value as Record<string, unknown>) : { value }

const INPUT_EXCERPT_LIMIT = 2048

type InputExcerpt = {
readonly text: string
readonly length: number
}

type InputFailureReason = "malformed" | "length" | "incomplete"

const INPUT_FAILURES = {
malformed: {
type: "tool.input-json",
message:
"The tool input could not be parsed as JSON, so the call was not executed. Reissue the call with valid JSON arguments.",
},
length: {
type: "tool.input-incomplete",
message:
"The tool call was not executed because the output token limit (stop reason: length) was reached before you could complete the arguments. Break large payloads into smaller tool calls and reissue the call with complete arguments.",
},
incomplete: {
type: "tool.input-incomplete",
message: "Tool call arguments were not completed and were not executed. Reissue the call with complete arguments.",
},
} satisfies Record<InputFailureReason, SessionError.Error>

/** Immutable fold of the durable facts a step's writer has recorded so far. */
export interface StepRecord {
/** The model produced visible output this attempt, which bars transparent retries and overflow recovery. */
Expand Down Expand Up @@ -83,6 +109,8 @@ export const createLLMEventPublisher = (bus: Pick<Bus.Interface, "publish">, inp
settled: boolean
providerExecuted: boolean
progress?: Tool.Metadata
input?: InputExcerpt
malformed?: boolean
}
const tools = new Map<string, ToolState>()
const failureSnapshot = (tool: { readonly progress?: Tool.Metadata }, metadata?: Tool.Metadata) => {
Expand Down Expand Up @@ -254,6 +282,7 @@ export const createLLMEventPublisher = (bus: Pick<Bus.Interface, "publish">, inp
Effect.gen(function* () {
const tool = tools.get(id)
if (!tool) return yield* Effect.die(new Error(`Tool input end before start: ${id}`))
tool.input = inputExcerpt(value)
yield* bus.publish(SessionEvent.Tool.Input.Ended, {
sessionID: input.sessionID,
assistantMessageID,
Expand Down Expand Up @@ -324,21 +353,21 @@ export const createLLMEventPublisher = (bus: Pick<Bus.Interface, "publish">, inp
if (tool.name !== event.name)
return yield* Effect.die(new Error(`Tool input name changed for ${event.id}: ${tool.name} -> ${event.name}`))
if (toolInput.has(event.id)) yield* endToolInput(event, event.raw)
tool.settled = true
yield* bus.publish(SessionEvent.Tool.Failed, {
sessionID: input.sessionID,
assistantMessageID,
id: event.id,
error: {
type: "tool.input-json",
message: "Tool call arguments were malformed JSON and were not executed. Retry with valid JSON.",
},
...failureSnapshot(tool),
executed: false,
})
tool.input = inputExcerpt(event.raw)
tool.malformed = true
})

const flush = Effect.fn("SessionRunner.flush")(flushFragments)
const failToolInput = (id: string, tool: ToolState, reason: Exclude<InputFailureReason, "length">) =>
failTool(id, inputFailure(tool.input, stepSettlement?.finish === "length" ? "length" : reason))

const flush = Effect.fn("SessionRunner.flush")(function* () {
yield* flushFragments()
// The finish reason distinguishes invalid JSON from arguments cut off by the output limit.
for (const [id, tool] of tools) {
if (!tool.malformed || tool.settled) continue
yield* failToolInput(id, tool, "malformed")
}
})

const failTool = Effect.fnUntraced(function* (id: string, error: SessionError.Error, metadata?: Tool.Metadata) {
const tool = tools.get(id)
Expand All @@ -363,11 +392,18 @@ export const createLLMEventPublisher = (bus: Pick<Bus.Interface, "publish">, inp
for (const [id, tool] of tools) {
if (tool.settled || (mode === "hosted" && !tool.providerExecuted) || (mode === "uncalled" && tool.called))
continue
failed = (yield* failTool(id, error)) || failed
if (yield* failTool(id, error)) failed = true
}
return failed
})

const failIncompleteTools = Effect.fn("SessionRunner.failIncompleteTools")(function* () {
for (const [id, tool] of tools) {
if (tool.called || tool.settled) continue
yield* failToolInput(id, tool, "incomplete")
}
})

const failAssistant = Effect.fnUntraced(function* (error: SessionError.Error) {
yield* flush()
yield* failTools(error, "uncalled")
Expand Down Expand Up @@ -534,7 +570,7 @@ export const createLLMEventPublisher = (bus: Pick<Bus.Interface, "publish">, inp
return
}
case "step-finish":
yield* flush()
yield* flushFragments()
if (stepSettlement) return yield* Effect.die(new Error("Duplicate step finish"))
stepSettlement = {
finish: event.reason.normalized,
Expand Down Expand Up @@ -604,6 +640,7 @@ export const createLLMEventPublisher = (bus: Pick<Bus.Interface, "publish">, inp
failTool,
publishStepFailure,
failUnsettledTools,
failIncompleteTools,
hasProviderError: () => providerFailed,
hasStarted: () => stepStarted,
/** Immutable snapshot of everything recorded for this step so far. */
Expand All @@ -621,3 +658,17 @@ export const createLLMEventPublisher = (bus: Pick<Bus.Interface, "publish">, inp
streamed,
}
}

function inputExcerpt(value: string): InputExcerpt {
return { text: value.slice(0, INPUT_EXCERPT_LIMIT), length: value.length }
}

function inputFailure(input: InputExcerpt | undefined, reason: InputFailureReason) {
const failure = INPUT_FAILURES[reason]
if (!input) return failure
const truncated = input.length > input.text.length ? "\n[truncated]" : ""
return {
...failure,
message: `${failure.message}\n\nInput excerpt (first ${input.text.length} of ${input.length} characters):\n${input.text}${truncated}`,
}
}
14 changes: 5 additions & 9 deletions packages/core/src/session/runner/step.ts
Original file line number Diff line number Diff line change
Expand Up @@ -38,6 +38,7 @@ export type Outcome = Data.TaggedEnum<{
}
RecoverFull: {}
Compacted: {}
OutputLimit: { readonly needsContinuation: boolean }
}>
export const Outcome = Data.taggedEnum<Outcome>()

Expand All @@ -61,11 +62,6 @@ interface Input {
const TOOLS_INTERRUPTED = { type: "aborted", message: "Tool execution interrupted" } as const
const STEP_INTERRUPTED = { type: "aborted", message: "Step interrupted" } as const
const RESULT_MISSING = { type: "tool.result-missing", message: "Provider did not return a tool result" } as const
const INPUT_INCOMPLETE = {
type: "tool.input-incomplete",
message:
"Tool call arguments were not completed and were not executed. Re-issue the tool call with complete arguments.",
} as const

/** Captures Location-scoped dependencies without introducing another service or execution loop. */
export const make = Effect.gen(function* () {
Expand Down Expand Up @@ -225,7 +221,7 @@ export const make = Effect.gen(function* () {
if (llmError || (Exit.isSuccess(stream) && !recorded.providerFailed)) {
const missing = yield* publisher.failUnsettledTools(RESULT_MISSING, "hosted")
if (missing && !llmError && !recorded.finish) yield* publisher.failAssistant(RESULT_MISSING)
yield* publisher.failUnsettledTools(INPUT_INCOMPLETE, "uncalled")
yield* publisher.failIncompleteTools()
}

const record = publisher.record()
Expand Down Expand Up @@ -274,9 +270,9 @@ export const make = Effect.gen(function* () {
if (tools.interrupted && tools.failure) return yield* Effect.failCause(tools.failure)
if (tools.interrupted && Exit.isFailure(joined)) return yield* Effect.failCause(joined.cause)
if (record.failure) return yield* new StepFailedError({ error: record.failure })
return Outcome.Completed({
needsContinuation: input.prepared.request.toolChoice?.type !== "none" && record.needsContinuation,
})
const needsContinuation = input.prepared.request.toolChoice?.type !== "none" && record.needsContinuation
if (record.finish?.finish === "length") return Outcome.OutputLimit({ needsContinuation })
return Outcome.Completed({ needsContinuation })
}),
)
}, Effect.scoped)
Expand Down
Loading
Loading