From bd44e03ca08d41965d87b7b6f704654d0380d97e Mon Sep 17 00:00:00 2001 From: foxytanuki Date: Wed, 30 Sep 2026 06:21:55 +0000 Subject: [PATCH 1/4] fix(claude): order devin signatures before messages answers --- .../src/content/docs/guides/claude-code.md | 7 + scripts/test-layout/layout.json | 1 + src/claude/devin-output-order.ts | 76 +++++++ src/server/responses/run-turn-execution.ts | 11 +- structure/clients/claude-desktop.md | 17 ++ structure/providers-and-adapters.md | 3 + structure/transports/byte-accounting.md | 4 + .../claude-devin-output-order.test.ts | 192 ++++++++++++++++++ tests/fixtures/test-layout-expected.json | 1 + 9 files changed, 310 insertions(+), 2 deletions(-) create mode 100644 src/claude/devin-output-order.ts create mode 100644 tests/claude-integration/claude-devin-output-order.test.ts diff --git a/docs-site/src/content/docs/guides/claude-code.md b/docs-site/src/content/docs/guides/claude-code.md index 1d3394163fe..6bee6138cb5 100644 --- a/docs-site/src/content/docs/guides/claude-code.md +++ b/docs-site/src/content/docs/guides/claude-code.md @@ -7,6 +7,13 @@ opencodex serves `POST /v1/messages` (plus `count_tokens`) alongside `/v1/respon Code can use every routed provider — OAuth logins, account pools, key failover and sidecars included — with zero extra auth work. +On Devin routes (including SWE-2), text and tool calls wait for the upstream turn to complete +so its late reasoning signature can precede the answer. This prevents Claude Code's final +result from becoming empty; reasoning and keepalive progress still flow during generation. +The buffer shares the request's 32 MiB translation limit and cancellation stops the producer. +This output-order fix does not resolve Cognition's separate refusal of some generated system +text. It preserves system instructions and safety constraints. + For an Anthropic route on stored OAuth or an Anthropic API key, native Fast is available on `claude-opus-5-5`, `claude-opus-5`, and `claude-opus-4-8`: pick the model's `--fast` row (listed when Fast rows are enabled) or set `fastMode: true`. Claude Code's own `/fast` toggle is not diff --git a/scripts/test-layout/layout.json b/scripts/test-layout/layout.json index 2e8719fa6d9..182c88493b8 100644 --- a/scripts/test-layout/layout.json +++ b/scripts/test-layout/layout.json @@ -2,6 +2,7 @@ "version": 1, "root": "tests", "explicit": { + "claude-devin-output-order.test.ts": "claude-integration", "service-desktop-startup-health.test.ts": "service", "service-desktop-startup.test.ts": "service", "startup-health-packaged-probe.test.ts": "server", diff --git a/src/claude/devin-output-order.ts b/src/claude/devin-output-order.ts new file mode 100644 index 00000000000..3b107a82f0a --- /dev/null +++ b/src/claude/devin-output-order.ts @@ -0,0 +1,76 @@ +import type { AdapterEvent } from "../types"; +import { + isTranslatorBudgetExceededError, + releaseTranslatedEvent, + retainTranslatedEvent, + type TranslatorBudget, +} from "../lib/translator-budget"; + +/** + * Cognition signs reasoning after text/tools. Claude Code treats a trailing empty + * thinking block as its final result, so Messages needs reasoning before content. + * Hold only this physical turn's semantic events; keep progress and reasoning live. + */ +export function createDevinMessagesOutputOrder( + emit: (event: AdapterEvent) => void, + budget: TranslatorBudget, + signal: AbortSignal, + abort: () => void, +) { + const held: AdapterEvent[] = []; + let ended = false; + const release = () => { + for (const event of held) releaseTranslatedEvent(event, budget); + held.length = 0; + }; + const cancel = () => { ended = true; release(); }; + signal.addEventListener("abort", cancel, { once: true }); + if (signal.aborted) cancel(); + const flush = () => { + if (ended) return; + try { + for (const event of held) { + emit(event); + releaseTranslatedEvent(event, budget); + } + } finally { release(); } + }; + return { + emit(event: AdapterEvent) { + if (ended) return; + if (event.type === "done" || event.type === "error" || event.type === "incomplete") { + flush(); + ended = true; + emit(event); + return; + } + if (event.type === "heartbeat" || event.type === "thinking_delta" + || event.type === "thinking_signature" || event.type === "redacted_thinking" + || event.type === "reasoning_raw_delta" || event.type === "kiro_redacted_reasoning") { + emit(event); + return; + } + const copy = { ...event } as AdapterEvent; + try { + retainTranslatedEvent(copy, budget, held.at(-1)); + held.push(copy); + } catch (error) { + if (!isTranslatorBudgetExceededError(error)) throw error; + ended = true; + release(); + emit({ type: "error", status: 413, errorType: "request_too_large", + code: "translation_buffer_limit", message: error.message }); + abort(); + return; + } + // Progress still feeds the stall watchdog while text/tool output waits for + // the signature. Preserve original replay-unsafe heartbeats above verbatim. + emit({ type: "heartbeat" }); + }, + flush, + dispose() { + signal.removeEventListener("abort", cancel); + cancel(); + }, + }; +} diff --git a/src/server/responses/run-turn-execution.ts b/src/server/responses/run-turn-execution.ts index 11cfa90d7fe..09694442137 100644 --- a/src/server/responses/run-turn-execution.ts +++ b/src/server/responses/run-turn-execution.ts @@ -50,6 +50,7 @@ import { undeclaredToolCallMessage } from "../responses-undeclared-tool-guard"; import { planWebSearch } from "../../web-search"; import { runTurnWebSearchInitialParsed, runTurnWebSearchLoop } from "../../web-search/run-turn-loop"; import { WEB_SEARCH_TOOL_NAME } from "../../web-search/synthetic-tool"; +import { createDevinMessagesOutputOrder } from "../../claude/devin-output-order"; // LOCAL PATCH (runturn-websearch): top-level fields route binding or the // adapter itself may write during a turn. Iteration-local `turnParsed` objects @@ -229,6 +230,7 @@ export async function executeResponsesRunTurn( options.onCompactionRecoveryAdapterEvent?.(event); targetQueue.push(event); }; + let messagesOutput: ReturnType | undefined; try { if (!pacingSlotAcquired) { pacingSlot = await waitForProviderRequestSlot(route.providerName, route.provider, route.modelId, runTurnAbort.signal); @@ -261,6 +263,9 @@ export async function executeResponsesRunTurn( turnScopedPacing: true, }, ); + if (inboundWire === "anthropic" && transportState.runTurnAdapter.name === "devin" && !routedCompaction) { + messagesOutput = createDevinMessagesOutputOrder(emit, translatorBudget, runTurnAbort.signal, () => runTurnAbort.abort()); + } await transportState.runTurnAdapter.runTurn?.( turnParsed, { @@ -282,7 +287,7 @@ export async function executeResponsesRunTurn( ), onRecoveryWithheld: noteAdapterRecoveryWithheld, }, - emit, + messagesOutput?.emit ?? emit, ); // LOCAL PATCH (runturn-websearch): adapters may write conversation/ // continuation state onto the object they received; merge it back so @@ -296,7 +301,7 @@ export async function executeResponsesRunTurn( Object.assign(parsed, routeState); } } catch (err) { - emit(err instanceof RequestPacingQueueOverloadError + (messagesOutput?.emit ?? emit)(err instanceof RequestPacingQueueOverloadError ? { type: "error", status: 429, @@ -321,6 +326,8 @@ export async function executeResponsesRunTurn( message: err instanceof Error ? err.message : String(err), }); } finally { + messagesOutput?.flush(); + messagesOutput?.dispose(); releaseProviderRequestSlot(pacingSlot); // Cursor assigns a stable conversation id inside runTurn on the first headerless // turn; backfill so Logs can filter/total that opening request (#330 / #522). diff --git a/structure/clients/claude-desktop.md b/structure/clients/claude-desktop.md index e256089a2a8..16c335050ff 100644 --- a/structure/clients/claude-desktop.md +++ b/structure/clients/claude-desktop.md @@ -29,6 +29,23 @@ Native OpenAI pool routing also accepts [Orca-linked accounts](../codex-home.md#orca-source-owned-account-import), whose source resolution belongs to the shared account store. The import CLI adds pool rows independently of Desktop profiles. +## Devin Messages output ordering + +`src/claude/devin-output-order.ts` orders each physical Devin turn at the emitter in +`src/server/responses/run-turn-execution.ts` when the original inbound wire is Anthropic +Messages. Provider names may be customized; the selected adapter determines applicability. +Text and tool events wait for that turn's terminal so Cognition's late reasoning signature +precedes them. Claude Code therefore receives a final text or tool block rather than an empty +signature-only thinking block. Reasoning and transport progress remain live; answer text and +tool dispatch incur turn-completion latency. Responses and Chat retain their original ordering, +and routed compaction is excluded. + +Retained events share the request translator budget and are released on terminal, cancellation, +overflow, or adapter EOF. Overflow emits one typed `translation_buffer_limit` error and aborts +the producer. Error and incomplete terminals retain partial output and their original usage. +Ordering occurs before hosted-search interception, independently for each physical iteration, +so one iteration's signature cannot be attached to another iteration's answer. + ## Desktop modes: gateway and first-party `src/claude/desktop-first-party.ts` owns the Desktop mode contract. Two modes exist and are diff --git a/structure/providers-and-adapters.md b/structure/providers-and-adapters.md index b8cbdb56e94..71bfe6793f8 100644 --- a/structure/providers-and-adapters.md +++ b/structure/providers-and-adapters.md @@ -1,5 +1,8 @@ # Providers And Adapters +Devin Messages follows the [per-turn output ordering contract](clients/claude-desktop.md#devin-messages-output-ordering), +preserving late signatures before text/tools without changing Responses or Chat ordering. + Anthropic account pause, model routes, and quota labels follow the [Anthropic account-pool contract](providers/anthropic-account-pool.md). diff --git a/structure/transports/byte-accounting.md b/structure/transports/byte-accounting.md index f514bcdcb5e..bab76ec8acd 100644 --- a/structure/transports/byte-accounting.md +++ b/structure/transports/byte-accounting.md @@ -56,6 +56,10 @@ lifecycle, cancellation races, protocol envelopes, and the real HTTP admission b ## Stream-buffer accounting +Devin's [Messages ordering buffer](../clients/claude-desktop.md#devin-messages-output-ordering) charges retained +semantic events to the shared translator budget until transfer to the independently bounded +adapter queue. Cancellation and overflow release held events before producer shutdown. + `src/web-search/run-turn-loop.ts` charges retained iteration events and generated replay history to the request translator budget. Each owner releases its own reservations on completion, error, cancellation or consumer closure; a buffer-limit failure terminates without another search. diff --git a/tests/claude-integration/claude-devin-output-order.test.ts b/tests/claude-integration/claude-devin-output-order.test.ts new file mode 100644 index 00000000000..57591c93bd2 --- /dev/null +++ b/tests/claude-integration/claude-devin-output-order.test.ts @@ -0,0 +1,192 @@ +import { afterAll, afterEach, beforeEach, describe, expect, mock, test } from "bun:test"; +import { createDevinMessagesOutputOrder } from "../../src/claude/devin-output-order"; +import { collectAnthropicMessage, responsesSseToAnthropicSse } from "../../src/claude/outbound"; +import { bridgeToResponsesSSE } from "../../src/bridge"; +import { createTestTranslatorBudget } from "../helpers/translator-budget"; +import { encodeDevinSignature } from "../../src/adapters/devin/reasoning-signature"; +import { mapOcxMessagesToDevin } from "../../src/adapters/devin"; +import { messagesToResponsesTranslation } from "../../src/protocols/codecs/messages"; +import { parseRequest } from "../../src/responses/parser"; +import type { AdapterEvent, OcxConfig, OcxProviderConfig } from "../../src/types"; +import type { ProviderAdapter } from "../../src/adapters/base"; +import { createTempHome } from "../helpers/temp-home"; +import { acquireOwnedSpendHome } from "../helpers/owned-spend-home"; +import { installIsolatedCodexHome } from "../helpers/isolated-codex-home"; + +const signature = encodeDevinSignature("sealed.v1.synthetic-attestation", "sealed"); +const usage = { inputTokens: 12, outputTokens: 3, totalTokens: 15 }; +const terminal: AdapterEvent = { type: "done", usage }; +let upstreamEvents: AdapterEvent[] = []; +const resolver = await import("../../src/server/adapter-resolve"); +const originalResolver = { ...resolver }; +mock.module("../../src/server/adapter-resolve", () => ({ ...originalResolver, + resolveAdapter(provider: OcxProviderConfig, cache?: "none" | "short" | "long") { + if (provider.adapter !== "devin") return originalResolver.resolveAdapter(provider, cache); + return { + name: "devin", + buildRequest: () => ({ url: provider.baseUrl, method: "POST", headers: {}, body: "" }), + async *parseStream() { yield terminal; }, + async runTurn(_parsed, _incoming, emit) { upstreamEvents.forEach(emit); }, + } satisfies ProviderAdapter; + }, +})); +const { handleClaudeMessages } = await import("../../src/server/claude-messages"); +const { handleResponses } = await import("../../src/server/responses"); +afterAll(() => { mock.module("../../src/server/adapter-resolve", () => originalResolver); }); + +let home: ReturnType; +let codexHome: ReturnType; +let releaseSpend: () => void; +beforeEach(() => { + home = createTempHome("ocx-claude-devin-output-"); + codexHome = installIsolatedCodexHome("ocx-claude-devin-codex-"); + releaseSpend = acquireOwnedSpendHome(); + upstreamEvents = []; +}); +afterEach(() => { releaseSpend(); codexHome.restore(); home.remove(); }); + +function ordered(events: AdapterEvent[]) { + const budget = createTestTranslatorBudget(); + const abort = new AbortController(); + const output: AdapterEvent[] = []; + const order = createDevinMessagesOutputOrder(event => output.push(event), budget, abort.signal, () => abort.abort()); + try { events.forEach(order.emit); order.flush(); } finally { order.dispose(); } + return { output, budget }; +} + +async function message(events: AdapterEvent[]) { + const { output, budget } = ordered(events); + const source = (async function* () { yield* output; })(); + const bridged = bridgeToResponsesSSE(source, "swe-2", undefined, undefined, undefined, undefined, 0, { translatorBudget: budget }); + return collectAnthropicMessage(responsesSseToAnthropicSse(bridged, "devin/swe-2", { + translatorBudget: budget, pingIntervalMs: 0, + }), "devin/swe-2", budget); +} + +describe("Devin Messages late signatures", () => { + test("late signature keeps the answer last, signed reasoning intact, and exact usage", async () => { + const result = await message([ + { type: "thinking_delta", thinking: "A thought" }, + { type: "text_delta", text: "O" }, { type: "text_delta", text: "K" }, + { type: "thinking_signature", signature }, terminal, + ]); + expect(result.content).toEqual([ + { type: "thinking", thinking: "A thought", signature }, { type: "text", text: "OK" }, + ]); + expect(result.stop_reason).toBe("end_turn"); + expect(result.usage).toMatchObject({ input_tokens: 12, output_tokens: 3 }); + }); + + test("signature-only turns preserve both the signature and final text", async () => { + const result = await message([ + { type: "text_delta", text: "OK" }, { type: "thinking_signature", signature }, terminal, + ]); + expect(result.content).toEqual([ + { type: "thinking", thinking: "", signature }, { type: "text", text: "OK" }, + ]); + }); + + test("a signed tool turn survives Messages inbound replay to the Devin prompt", async () => { + const result = await message([ + { type: "thinking_delta", thinking: "Read the file" }, + { type: "tool_call_start", id: "call_read", name: "Read" }, + { type: "tool_call_delta", arguments: '{"file_path":"/synthetic/file.txt"}' }, + { type: "tool_call_end" }, { type: "thinking_signature", signature }, terminal, + ]); + expect(result.stop_reason).toBe("tool_use"); + expect(result.content.at(-1)).toMatchObject({ type: "tool_use", id: "call_read", name: "Read", + input: { file_path: "/synthetic/file.txt" } }); + const translated = messagesToResponsesTranslation({ model: "devin/swe-2", max_tokens: 64, + messages: [{ role: "user", content: "Read" }, { role: "assistant", content: result.content }, + { role: "user", content: [{ type: "tool_result", tool_use_id: "call_read", content: "marker" }] }], + }, undefined, createTestTranslatorBudget()); + const history = mapOcxMessagesToDevin(parseRequest(translated.body)); + expect(history.find(row => row.role === "assistant")).toMatchObject({ + thinking: "Read the file", signature: "sealed.v1.synthetic-attestation", signature_type: "sealed", + }); + expect(history.find(row => row.role === "tool")?.content).toBe("marker"); + }); + + test("independently signed reasoning blocks retain their associations", async () => { + const result = await message([ + { type: "thinking_delta", thinking: "First" }, { type: "thinking_signature", signature: "sig-first" }, + { type: "text_delta", text: "answer" }, + { type: "thinking_delta", thinking: "Second" }, { type: "thinking_signature", signature: "sig-second" }, terminal, + ]); + expect(result.content).toEqual([ + { type: "thinking", thinking: "First", signature: "sig-first" }, + { type: "thinking", thinking: "Second", signature: "sig-second" }, { type: "text", text: "answer" }, + ]); + }); + + for (const ending of [terminal, { type: "incomplete", reason: "max_tokens", usage }, + { type: "error", status: 502, message: "synthetic reset", usage }] as AdapterEvent[]) { + test(`${ending.type} preserves partial content and the original terminal`, () => { + const { output, budget } = ordered([{ type: "text_delta", text: "partial" }, ending]); + expect(output.filter(event => event.type !== "heartbeat")).toEqual([{ type: "text_delta", text: "partial" }, ending]); + expect(budget.snapshot().currentBytes).toBe(0); + }); + } + + test("holding output keeps progress live and cancellation releases it immediately", () => { + const budget = createTestTranslatorBudget(); + const abort = new AbortController(); + const output: AdapterEvent[] = []; + const order = createDevinMessagesOutputOrder(event => output.push(event), budget, abort.signal, () => abort.abort()); + order.emit({ type: "text_delta", text: "still generating" }); + expect(output).toEqual([{ type: "heartbeat" }]); + expect(budget.snapshot().currentBytes).toBeGreaterThan(0); + order.emit({ type: "heartbeat", replayUnsafe: true }); + expect(output.at(-1)).toEqual({ type: "heartbeat", replayUnsafe: true }); + abort.abort(); + expect(budget.snapshot().currentBytes).toBe(0); + order.emit(terminal); + order.dispose(); + expect(output.some(event => event.type === "text_delta" || event.type === "done")).toBe(false); + }); + + test("overflow emits one typed error, aborts the producer, and releases held output", () => { + const budget = createTestTranslatorBudget({ maxTurnBytes: 160 }); + const abort = new AbortController(); + const output: AdapterEvent[] = []; + const order = createDevinMessagesOutputOrder(event => output.push(event), budget, abort.signal, () => abort.abort()); + order.emit({ type: "text_delta", text: "small" }); + order.emit({ type: "text_delta", text: "x".repeat(200) }); + order.emit(terminal); + expect(output.filter(event => event.type === "error")).toHaveLength(1); + expect(output.at(-1)).toMatchObject({ type: "error", status: 413, code: "translation_buffer_limit" }); + expect(abort.signal.aborted).toBe(true); + expect(budget.snapshot().currentBytes).toBe(0); + order.dispose(); + }); +}); + +const config = (): OcxConfig => ({ port: 0, defaultProvider: "cognition-custom", claudeCode: { enabled: true }, + providers: { "cognition-custom": { adapter: "devin", baseUrl: "https://synthetic.invalid", apiKey: "synthetic-key", models: ["swe-2"] } }, +}); + +for (const stream of [true, false]) { + test(`actual Messages ingress orders a renamed Devin provider, stream=${stream}`, async () => { + upstreamEvents = [{ type: "text_delta", text: "OK" }, { type: "thinking_signature", signature }, terminal]; + const response = await handleClaudeMessages(new Request("http://localhost/v1/messages", { + method: "POST", headers: { "content-type": "application/json" }, + body: JSON.stringify({ model: "cognition-custom/swe-2", max_tokens: 64, stream, + messages: [{ role: "user", content: "Reply OK" }] }), + }), config(), {}); + expect(response.status).toBe(200); + const result = stream ? await collectAnthropicMessage(response.body!, "cognition-custom/swe-2", createTestTranslatorBudget()) + : await response.json(); + expect(result.content.at(-1)).toEqual({ type: "text", text: "OK" }); + expect(result.content[0]).toEqual({ type: "thinking", thinking: "", signature }); + }); +} + +test("the Responses ingress retains incremental text before the late signature", async () => { + upstreamEvents = [{ type: "text_delta", text: "OK" }, { type: "thinking_signature", signature }, terminal]; + const response = await handleResponses(new Request("http://localhost/v1/responses", { + method: "POST", headers: { "content-type": "application/json" }, + body: JSON.stringify({ model: "cognition-custom/swe-2", input: "Reply OK", stream: true }), + }), config(), {}); + const text = await response.text(); + expect(text.indexOf("response.output_text.delta")).toBeLessThan(text.indexOf("\"type\":\"reasoning\"")); +}); diff --git a/tests/fixtures/test-layout-expected.json b/tests/fixtures/test-layout-expected.json index f021ccb19b0..e4860dcabe6 100644 --- a/tests/fixtures/test-layout-expected.json +++ b/tests/fixtures/test-layout-expected.json @@ -1,4 +1,5 @@ { + "claude-devin-output-order.test.ts": "claude-integration", "service-desktop-startup-health.test.ts": "service", "service-desktop-startup.test.ts": "service", "startup-health-packaged-probe.test.ts": "server", From eb76c9608e73558d2e37054df7eb3737b791fe70 Mon Sep 17 00:00:00 2001 From: foxytanuki Date: Wed, 30 Sep 2026 06:41:20 +0000 Subject: [PATCH 2/4] fix(claude): drain devin messages after raw-event preflight --- src/claude/devin-output-order.ts | 99 +++++----- src/server/responses/run-turn-execution.ts | 36 ++-- structure/clients/claude-desktop.md | 11 +- structure/transports/byte-accounting.md | 4 +- .../claude-devin-output-order.test.ts | 178 +++++++++++++++--- 5 files changed, 241 insertions(+), 87 deletions(-) diff --git a/src/claude/devin-output-order.ts b/src/claude/devin-output-order.ts index 3b107a82f0a..555f9b87111 100644 --- a/src/claude/devin-output-order.ts +++ b/src/claude/devin-output-order.ts @@ -9,46 +9,56 @@ import { /** * Cognition signs reasoning after text/tools. Claude Code treats a trailing empty * thinking block as its final result, so Messages needs reasoning before content. - * Hold only this physical turn's semantic events; keep progress and reasoning live. + * Consume after raw-event preflight and drain on demand, avoiding a queue burst. */ -export function createDevinMessagesOutputOrder( - emit: (event: AdapterEvent) => void, +export async function* orderDevinMessagesOutput( + source: AsyncIterable, budget: TranslatorBudget, signal: AbortSignal, - abort: () => void, -) { - const held: AdapterEvent[] = []; - let ended = false; + abortProducer: () => void, +): AsyncGenerator { + const held: Array = []; + let cancelled = signal.aborted; + let terminalDelivered = false; + let delivering: AdapterEvent | undefined; const release = () => { - for (const event of held) releaseTranslatedEvent(event, budget); + for (const event of held) if (event) releaseTranslatedEvent(event, budget); held.length = 0; + if (delivering) releaseTranslatedEvent(delivering, budget); }; - const cancel = () => { ended = true; release(); }; + const cancel = () => { cancelled = true; release(); }; signal.addEventListener("abort", cancel, { once: true }); - if (signal.aborted) cancel(); - const flush = () => { - if (ended) return; - try { - for (const event of held) { - emit(event); - releaseTranslatedEvent(event, budget); + async function* drain() { + for (let index = 0; index < held.length && !cancelled; index++) { + delivering = held[index]; + held[index] = undefined; + try { + if (delivering) { + terminalDelivered = delivering.type === "done" || delivering.type === "error" || delivering.type === "incomplete"; + yield delivering; + } } - } finally { release(); } - }; - return { - emit(event: AdapterEvent) { - if (ended) return; - if (event.type === "done" || event.type === "error" || event.type === "incomplete") { - flush(); - ended = true; - emit(event); - return; + finally { + if (delivering) releaseTranslatedEvent(delivering, budget); + delivering = undefined; + } + } + held.length = 0; + } + try { + for await (const event of source) { + const terminal = event.type === "done" || event.type === "error" || event.type === "incomplete"; + if (cancelled) { + // The adapter's cancellation terminal retains its measured usage. Client + // stream cancellation instead returns this iterator and stops consumption. + if (terminal) { yield event; return; } + continue; } if (event.type === "heartbeat" || event.type === "thinking_delta" || event.type === "thinking_signature" || event.type === "redacted_thinking" || event.type === "reasoning_raw_delta" || event.type === "kiro_redacted_reasoning") { - emit(event); - return; + yield event; + continue; } const copy = { ...event } as AdapterEvent; try { @@ -56,21 +66,26 @@ export function createDevinMessagesOutputOrder( held.push(copy); } catch (error) { if (!isTranslatorBudgetExceededError(error)) throw error; - ended = true; release(); - emit({ type: "error", status: 413, errorType: "request_too_large", - code: "translation_buffer_limit", message: error.message }); - abort(); + // Stop only the producer: aborting the request signal here would let the + // hosted-search loop replace the typed overflow with a client-cancel error. + abortProducer(); + yield { type: "error", status: 413, errorType: "request_too_large", + code: "translation_buffer_limit", message: error.message }; return; } - // Progress still feeds the stall watchdog while text/tool output waits for - // the signature. Preserve original replay-unsafe heartbeats above verbatim. - emit({ type: "heartbeat" }); - }, - flush, - dispose() { - signal.removeEventListener("abort", cancel); - cancel(); - }, - }; + if (terminal) { + yield* drain(); + if (cancelled && !terminalDelivered) yield event; + return; + } + // Raw preflight already observed real output; feed the stream watchdog while + // the answer waits for its signature. Original heartbeats stay verbatim. + yield { type: "heartbeat" }; + } + if (!cancelled) yield* drain(); + } finally { + signal.removeEventListener("abort", cancel); + release(); + } } diff --git a/src/server/responses/run-turn-execution.ts b/src/server/responses/run-turn-execution.ts index 09694442137..d74aeaa5fb4 100644 --- a/src/server/responses/run-turn-execution.ts +++ b/src/server/responses/run-turn-execution.ts @@ -50,7 +50,7 @@ import { undeclaredToolCallMessage } from "../responses-undeclared-tool-guard"; import { planWebSearch } from "../../web-search"; import { runTurnWebSearchInitialParsed, runTurnWebSearchLoop } from "../../web-search/run-turn-loop"; import { WEB_SEARCH_TOOL_NAME } from "../../web-search/synthetic-tool"; -import { createDevinMessagesOutputOrder } from "../../claude/devin-output-order"; +import { orderDevinMessagesOutput } from "../../claude/devin-output-order"; // LOCAL PATCH (runturn-websearch): top-level fields route binding or the // adapter itself may write during a turn. Iteration-local `turnParsed` objects @@ -182,6 +182,13 @@ export async function executeResponsesRunTurn( const runTurnAbort = new AbortController(); const cleanupRunTurnAbort = linkAbortSignal(runTurnAbort, options.abortSignal); + const devinProducers = new Set(); + const orderMessagesOutput = (source: AsyncIterable): AsyncIterable => + inboundWire === "anthropic" && transportState.runTurnAdapter.name === "devin" && !routedCompaction + ? orderDevinMessagesOutput(source, translatorBudget, runTurnAbort.signal, () => { + for (const producer of devinProducers) producer.abort(); + }) + : source; const queue = createAdapterEventQueue({ onBacklogExceeded: () => runTurnAbort.abort(), }); @@ -230,7 +237,8 @@ export async function executeResponsesRunTurn( options.onCompactionRecoveryAdapterEvent?.(event); targetQueue.push(event); }; - let messagesOutput: ReturnType | undefined; + let producerAbort = runTurnAbort; + let cleanupProducerAbort: (() => void) | undefined; try { if (!pacingSlotAcquired) { pacingSlot = await waitForProviderRequestSlot(route.providerName, route.provider, route.modelId, runTurnAbort.signal); @@ -264,13 +272,15 @@ export async function executeResponsesRunTurn( }, ); if (inboundWire === "anthropic" && transportState.runTurnAdapter.name === "devin" && !routedCompaction) { - messagesOutput = createDevinMessagesOutputOrder(emit, translatorBudget, runTurnAbort.signal, () => runTurnAbort.abort()); + producerAbort = new AbortController(); + cleanupProducerAbort = linkAbortSignal(producerAbort, runTurnAbort.signal); + devinProducers.add(producerAbort); } await transportState.runTurnAdapter.runTurn?.( turnParsed, { headers: requestState.selectedForwardHeaders, - abortSignal: runTurnAbort.signal, + abortSignal: producerAbort.signal, comboAttempt: options.comboAttempt === true, translatorBudget, providerFetch: runTurnProviderFetch, @@ -287,7 +297,7 @@ export async function executeResponsesRunTurn( ), onRecoveryWithheld: noteAdapterRecoveryWithheld, }, - messagesOutput?.emit ?? emit, + emit, ); // LOCAL PATCH (runturn-websearch): adapters may write conversation/ // continuation state onto the object they received; merge it back so @@ -301,7 +311,7 @@ export async function executeResponsesRunTurn( Object.assign(parsed, routeState); } } catch (err) { - (messagesOutput?.emit ?? emit)(err instanceof RequestPacingQueueOverloadError + emit(err instanceof RequestPacingQueueOverloadError ? { type: "error", status: 429, @@ -326,8 +336,8 @@ export async function executeResponsesRunTurn( message: err instanceof Error ? err.message : String(err), }); } finally { - messagesOutput?.flush(); - messagesOutput?.dispose(); + devinProducers.delete(producerAbort); + cleanupProducerAbort?.(); releaseProviderRequestSlot(pacingSlot); // Cursor assigns a stable conversation id inside runTurn on the first headerless // turn; backfill so Logs can filter/total that opening request (#330 / #522). @@ -352,14 +362,14 @@ export async function executeResponsesRunTurn( }); void runTurnAttempt(iterQueue, undefined, false, iterParsed); const stream = iterQueue.stream(); - if (!runTurnFailoverArmed()) return stream; + if (!runTurnFailoverArmed()) return orderMessagesOutput(stream); // LOCAL PATCH (runturn-websearch): a post-search iteration can open on a // 429 too — the search cells already reached the client, so only this // answer call rotates. Preflight replays it on the next account with the // grown history intact; a mid-stream error still ends the turn as before. - return (async function* () { + return orderMessagesOutput((async function* () { yield* await preflightRunTurnFailover(stream, iterParsed); - })(); + })()); }; // Rebind the turn to an admitted account. The failed attempt emitted no client-visible bytes, // so replay is safe, but a Cursor conversation/checkpoint is credential-scoped: carrying its @@ -621,7 +631,7 @@ export async function executeResponsesRunTurn( onBacklogExceeded: () => runTurnAbort.abort(), }); void runTurnAttempt(retryQueue, "empty-completion"); - return retryQueue.stream(); + return orderMessagesOutput(retryQueue.stream()); }; const { toolNsMap, declaredToolNames, toolParameterSchemas, freeformToolNames, bareCustomToolNames, toolSearchToolNames } = toolBridgeMaps; @@ -706,6 +716,8 @@ export async function executeResponsesRunTurn( } eventSource = preflight.stream; } + // Preflight sees raw output before Messages delays text for the late signature. + eventSource = orderMessagesOutput(eventSource); // LOCAL PATCH (runturn-websearch): intercept web_search calls across // iterations; terminal output keeps flowing through the same queue/bridge. if (wsPlan) { diff --git a/structure/clients/claude-desktop.md b/structure/clients/claude-desktop.md index 16c335050ff..de6a833d09c 100644 --- a/structure/clients/claude-desktop.md +++ b/structure/clients/claude-desktop.md @@ -31,7 +31,7 @@ belongs to the shared account store. The import CLI adds pool rows independently ## Devin Messages output ordering -`src/claude/devin-output-order.ts` orders each physical Devin turn at the emitter in +`src/claude/devin-output-order.ts` orders each physical Devin turn after raw-event preflight in `src/server/responses/run-turn-execution.ts` when the original inbound wire is Anthropic Messages. Provider names may be customized; the selected adapter determines applicability. Text and tool events wait for that turn's terminal so Cognition's late reasoning signature @@ -40,9 +40,12 @@ signature-only thinking block. Reasoning and transport progress remain live; ans tool dispatch incur turn-completion latency. Responses and Chat retain their original ordering, and routed compaction is excluded. -Retained events share the request translator budget and are released on terminal, cancellation, -overflow, or adapter EOF. Overflow emits one typed `translation_buffer_limit` error and aborts -the producer. Error and incomplete terminals retain partial output and their original usage. +Retained events share the request translator budget and drain on demand without a synchronous +burst into the adapter queue. They are released on terminal, cancellation, overflow, or adapter +EOF. Overflow emits one typed `translation_buffer_limit` error and aborts only active Devin +producers, preserving error classification through hosted search. Cancellation drops held +semantic output but forwards the adapter's terminal and usage when the consumer keeps reading. +Adapter error and incomplete terminals retain partial output and their original usage. Ordering occurs before hosted-search interception, independently for each physical iteration, so one iteration's signature cannot be attached to another iteration's answer. diff --git a/structure/transports/byte-accounting.md b/structure/transports/byte-accounting.md index bab76ec8acd..d8fcc55a22f 100644 --- a/structure/transports/byte-accounting.md +++ b/structure/transports/byte-accounting.md @@ -57,8 +57,8 @@ lifecycle, cancellation races, protocol envelopes, and the real HTTP admission b ## Stream-buffer accounting Devin's [Messages ordering buffer](../clients/claude-desktop.md#devin-messages-output-ordering) charges retained -semantic events to the shared translator budget until transfer to the independently bounded -adapter queue. Cancellation and overflow release held events before producer shutdown. +semantic events consumed from the independently bounded adapter queue to the shared translator +budget until downstream delivery. Cancellation and overflow release held events before producer shutdown. `src/web-search/run-turn-loop.ts` charges retained iteration events and generated replay history to the request translator budget. Each owner releases its own reservations on completion, error, diff --git a/tests/claude-integration/claude-devin-output-order.test.ts b/tests/claude-integration/claude-devin-output-order.test.ts index 57591c93bd2..c7709cdec36 100644 --- a/tests/claude-integration/claude-devin-output-order.test.ts +++ b/tests/claude-integration/claude-devin-output-order.test.ts @@ -1,5 +1,5 @@ import { afterAll, afterEach, beforeEach, describe, expect, mock, test } from "bun:test"; -import { createDevinMessagesOutputOrder } from "../../src/claude/devin-output-order"; +import { orderDevinMessagesOutput } from "../../src/claude/devin-output-order"; import { collectAnthropicMessage, responsesSseToAnthropicSse } from "../../src/claude/outbound"; import { bridgeToResponsesSSE } from "../../src/bridge"; import { createTestTranslatorBudget } from "../helpers/translator-budget"; @@ -9,6 +9,10 @@ import { messagesToResponsesTranslation } from "../../src/protocols/codecs/messa import { parseRequest } from "../../src/responses/parser"; import type { AdapterEvent, OcxConfig, OcxProviderConfig } from "../../src/types"; import type { ProviderAdapter } from "../../src/adapters/base"; +import { createAdapterEventQueue } from "../../src/adapters/run-turn-queue"; +import { runTurnWebSearchLoop } from "../../src/web-search/run-turn-loop"; +import { DEVIN_CLI_CREDENTIALS_ENV } from "../../src/oauth/devin/cli-import"; +import { saveCredential } from "../../src/oauth/store"; import { createTempHome } from "../helpers/temp-home"; import { acquireOwnedSpendHome } from "../helpers/owned-spend-home"; import { installIsolatedCodexHome } from "../helpers/isolated-codex-home"; @@ -17,6 +21,8 @@ const signature = encodeDevinSignature("sealed.v1.synthetic-attestation", "seale const usage = { inputTokens: 12, outputTokens: 3, totalTokens: 15 }; const terminal: AdapterEvent = { type: "done", usage }; let upstreamEvents: AdapterEvent[] = []; +let afterFirstEvent: (() => Promise) | undefined; +let paceFragments = false; const resolver = await import("../../src/server/adapter-resolve"); const originalResolver = { ...resolver }; mock.module("../../src/server/adapter-resolve", () => ({ ...originalResolver, @@ -26,7 +32,13 @@ mock.module("../../src/server/adapter-resolve", () => ({ ...originalResolver, name: "devin", buildRequest: () => ({ url: provider.baseUrl, method: "POST", headers: {}, body: "" }), async *parseStream() { yield terminal; }, - async runTurn(_parsed, _incoming, emit) { upstreamEvents.forEach(emit); }, + async runTurn(_parsed, _incoming, emit) { + for (let index = 0; index < upstreamEvents.length; index++) { + emit(upstreamEvents[index]!); + if (index === 0) await afterFirstEvent?.(); + if (paceFragments && index % 16 === 0) await new Promise(resolve => setImmediate(resolve)); + } + }, } satisfies ProviderAdapter; }, })); @@ -37,25 +49,34 @@ afterAll(() => { mock.module("../../src/server/adapter-resolve", () => originalR let home: ReturnType; let codexHome: ReturnType; let releaseSpend: () => void; +let previousCliCredentialsPath: string | undefined; beforeEach(() => { home = createTempHome("ocx-claude-devin-output-"); codexHome = installIsolatedCodexHome("ocx-claude-devin-codex-"); releaseSpend = acquireOwnedSpendHome(); + previousCliCredentialsPath = process.env[DEVIN_CLI_CREDENTIALS_ENV]; + process.env[DEVIN_CLI_CREDENTIALS_ENV] = home.path("absent-devin-cli.toml"); upstreamEvents = []; + afterFirstEvent = undefined; + paceFragments = false; +}); +afterEach(() => { + if (previousCliCredentialsPath === undefined) delete process.env[DEVIN_CLI_CREDENTIALS_ENV]; + else process.env[DEVIN_CLI_CREDENTIALS_ENV] = previousCliCredentialsPath; + releaseSpend(); codexHome.restore(); home.remove(); }); -afterEach(() => { releaseSpend(); codexHome.restore(); home.remove(); }); -function ordered(events: AdapterEvent[]) { +async function ordered(events: AdapterEvent[]) { const budget = createTestTranslatorBudget(); const abort = new AbortController(); const output: AdapterEvent[] = []; - const order = createDevinMessagesOutputOrder(event => output.push(event), budget, abort.signal, () => abort.abort()); - try { events.forEach(order.emit); order.flush(); } finally { order.dispose(); } + for await (const event of orderDevinMessagesOutput((async function* () { yield* events; })(), + budget, abort.signal, () => abort.abort())) output.push(event); return { output, budget }; } async function message(events: AdapterEvent[]) { - const { output, budget } = ordered(events); + const { output, budget } = await ordered(events); const source = (async function* () { yield* output; })(); const bridged = bridgeToResponsesSSE(source, "swe-2", undefined, undefined, undefined, undefined, 0, { translatorBudget: budget }); return collectAnthropicMessage(responsesSseToAnthropicSse(bridged, "devin/swe-2", { @@ -121,44 +142,50 @@ describe("Devin Messages late signatures", () => { for (const ending of [terminal, { type: "incomplete", reason: "max_tokens", usage }, { type: "error", status: 502, message: "synthetic reset", usage }] as AdapterEvent[]) { - test(`${ending.type} preserves partial content and the original terminal`, () => { - const { output, budget } = ordered([{ type: "text_delta", text: "partial" }, ending]); + test(`${ending.type} preserves partial content and the original terminal`, async () => { + const { output, budget } = await ordered([{ type: "text_delta", text: "partial" }, ending]); expect(output.filter(event => event.type !== "heartbeat")).toEqual([{ type: "text_delta", text: "partial" }, ending]); expect(budget.snapshot().currentBytes).toBe(0); }); } - test("holding output keeps progress live and cancellation releases it immediately", () => { + test("holding output keeps progress live and cancellation preserves the adapter terminal and usage", async () => { const budget = createTestTranslatorBudget(); const abort = new AbortController(); - const output: AdapterEvent[] = []; - const order = createDevinMessagesOutputOrder(event => output.push(event), budget, abort.signal, () => abort.abort()); - order.emit({ type: "text_delta", text: "still generating" }); - expect(output).toEqual([{ type: "heartbeat" }]); + const queue = createAdapterEventQueue(); + const iterator = orderDevinMessagesOutput(queue.stream(), budget, abort.signal, () => abort.abort()); + queue.push({ type: "text_delta", text: "still generating" }); + expect((await iterator.next()).value).toEqual({ type: "heartbeat" }); expect(budget.snapshot().currentBytes).toBeGreaterThan(0); - order.emit({ type: "heartbeat", replayUnsafe: true }); - expect(output.at(-1)).toEqual({ type: "heartbeat", replayUnsafe: true }); + queue.push({ type: "heartbeat", replayUnsafe: true }); + expect((await iterator.next()).value).toEqual({ type: "heartbeat", replayUnsafe: true }); abort.abort(); expect(budget.snapshot().currentBytes).toBe(0); - order.emit(terminal); - order.dispose(); - expect(output.some(event => event.type === "text_delta" || event.type === "done")).toBe(false); + const cancelled: AdapterEvent = { type: "error", status: 499, message: "client closed request", usage }; + queue.push(cancelled); + queue.close(); + expect((await iterator.next()).value).toEqual(cancelled); + expect((await iterator.next()).done).toBe(true); }); - test("overflow emits one typed error, aborts the producer, and releases held output", () => { + test("overflow emits one typed error and stops only the producer, preserving search error classification", async () => { const budget = createTestTranslatorBudget({ maxTurnBytes: 160 }); - const abort = new AbortController(); + const requestAbort = new AbortController(); + const producerAbort = new AbortController(); const output: AdapterEvent[] = []; - const order = createDevinMessagesOutputOrder(event => output.push(event), budget, abort.signal, () => abort.abort()); - order.emit({ type: "text_delta", text: "small" }); - order.emit({ type: "text_delta", text: "x".repeat(200) }); - order.emit(terminal); + const source = (async function* () { + yield { type: "text_delta", text: "small" } as AdapterEvent; + yield { type: "text_delta", text: "x".repeat(200) } as AdapterEvent; + yield terminal; + })(); + for await (const event of orderDevinMessagesOutput(source, budget, requestAbort.signal, () => producerAbort.abort())) output.push(event); expect(output.filter(event => event.type === "error")).toHaveLength(1); expect(output.at(-1)).toMatchObject({ type: "error", status: 413, code: "translation_buffer_limit" }); - expect(abort.signal.aborted).toBe(true); + expect(producerAbort.signal.aborted).toBe(true); + expect(requestAbort.signal.aborted).toBe(false); expect(budget.snapshot().currentBytes).toBe(0); - order.dispose(); }); + }); const config = (): OcxConfig => ({ port: 0, defaultProvider: "cognition-custom", claudeCode: { enabled: true }, @@ -190,3 +217,100 @@ test("the Responses ingress retains incremental text before the late signature", const text = await response.text(); expect(text.indexOf("response.output_text.delta")).toBeLessThan(text.indexOf("\"type\":\"reasoning\"")); }); + + +test("a tool call with 3,000 streamed argument fragments drains without overflowing the real queue", async () => { + const argumentsText = JSON.stringify({ file_path: "/synthetic/file.txt", content: "x".repeat(3000) }); + paceFragments = true; + upstreamEvents = [{ type: "tool_call_start", id: "call_write", name: "Write" }, + ...[...argumentsText].map(fragment => ({ type: "tool_call_delta", arguments: fragment }) as AdapterEvent), + { type: "tool_call_end" }, { type: "thinking_signature", signature }, terminal]; + const response = await handleClaudeMessages(new Request("http://localhost/v1/messages", { + method: "POST", headers: { "content-type": "application/json" }, + body: JSON.stringify({ model: "cognition-custom/swe-2", max_tokens: 4096, stream: false, + messages: [{ role: "user", content: "Write" }], + tools: [{ name: "Write", input_schema: { type: "object", properties: { content: { type: "string" } } } }], + }), + }), config(), {}); + const result = await response.json(); + expect(result.stop_reason).toBe("tool_use"); + expect(result.content.at(-1)).toMatchObject({ type: "tool_use", input: JSON.parse(argumentsText) }); + expect(result.usage).toMatchObject({ input_tokens: 12, output_tokens: 3 }); +}); + +test("OAuth preflight releases Messages headers at the first text, before a late signature", async () => { + await saveCredential("devin", { access: "devin-session-token$synthetic", refresh: "synthetic", + expires: Number.MAX_SAFE_INTEGER, accountId: "synthetic-account", source: "oauth", + apiBaseUrl: "https://server.codeium.com" }); + let finish!: () => void; + afterFirstEvent = () => new Promise(resolve => { finish = resolve; }); + upstreamEvents = [{ type: "text_delta", text: "OK" }, { type: "thinking_signature", signature }, terminal]; + const oauthConfig: OcxConfig = { port: 0, defaultProvider: "devin", claudeCode: { enabled: true }, + providers: { devin: { adapter: "devin", authMode: "oauth", baseUrl: "https://server.codeium.com", models: ["swe-2"] } } }; + const responsePromise = handleClaudeMessages(new Request("http://localhost/v1/messages", { + method: "POST", headers: { "content-type": "application/json" }, + body: JSON.stringify({ model: "devin/swe-2", max_tokens: 64, stream: true, + messages: [{ role: "user", content: "Reply OK" }] }), + }), oauthConfig, {}); + let headerTimer: ReturnType | undefined; + try { + const response = await Promise.race([responsePromise, + new Promise((_, reject) => { headerTimer = setTimeout(() => reject(new Error("headers waited for terminal")), 1000); })]); + expect(response.status).toBe(200); + expect(finish).toBeDefined(); + finish(); + const result = await collectAnthropicMessage(response.body!, "devin/swe-2", createTestTranslatorBudget()); + expect(result.content.at(-1)).toEqual({ type: "text", text: "OK" }); + } finally { clearTimeout(headerTimer); finish?.(); } +}); + + +test("hosted search forwards the typed overflow instead of inventing a client cancellation", async () => { + const budget = createTestTranslatorBudget({ maxTurnBytes: 160 }); + const requestAbort = new AbortController(); + const producerAbort = new AbortController(); + const source = orderDevinMessagesOutput((async function* () { + yield { type: "text_delta", text: "x".repeat(200) } as AdapterEvent; + })(), budget, requestAbort.signal, () => producerAbort.abort()); + const loop = runTurnWebSearchLoop(source, { + parsed: { modelId: "swe-2", stream: true, options: {}, context: { messages: [], tools: [] } }, + plan: { backend: "exa", hostedTool: { type: "web_search" }, maxSearches: 1, + settings: { model: "fixture", reasoning: "low", timeoutMs: 100 }, + routedModelStallTimeoutMs: 100, stallTimeoutSec: 1, streamRoutedModelOutput: true }, + abortSignal: requestAbort.signal, translatorBudget: budget, + dispatch: async function* () { throw new Error("no search dispatch expected"); }, + }); + const output: AdapterEvent[] = []; + for await (const event of loop) output.push(event); + expect(output).toEqual([expect.objectContaining({ type: "error", status: 413, code: "translation_buffer_limit" })]); + expect(producerAbort.signal.aborted).toBe(true); + expect(budget.snapshot().currentBytes).toBe(0); +}); + +test("cancelling during the terminal drain preserves usage and clears retained semantic events", async () => { + const budget = createTestTranslatorBudget(); + const abort = new AbortController(); + const iterator = orderDevinMessagesOutput((async function* () { + yield { type: "text_delta", text: "partial" } as AdapterEvent; + yield terminal; + })(), budget, abort.signal, () => abort.abort()); + expect((await iterator.next()).value).toEqual({ type: "heartbeat" }); + expect((await iterator.next()).value).toEqual({ type: "text_delta", text: "partial" }); + abort.abort(); + expect(budget.snapshot().currentBytes).toBe(0); + expect((await iterator.next()).value).toEqual(terminal); + expect((await iterator.next()).done).toBe(true); +}); + +test("a consumer return releases the ordering buffer without a terminal", async () => { + const budget = createTestTranslatorBudget(); + const abort = new AbortController(); + const iterator = orderDevinMessagesOutput((async function* () { + yield { type: "text_delta", text: "held" } as AdapterEvent; + yield terminal; + })(), budget, abort.signal, () => abort.abort()); + await iterator.next(); + expect(budget.snapshot().currentBytes).toBeGreaterThan(0); + await iterator.return(); + expect(budget.snapshot().currentBytes).toBe(0); +}); From 794d61af74a0a670e0981975da9c72052257233f Mon Sep 17 00:00:00 2001 From: foxytanuki Date: Wed, 30 Sep 2026 06:59:31 +0000 Subject: [PATCH 3/4] fix(claude): preserve cancellation during devin output drain --- src/claude/devin-output-order.ts | 9 +++- src/server/responses/run-turn-execution.ts | 1 + structure/clients/claude-desktop.md | 4 +- .../claude-devin-output-order.test.ts | 47 ++++++++++++++++++- 4 files changed, 57 insertions(+), 4 deletions(-) diff --git a/src/claude/devin-output-order.ts b/src/claude/devin-output-order.ts index 555f9b87111..69be3f6a9fc 100644 --- a/src/claude/devin-output-order.ts +++ b/src/claude/devin-output-order.ts @@ -28,6 +28,11 @@ export async function* orderDevinMessagesOutput( }; const cancel = () => { cancelled = true; release(); }; signal.addEventListener("abort", cancel, { once: true }); + const cancelledTerminal = (event: AdapterEvent): AdapterEvent => { + if (event.type !== "done" && event.type !== "incomplete") return event; + return { type: "error", status: 499, message: "client closed request", retryable: false, + ...(event.usage ? { usage: event.usage } : {}) }; + }; async function* drain() { for (let index = 0; index < held.length && !cancelled; index++) { delivering = held[index]; @@ -51,7 +56,7 @@ export async function* orderDevinMessagesOutput( if (cancelled) { // The adapter's cancellation terminal retains its measured usage. Client // stream cancellation instead returns this iterator and stops consumption. - if (terminal) { yield event; return; } + if (terminal) { yield cancelledTerminal(event); return; } continue; } if (event.type === "heartbeat" || event.type === "thinking_delta" @@ -76,7 +81,7 @@ export async function* orderDevinMessagesOutput( } if (terminal) { yield* drain(); - if (cancelled && !terminalDelivered) yield event; + if (cancelled && !terminalDelivered) yield cancelledTerminal(event); return; } // Raw preflight already observed real output; feed the stream watchdog while diff --git a/src/server/responses/run-turn-execution.ts b/src/server/responses/run-turn-execution.ts index d74aeaa5fb4..b038bffb191 100644 --- a/src/server/responses/run-turn-execution.ts +++ b/src/server/responses/run-turn-execution.ts @@ -668,6 +668,7 @@ export async function executeResponsesRunTurn( retryAfter: resolveClientRetryAfter({ status: httpStatus, message: error.message }), }); }; + // Messages ingress forces internal streaming, including buffered client requests. if (parsed.stream) { try { void runTurn(); diff --git a/structure/clients/claude-desktop.md b/structure/clients/claude-desktop.md index de6a833d09c..fcaa1c50daa 100644 --- a/structure/clients/claude-desktop.md +++ b/structure/clients/claude-desktop.md @@ -44,7 +44,9 @@ Retained events share the request translator budget and drain on demand without burst into the adapter queue. They are released on terminal, cancellation, overflow, or adapter EOF. Overflow emits one typed `translation_buffer_limit` error and aborts only active Devin producers, preserving error classification through hosted search. Cancellation drops held -semantic output but forwards the adapter's terminal and usage when the consumer keeps reading. +semantic output; this consumer preserves error terminals and maps cancelled success/incomplete +terminals to a 499 with their original usage, so partial output cannot commit completed replay +state. Hosted search retains its existing independent cancellation mapping. Adapter error and incomplete terminals retain partial output and their original usage. Ordering occurs before hosted-search interception, independently for each physical iteration, so one iteration's signature cannot be attached to another iteration's answer. diff --git a/tests/claude-integration/claude-devin-output-order.test.ts b/tests/claude-integration/claude-devin-output-order.test.ts index c7709cdec36..96f23b1066c 100644 --- a/tests/claude-integration/claude-devin-output-order.test.ts +++ b/tests/claude-integration/claude-devin-output-order.test.ts @@ -298,7 +298,7 @@ test("cancelling during the terminal drain preserves usage and clears retained s expect((await iterator.next()).value).toEqual({ type: "text_delta", text: "partial" }); abort.abort(); expect(budget.snapshot().currentBytes).toBe(0); - expect((await iterator.next()).value).toEqual(terminal); + expect((await iterator.next()).value).toMatchObject({ type: "error", status: 499, usage }); expect((await iterator.next()).done).toBe(true); }); @@ -314,3 +314,48 @@ test("a consumer return releases the ordering buffer without a terminal", async await iterator.return(); expect(budget.snapshot().currentBytes).toBe(0); }); + + +for (const queuedTerminal of [terminal, { type: "incomplete", reason: "max_tokens", usage }] as AdapterEvent[]) { + test(`cancel before consuming a queued ${queuedTerminal.type} keeps usage but forbids a successful terminal`, async () => { + const budget = createTestTranslatorBudget(); + const abort = new AbortController(); + const queue = createAdapterEventQueue(); + const iterator = orderDevinMessagesOutput(queue.stream(), budget, abort.signal, () => abort.abort()); + queue.push({ type: "text_delta", text: "held" }); + expect((await iterator.next()).value).toEqual({ type: "heartbeat" }); + queue.push(queuedTerminal); + queue.close(); + abort.abort(); + expect((await iterator.next()).value).toMatchObject({ type: "error", status: 499, retryable: false, usage }); + expect((await iterator.next()).done).toBe(true); + expect(budget.snapshot().currentBytes).toBe(0); + }); +} + +test("cancelled partial drain never commits a completed response or replay state", async () => { + const budget = createTestTranslatorBudget(); + const abort = new AbortController(); + const orderedEvents = orderDevinMessagesOutput((async function* () { + yield { type: "text_delta", text: "partial" } as AdapterEvent; + yield { type: "tool_call_start", id: "call_partial", name: "Write" } as AdapterEvent; + yield { type: "tool_call_delta", arguments: '{"content":' } as AdapterEvent; + yield terminal; + })(), budget, abort.signal, () => abort.abort()); + const cancelDuringDrain = (async function* () { + for await (const event of orderedEvents) { + yield event; + if (event.type === "text_delta") abort.abort(); + } + })(); + let completions = 0; + const observedUsage: unknown[] = []; + const bridge = bridgeToResponsesSSE(cancelDuringDrain, "swe-2", undefined, undefined, undefined, undefined, 0, { + translatorBudget: budget, onCompletedResponse: () => { completions++; }, onUsage: value => observedUsage.push(value), + }); + const bytes = await new Response(bridge).text(); + expect(bytes).not.toContain("response.completed"); + expect(completions).toBe(0); + expect(observedUsage).toContainEqual(usage); + expect(bytes).toContain("client closed request"); +}); From b67d38ce3f2b5298ec0651c46b6ca8c280d9c907 Mon Sep 17 00:00:00 2001 From: foxytanuki Date: Thu, 1 Oct 2026 01:24:31 +0900 Subject: [PATCH 4/4] docs(claude): scope Devin buffering to Messages API --- docs-site/src/content/docs/guides/claude-code.md | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/docs-site/src/content/docs/guides/claude-code.md b/docs-site/src/content/docs/guides/claude-code.md index 6bee6138cb5..8ba8e33f281 100644 --- a/docs-site/src/content/docs/guides/claude-code.md +++ b/docs-site/src/content/docs/guides/claude-code.md @@ -7,8 +7,8 @@ opencodex serves `POST /v1/messages` (plus `count_tokens`) alongside `/v1/respon Code can use every routed provider — OAuth logins, account pools, key failover and sidecars included — with zero extra auth work. -On Devin routes (including SWE-2), text and tool calls wait for the upstream turn to complete -so its late reasoning signature can precede the answer. This prevents Claude Code's final +On Devin routes (including SWE-2) reached through the Messages API, text and tool calls wait +for the upstream turn to complete so its late reasoning signature can precede the answer. This prevents Claude Code's final result from becoming empty; reasoning and keepalive progress still flow during generation. The buffer shares the request's 32 MiB translation limit and cancellation stops the producer. This output-order fix does not resolve Cognition's separate refusal of some generated system