diff --git a/docs-site/src/content/docs/guides/claude-code.md b/docs-site/src/content/docs/guides/claude-code.md index 1d3394163f..8ba8e33f28 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) 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 +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 2c11df1cbb..b4d27b9425 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", "codex-config-preservation.test.ts": "codex-integration", "codex-credits.test.ts": "codex-integration", "codex-credits-settings.test.ts": "codex-integration", diff --git a/src/claude/devin-output-order.ts b/src/claude/devin-output-order.ts new file mode 100644 index 0000000000..ef4ee7e293 --- /dev/null +++ b/src/claude/devin-output-order.ts @@ -0,0 +1,98 @@ +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. + * Consume after raw-event preflight and drain on demand, avoiding a queue burst. + */ +export async function* orderDevinMessagesOutput( + source: AsyncIterable, + budget: TranslatorBudget, + signal: AbortSignal, + abortProducer: () => void, +): AsyncGenerator { + const held: Array = []; + let cancelled = signal.aborted; + let terminalDelivered = false; + let delivering: AdapterEvent | undefined; + const release = () => { + for (const event of held) if (event) releaseTranslatedEvent(event, budget); + held.length = 0; + if (delivering) releaseTranslatedEvent(delivering, budget); + }; + 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]; + held[index] = undefined; + try { + if (delivering) { + terminalDelivered = delivering.type === "done" || delivering.type === "error" || delivering.type === "incomplete"; + yield delivering; + } + } + 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 cancelledTerminal(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") { + yield event; + continue; + } + // Retention owns a snapshot, including nested usage, so later producer or + // consumer mutations cannot change the bytes measured by the budget. + const copy = structuredClone(event); + try { + retainTranslatedEvent(copy, budget, held.at(-1)); + held.push(copy); + } catch (error) { + if (!isTranslatorBudgetExceededError(error)) throw error; + release(); + // 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; + } + if (terminal) { + yield* drain(); + if (cancelled && !terminalDelivered) yield cancelledTerminal(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 11cfa90d7f..b038bffb19 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 { 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 @@ -181,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(), }); @@ -229,6 +237,8 @@ export async function executeResponsesRunTurn( options.onCompactionRecoveryAdapterEvent?.(event); targetQueue.push(event); }; + let producerAbort = runTurnAbort; + let cleanupProducerAbort: (() => void) | undefined; try { if (!pacingSlotAcquired) { pacingSlot = await waitForProviderRequestSlot(route.providerName, route.provider, route.modelId, runTurnAbort.signal); @@ -261,11 +271,16 @@ export async function executeResponsesRunTurn( turnScopedPacing: true, }, ); + if (inboundWire === "anthropic" && transportState.runTurnAdapter.name === "devin" && !routedCompaction) { + 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, @@ -321,6 +336,8 @@ export async function executeResponsesRunTurn( message: err instanceof Error ? err.message : String(err), }); } finally { + 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). @@ -345,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 @@ -614,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; @@ -651,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(); @@ -699,6 +717,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 97d24db480..4d433a3ab4 100644 --- a/structure/clients/claude-desktop.md +++ b/structure/clients/claude-desktop.md @@ -29,6 +29,29 @@ 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 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 +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 own deep snapshots, including nested usage, so producer/consumer mutations +cannot alter their measured payload. They 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; 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. + ## 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 991a3723fa..f6586964ea 100644 --- a/structure/providers-and-adapters.md +++ b/structure/providers-and-adapters.md @@ -1,7 +1,7 @@ # Providers And Adapters Anthropic account pause, model routes, and quota labels follow the -[Anthropic account-pool contract](providers/anthropic-account-pool.md). +[Anthropic account-pool contract](providers/anthropic-account-pool.md). 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. Per-account usage thresholds follow the [Anthropic account thresholds contract](providers/anthropic-account-thresholds.md). An Anthropic 429 records the served account's cooldown even when the request has used its allowed retry sends. That final account remains excluded on the next request; combo target cooling is skipped only after the matching account cooldown is present. diff --git a/structure/transports/byte-accounting.md b/structure/transports/byte-accounting.md index 72f3945eb1..314fbee80b 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 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, 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 0000000000..982a840fe0 --- /dev/null +++ b/tests/claude-integration/claude-devin-output-order.test.ts @@ -0,0 +1,422 @@ +import { afterAll, afterEach, beforeEach, describe, expect, mock, test } from "bun:test"; +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"; +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 { 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"; + +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[] = []; +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, + 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) { + 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; + }, +})); +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; +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(); +}); + +async function ordered(events: AdapterEvent[]) { + const budget = createTestTranslatorBudget(); + const abort = new AbortController(); + const output: AdapterEvent[] = []; + 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 } = 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", { + 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`, 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 preserves the adapter terminal and usage", 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: "still generating" }); + expect((await iterator.next()).value).toEqual({ type: "heartbeat" }); + expect(budget.snapshot().currentBytes).toBeGreaterThan(0); + queue.push({ type: "heartbeat", replayUnsafe: true }); + expect((await iterator.next()).value).toEqual({ type: "heartbeat", replayUnsafe: true }); + abort.abort(); + expect(budget.snapshot().currentBytes).toBe(0); + 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 and stops only the producer, preserving search error classification", async () => { + const budget = createTestTranslatorBudget({ maxTurnBytes: 160 }); + const requestAbort = new AbortController(); + const producerAbort = new AbortController(); + const output: AdapterEvent[] = []; + 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(producerAbort.signal.aborted).toBe(true); + expect(requestAbort.signal.aborted).toBe(false); + expect(budget.snapshot().currentBytes).toBe(0); + }); + +}); + +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"] } }, +}); + +test("retained terminal usage is isolated from producer and consumer mutations", async () => { + const budget = createTestTranslatorBudget(); + const abort = new AbortController(); + const originalUsage = { ...usage, rawUsage: { detail: { tokens: 3 } } }; + const ending: AdapterEvent = { type: "done", usage: originalUsage }; + const iterator = orderDevinMessagesOutput((async function* () { + yield { type: "text_delta", text: "OK" } as AdapterEvent; + yield ending; + })(), budget, abort.signal, () => abort.abort()); + await iterator.next(); // Progress while the answer is retained. + expect((await iterator.next()).value).toEqual({ type: "text_delta", text: "OK" }); + const measuredBytes = budget.snapshot().currentBytes; + originalUsage.outputTokens = 999; + originalUsage.rawUsage.detail.tokens = 999; + const result = (await iterator.next()).value; + expect(result).toEqual({ type: "done", usage: { ...usage, rawUsage: { detail: { tokens: 3 } } } }); + if (!result || result.type !== "done" || !result.usage) throw new Error("missing terminal usage"); + result.usage.outputTokens = 1; + expect(originalUsage.outputTokens).toBe(999); + expect(budget.snapshot().currentBytes).toBeLessThan(measuredBytes); + expect((await iterator.next()).done).toBe(true); + expect(budget.snapshot().currentBytes).toBe(0); +}); + +test("a stalled Devin turn sends repeated Messages pings before releasing the held answer", async () => { + const budget = createTestTranslatorBudget(); + const abort = new AbortController(); + const queue = createAdapterEventQueue(); + const events = orderDevinMessagesOutput(queue.stream(), budget, abort.signal, () => abort.abort()); + const bridge = bridgeToResponsesSSE(events, "swe-2", undefined, undefined, undefined, undefined, 2000, { translatorBudget: budget }); + const reader = responsesSseToAnthropicSse(bridge, "devin/swe-2", { translatorBudget: budget, pingIntervalMs: 10 }).getReader(); + const chunks: string[] = []; + const decoder = new TextDecoder(); + let timer: ReturnType | undefined; + queue.push({ type: "text_delta", text: "OK" }); + try { + await Promise.race([(async () => { + // Include a periodic ping after the progress ping from held text. + while ((chunks.join("").match(/event: ping/g) ?? []).length < 3) { + const next = await reader.read(); + if (next.done) throw new Error("stream ended before the signature"); + chunks.push(decoder.decode(next.value, { stream: true })); + } + })(), new Promise((_, reject) => { timer = setTimeout(() => reject(new Error("stalled Messages stream sent no pings")), 1000); })]); + expect(chunks.join("")).not.toContain("text_delta"); + expect(budget.snapshot().currentBytes).toBeGreaterThan(0); + queue.push({ type: "thinking_signature", signature }); + queue.push(terminal); + queue.close(); + for (;;) { + const next = await reader.read(); + if (next.done) break; + chunks.push(decoder.decode(next.value, { stream: true })); + } + const wire = chunks.join(""); + expect(wire).toContain('"text":"OK"'); + expect(wire.indexOf("signature_delta")).toBeLessThan(wire.indexOf("text_delta")); + expect(wire).toContain("event: message_stop"); + } finally { clearTimeout(timer); queue.close(); await reader.cancel(); budget.dispose(); } +}); + +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\"")); +}); + + +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).toMatchObject({ type: "error", status: 499, usage }); + 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); +}); + + +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"); +}); diff --git a/tests/fixtures/test-layout-expected.json b/tests/fixtures/test-layout-expected.json index 3bf189236e..83170e64b6 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", "codex-config-preservation.test.ts": "codex-integration", "codex-credits.test.ts": "codex-integration", "codex-credits-settings.test.ts": "codex-integration",