diff --git a/packages/core/src/session/runner/llm.ts b/packages/core/src/session/runner/llm.ts index 3ffe4ed6e900..a6dd35c80e3d 100644 --- a/packages/core/src/session/runner/llm.ts +++ b/packages/core/src/session/runner/llm.ts @@ -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" @@ -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* () { @@ -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() diff --git a/packages/core/src/session/runner/publish-llm-event.ts b/packages/core/src/session/runner/publish-llm-event.ts index c3246b287c6f..951a435f4ad1 100644 --- a/packages/core/src/session/runner/publish-llm-event.ts +++ b/packages/core/src/session/runner/publish-llm-event.ts @@ -28,6 +28,32 @@ type Input = { const asRecord = (value: unknown): Record => typeof value === "object" && value !== null && !Array.isArray(value) ? (value as Record) : { 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 + /** 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. */ @@ -83,6 +109,8 @@ export const createLLMEventPublisher = (bus: Pick, inp settled: boolean providerExecuted: boolean progress?: Tool.Metadata + input?: InputExcerpt + malformed?: boolean } const tools = new Map() const failureSnapshot = (tool: { readonly progress?: Tool.Metadata }, metadata?: Tool.Metadata) => { @@ -254,6 +282,7 @@ export const createLLMEventPublisher = (bus: Pick, 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, @@ -324,21 +353,21 @@ export const createLLMEventPublisher = (bus: Pick, 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) => + 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) @@ -363,11 +392,18 @@ export const createLLMEventPublisher = (bus: Pick, 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") @@ -534,7 +570,7 @@ export const createLLMEventPublisher = (bus: Pick, 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, @@ -604,6 +640,7 @@ export const createLLMEventPublisher = (bus: Pick, inp failTool, publishStepFailure, failUnsettledTools, + failIncompleteTools, hasProviderError: () => providerFailed, hasStarted: () => stepStarted, /** Immutable snapshot of everything recorded for this step so far. */ @@ -621,3 +658,17 @@ export const createLLMEventPublisher = (bus: Pick, 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}`, + } +} diff --git a/packages/core/src/session/runner/step.ts b/packages/core/src/session/runner/step.ts index afe188f52516..ebdc5581c429 100644 --- a/packages/core/src/session/runner/step.ts +++ b/packages/core/src/session/runner/step.ts @@ -38,6 +38,7 @@ export type Outcome = Data.TaggedEnum<{ } RecoverFull: {} Compacted: {} + OutputLimit: { readonly needsContinuation: boolean } }> export const Outcome = Data.taggedEnum() @@ -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* () { @@ -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() @@ -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) diff --git a/packages/core/test/session-runner.test.ts b/packages/core/test/session-runner.test.ts index 62a4b7c1be4c..dc34d76ef0bf 100644 --- a/packages/core/test/session-runner.test.ts +++ b/packages/core/test/session-runner.test.ts @@ -5541,6 +5541,99 @@ describe("SessionRunnerLLM", () => { expect(s.requests).toHaveLength(2) }) + for (const type of ["text", "reasoning"] as const) { + scenario(`continues ${type} after an output token limit`, function* (s) { + const nudge = + "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 partial = + type === "text" + ? [ + LLMEvent.textStart({ id: "partial" }), + LLMEvent.textDelta({ id: "partial", text: "Partial" }), + LLMEvent.textEnd({ id: "partial" }), + ] + : [ + LLMEvent.reasoningStart({ id: "partial" }), + LLMEvent.reasoningDelta({ id: "partial", text: "Partial" }), + LLMEvent.reasoningEnd({ id: "partial" }), + ] + yield* s.llm.push( + TestLLM.complete({ reason: { normalized: "length" } }, ...partial), + TestLLM.text("Finished", "finished"), + ) + + yield* s.runPrompt("Complete the response") + + expect(s.requests).toHaveLength(2) + expect(s.requests[1]?.messages.at(-2)).toMatchObject({ + role: "assistant", + content: [expect.objectContaining({ text: "Partial" })], + }) + expect(s.requests[1]?.messages.at(-1)).toMatchObject({ + role: "user", + content: [{ type: "text", text: nudge }], + }) + expect(yield* s.context).toMatchObject([ + { type: "user" }, + Expected.assistant({ finish: "length" }, [ + type === "text" ? Expected.text("Partial") : Expected.reasoning("Partial"), + ]), + { type: "synthetic", text: nudge }, + Expected.assistant({ finish: "stop" }, [Expected.text("Finished")]), + ]) + expect(yield* recordedEventTypes(sessionID)).not.toContain("session.retry.scheduled.1") + }) + } + + for (const tools of ["hosted", "unfinished"] as const) { + scenario(`caps output token limit continuations and resets for new input (${tools})`, function* (s) { + const toolEvents = + tools === "hosted" + ? [ + hostedCall("hosted-search", "Search"), + LLMEvent.toolResult({ + id: "hosted-search", + name: "web_search", + providerExecuted: true, + result: { type: "json", value: [] }, + }), + ] + : [ + LLMEvent.toolInputStart({ id: "call-incomplete", name: "echo" }), + LLMEvent.toolInputDelta({ id: "call-incomplete", name: "echo", text: '{"text":"partial' }), + ] + const truncated = () => + TestLLM.complete( + { reason: { normalized: "length" } }, + ...toolEvents, + LLMEvent.textStart({ id: "partial" }), + LLMEvent.textDelta({ id: "partial", text: "Partial" }), + LLMEvent.textEnd({ id: "partial" }), + ) + const syntheticCount = s.context.pipe( + Effect.map((messages) => messages.filter((message) => message.type === "synthetic").length), + ) + const expectedNudges = (count: number) => (tools === "hosted" ? count : 0) + yield* s.llm.push(...Array.from({ length: 3 }, truncated)) + + expect(yield* s.runPrompt("Keep going").pipe(Effect.flip)).toMatchObject({ error: { type: "output-limit" } }) + expect(s.requests).toHaveLength(3) + expect(yield* syntheticCount).toBe(expectedNudges(2)) + expect(s.executions).toEqual([]) + + yield* replaySessionProjection(sessionID) + yield* s.llm.push(truncated()) + expect(yield* s.resume.pipe(Effect.flip)).toMatchObject({ error: { type: "output-limit" } }) + expect(s.requests).toHaveLength(4) + expect(yield* syntheticCount).toBe(expectedNudges(2)) + + yield* s.llm.push(truncated(), TestLLM.text("Finished", "finished")) + yield* s.runPrompt("Try a new response") + expect(s.requests).toHaveLength(6) + expect(yield* syntheticCount).toBe(expectedNudges(3)) + }) + } + scenario("continues an incomplete stream after observable text", function* (s) { const failure = incompleteStream() yield* s.admit("Continue partial output") @@ -6046,103 +6139,127 @@ describe("SessionRunnerLLM", () => { expect(s.requests).toHaveLength(2) expect(s.executions).toEqual([]) + expect((yield* s.context).some((message) => message.type === "synthetic")).toBe(false) expect(requireAssistant(yield* s.context).content).toMatchObject([ { type: "tool", id: "call-incomplete", executed: false, - state: { status: "error", error: { type: "tool.input-incomplete" } }, + state: { + status: "error", + input: {}, + error: { + type: "tool.input-incomplete", + message: expect.stringContaining('Input excerpt (first 16 of 16 characters):\n{"text":"partial'), + }, + }, }, ]) }) - scenario("continues after malformed local tool input without exposing raw arguments", function* (s) { - const marker = "raw-malformed-marker" - const raw = `{"text":"${marker}` - yield* s.llm.push( - TestLLM.toolCalls( - LLMEvent.toolInputStart({ id: "call-malformed", name: "echo" }), - LLMEvent.toolInputDelta({ id: "call-malformed", name: "echo", text: raw }), - LLMEvent.toolInputEnd({ id: "call-malformed", name: "echo" }), - LLMEvent.toolInputError({ - id: "call-malformed", - name: "echo", - raw, - }), - ), - TestLLM.stop(), - ) + for (const finish of ["tool-calls", "length"] as const) { + scenario(`continues malformed local tool input with a bounded excerpt (${finish})`, function* (s) { + const marker = "raw-malformed-marker" + const raw = `{"text":"${marker}${"x".repeat(3000)}omitted-marker` + const guidance = + finish === "length" + ? "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." + : "The tool input could not be parsed as JSON, so the call was not executed. Reissue the call with valid JSON arguments." + const message = [ + guidance, + "", + `Input excerpt (first 2048 of ${raw.length} characters):`, + raw.slice(0, 2048), + "[truncated]", + ].join("\n") + yield* s.llm.push( + TestLLM.complete( + { reason: { normalized: finish } }, + LLMEvent.toolInputStart({ id: "call-malformed", name: "echo" }), + LLMEvent.toolInputDelta({ id: "call-malformed", name: "echo", text: raw }), + LLMEvent.toolInputEnd({ id: "call-malformed", name: "echo" }), + LLMEvent.toolInputError({ + id: "call-malformed", + name: "echo", + raw, + }), + ), + TestLLM.stop(), + ) - yield* s.runPrompt("Recover malformed tool input") + yield* s.runPrompt("Recover malformed tool input") - expect(s.requests).toHaveLength(2) - expect(s.executions).toEqual([]) - expect(JSON.stringify(s.requests[1])).not.toContain(marker) - expect(s.requests[1]?.messages).toEqual( - expect.arrayContaining([ - expect.objectContaining({ - role: "assistant", - content: expect.arrayContaining([ - expect.objectContaining({ type: "tool-call", id: "call-malformed", name: "echo", input: {} }), - ]), - }), - expect.objectContaining({ - role: "tool", - content: expect.arrayContaining([ - expect.objectContaining({ - type: "tool-result", - id: "call-malformed", - result: expect.objectContaining({ - type: "error", - value: expect.objectContaining({ - error: expect.objectContaining({ - message: "Tool call arguments were malformed JSON and were not executed. Retry with valid JSON.", + expect(s.requests).toHaveLength(2) + expect(s.executions).toEqual([]) + expect(JSON.stringify(s.requests[1])).toContain(marker) + expect(JSON.stringify(s.requests[1])).not.toContain("omitted-marker") + expect((yield* s.context).some((message) => message.type === "synthetic")).toBe(false) + expect(s.requests[1]?.messages).toEqual( + expect.arrayContaining([ + expect.objectContaining({ + role: "assistant", + content: expect.arrayContaining([ + expect.objectContaining({ type: "tool-call", id: "call-malformed", name: "echo", input: {} }), + ]), + }), + expect.objectContaining({ + role: "tool", + content: expect.arrayContaining([ + expect.objectContaining({ + type: "tool-result", + id: "call-malformed", + result: expect.objectContaining({ + type: "error", + value: expect.objectContaining({ + error: expect.objectContaining({ + message, + }), }), }), }), - }), - ]), - }), - ]), - ) - const context = yield* s.context - const failed = context.find( - (message): message is SessionMessage.Assistant => - message.type === "assistant" && message.content.some((item) => item.type === "tool"), - ) - expect(failed).toMatchObject({ - content: [ - Expected.failedTool( - { id: "call-malformed", executed: false }, - { - input: {}, - error: { - type: "tool.input-json", - message: "Tool call arguments were malformed JSON and were not executed. Retry with valid JSON.", + ]), + }), + ]), + ) + const context = yield* s.context + const failed = context.find( + (message): message is SessionMessage.Assistant => + message.type === "assistant" && message.content.some((item) => item.type === "tool"), + ) + expect(failed).toMatchObject({ + content: [ + Expected.failedTool( + { id: "call-malformed", executed: false }, + { + input: {}, + error: { + type: finish === "length" ? "tool.input-incomplete" : "tool.input-json", + message, + }, }, - }, - ), - ], - }) - if (!failed) throw new Error("Malformed tool assistant missing") - expect(failed.error).toBeUndefined() - expect(yield* recordedStepSettlementTypes(sessionID, failed.id)).toEqual([ - "session.step.started.1", - "session.tool.failed.2", - "session.step.ended.1", - ]) + ), + ], + }) + if (!failed) throw new Error("Malformed tool assistant missing") + expect(failed.error).toBeUndefined() + expect(yield* recordedStepSettlementTypes(sessionID, failed.id)).toEqual([ + "session.step.started.1", + "session.tool.failed.2", + "session.step.ended.1", + ]) - const durable = yield* s.db - .select({ type: EventTable.type, data: EventTable.data }) - .from(EventTable) - .where(eq(EventTable.aggregate_id, sessionID)) - .all() - .pipe(Effect.orDie) - expect(durable.find((event) => event.type === "session.tool.input.ended.1")?.data).toMatchObject({ - id: "call-malformed", - text: raw, + const durable = yield* s.db + .select({ type: EventTable.type, data: EventTable.data }) + .from(EventTable) + .where(eq(EventTable.aggregate_id, sessionID)) + .all() + .pipe(Effect.orDie) + expect(durable.find((event) => event.type === "session.tool.input.ended.1")?.data).toMatchObject({ + id: "call-malformed", + text: raw, + }) }) - }) + } scenario("settles a valid sibling before recovering malformed tool input", function* (s) { yield* s.admit("Run parallel tools")