diff --git a/scripts/test-layout/layout.json b/scripts/test-layout/layout.json index f3f16fae58b..43ef8b9e146 100644 --- a/scripts/test-layout/layout.json +++ b/scripts/test-layout/layout.json @@ -920,6 +920,8 @@ "devin-provider-merge-migration.test.ts": "providers", "devin-stated-reset-hardening.test.ts": "providers", "devin-stated-reset-retry.test.ts": "providers", + "devin-anthropic-signature-fallback.test.ts": "providers", + "devin-reasoning-continuation.test.ts": "providers", "devin-stream-deadline.test.ts": "providers", "digitalocean-scaleway-provider.test.ts": "providers", "docs-429-failover-claims.test.ts": "ci-workflows", diff --git a/src/adapters/devin.ts b/src/adapters/devin.ts index 6c7c8e3ff6d..7b148d09334 100644 --- a/src/adapters/devin.ts +++ b/src/adapters/devin.ts @@ -10,12 +10,12 @@ import type { AdapterEvent, OcxAssistantMessage, OcxContentPart, OcxMessage, Ocx import { namespacedToolName } from "../types"; import type { IncomingMeta, ProviderAdapter } from "./base"; import { streamChatEventsWithResetRetry, devinStatedResetWaitMs, allocateCascadeId, CloudChatError, type ChatHistoryItem, type ToolDef } from "./devin/cloud-direct"; -import type { ContentPart } from "./devin/cloud-direct/chat"; +import type { CloudChatEvent, ContentPart } from "./devin/cloud-direct/chat"; import { getCachedCatalog, type CacheEntry } from "./devin/cloud-direct/catalog"; import { collapseDevinModelUid } from "./devin/live-models"; import { buildNonOpenAIToolCatalogNudgeForTools } from "./tool-catalog-nudge"; import { DEVIN_DEFAULT_API_SERVER, resolveDevinApiServer } from "../oauth/devin"; -import { isProviderIssuedThinkingSignature } from "../responses/reasoning-envelope"; +import { devinAssistantReasoning, encodeDevinSignature, hasAnthropicSignature } from "./devin/reasoning-signature"; import { SendBudgetExhaustedError } from "../lib/upstream-retry"; /** @@ -54,6 +54,22 @@ export function mergeDevinUsage(previous: OcxUsage, next: OcxUsage): OcxUsage { */ const DEVIN_CLIENT_CLOSED_MESSAGE = "client closed request"; +/** Below the bridge's upstream stall deadline, so held reasoning never reads as a stall. */ +const HELD_REASONING_HEARTBEAT_MS = 15_000; + +type DevinUsageEvent = Extract; + +/** The retry's cumulative usage plus the refused attempt's final counts. */ +function addDevinUsage(event: DevinUsageEvent, prior: DevinUsageEvent): DevinUsageEvent { + const sum = (a?: number, b?: number) => (a === undefined && b === undefined ? undefined : (a ?? 0) + (b ?? 0)); + const out: DevinUsageEvent = { ...event }; + for (const key of ["promptTokens", "completionTokens", "totalTokens", "cachedInputTokens", "cacheCreationInputTokens", "reasoningTokens"] as const) { + const value = sum(event[key], prior[key]); + if (value !== undefined) out[key] = value; + } + return out; +} + /** Map a cloud-direct failure onto the structured fields the error event carries. */ export function devinErrorClassification(error: unknown): { status?: number; errorType?: string; retryable?: boolean } { const status = error instanceof CloudChatError ? error.status : undefined; @@ -364,45 +380,16 @@ function assistantText(message: OcxAssistantMessage): string { // Thinking stays out of the replayed TEXT: folding chain-of-thought into // assistant text sends it back as visible prior output, which the model // then treats as something it said to the user. It is replayed in its own - // field instead — see assistantThinking below. + // field instead — see devinAssistantReasoning. .map((part) => (part.type === "text" ? part.text : "")) .filter(Boolean) .join("\n"); } -/** - * The assistant turn's own reasoning, for replay in ChatMessagePrompt #11. - * - * This adapter previously asserted that Cognition has no reasoning-replay - * field and dropped the thinking outright, so a reasoning model restarted its - * chain on every turn of a tool loop. The field exists: two independent - * clients of the same service write #11 thinking with #12 signature and #18 - * signature_type on the assistant prompt. - * - * Field #12 attests the exact text at #11, and the wire has room for one pair. - * Every block that carries text is replayed, so the chain stays intact; the - * signature rides along only when the text being replayed IS the text it - * attests, which is exactly the single-block case. Several independently signed - * blocks send an unsigned prompt rather than pairing one block's attestation - * with another block's words. A signature-only block attests encrypted thinking - * that is not being replayed at all, so it is not one of these blocks and - * cannot contribute the pair. - */ -function assistantThinking( - message: OcxAssistantMessage, -): { thinking?: string; signature?: string } { - const blocks = message.content.filter( - (part): part is Extract => part.type === "thinking", - ).filter(part => Boolean(part.thinking)); - if (blocks.length === 0) return {}; - const signature = blocks.length === 1 ? blocks[0]!.signature : undefined; - return { - thinking: blocks.map(part => part.thinking).join("\n"), - ...(isProviderIssuedThinkingSignature(signature) ? { signature } : {}), - }; -} - -export function mapOcxMessagesToDevin(parsed: OcxParsedRequest): ChatHistoryItem[] { +export function mapOcxMessagesToDevin( + parsed: OcxParsedRequest, + options: { withholdAnthropicSignatures?: boolean } = {}, +): ChatHistoryItem[] { const items: ChatHistoryItem[] = []; // Cognition is not an OpenAI host, and this adapter does advertise a real // client tool catalog (proto #10 via `mapOcxToolsToDevin`), so the same @@ -421,13 +408,17 @@ export function mapOcxMessagesToDevin(parsed: OcxParsedRequest): ChatHistoryItem if (system) items.push({ role: "system", content: system }); for (const message of parsed.context.messages) { - const mapped = mapOneMessage(message); + const mapped = mapOneMessage(message, parsed.modelId, options); if (mapped) items.push(mapped); } return items; } -function mapOneMessage(message: OcxMessage): ChatHistoryItem | undefined { +function mapOneMessage( + message: OcxMessage, + modelId: string, + options: { withholdAnthropicSignatures?: boolean }, +): ChatHistoryItem | undefined { if (message.role === "user" || message.role === "developer") { const content = mapOcxContentToWire(message.content); // An image with no caption text is a complete user message on its own. @@ -439,10 +430,10 @@ function mapOneMessage(message: OcxMessage): ChatHistoryItem | undefined { if (message.role === "assistant") { const toolCalls = assistantToolCalls(message); const text = assistantText(message); - const reasoning = assistantThinking(message); + const reasoning = devinAssistantReasoning(message, modelId, options.withholdAnthropicSignatures === true); // A turn that produced only reasoning is still worth replaying: dropping it // is what makes the next turn re-derive the same chain. - if (!text && toolCalls.length === 0 && !reasoning.thinking) return undefined; + if (!text && toolCalls.length === 0 && !reasoning.thinking && !reasoning.signature) return undefined; return { role: "assistant", content: text || "", @@ -648,12 +639,20 @@ export function createDevinAdapter( // An admitted HTTP turn owns globally shared capacity until this call // emits. Without an explicit wait allowance, preserve the typed reset // delay in generated diagnostic wording and return immediately. - for await (const event of streamChatEventsWithResetRetry({ + const signedMessages = mapOcxMessagesToDevin(parsed); + // A Claude signature is replayed because it is what carries the reasoning into this + // turn, but Cognition streams Claude's thinking as a summary the signature does not + // cover, and some replays are refused with invalid_argument before any output. That + // refusal is retried once with the Anthropic signatures withheld and the text kept. + const unsignedMessages = hasAnthropicSignature(signedMessages, parsed.modelId) + ? mapOcxMessagesToDevin(parsed, { withholdAnthropicSignatures: true }) + : undefined; + const request = (messages: ChatHistoryItem[]) => streamChatEventsWithResetRetry({ apiKey, apiServerUrl: host, modelUid, catalog, - messages: mapOcxMessagesToDevin(parsed), + messages, tools: mapOcxToolsToDevin(parsed.context.tools), cascadeId, // Input and output ceilings are separate wire fields. Omitting the @@ -675,7 +674,57 @@ export function createDevinAdapter( onPhysicalSend: incoming.onPhysicalSend, onRecoveryWithheld: incoming.onRecoveryWithheld, }, - })) { + }); + async function* withSignatureFallback() { + if (!unsignedMessages) { + yield* request(signedMessages); + return; + } + // Events from the signed attempt are held until its outcome is known: a refusal + // after reasoning would otherwise leave the client with the refused attempt's + // reasoning and signature, and the next turn would replay that signature against + // the retry's thinking. + const held: CloudChatEvent[] = []; + // Usage is still real: the refused attempt was processed, so its final counts are + // added to every usage frame of the retry (frames are cumulative per request). + let refusedUsage: Extract | undefined; + let visible = false; + let lastHeartbeat = Date.now(); + try { + for await (const event of request(signedMessages)) { + // Only visible output makes a retry unsafe. Live, the refusal often lands after the + // model has streamed its reasoning, its signature and a finish frame, and nothing else. + if (!visible && (event.kind === "text" || event.kind === "tool_call_start" || event.kind === "tool_call_args")) { + visible = true; + yield* held.splice(0); + } + if (visible) { + yield event; + continue; + } + held.push(event); + if (event.kind === "usage") refusedUsage = event; + // Held reasoning must not look like a stalled upstream to the bridge. + if (Date.now() - lastHeartbeat >= HELD_REASONING_HEARTBEAT_MS) { + lastHeartbeat = Date.now(); + emit({ type: "heartbeat" }); + } + } + } catch (error) { + if (visible || !(error instanceof CloudChatError && error.code === "invalid_argument")) { + yield* held.splice(0); + throw error; + } + // Emitted first so the counts survive a retry that reports no usage or fails early. + if (refusedUsage) yield refusedUsage; + for await (const event of request(unsignedMessages)) { + yield event.kind === "usage" && refusedUsage ? addDevinUsage(event, refusedUsage) : event; + } + return; + } + yield* held.splice(0); + } + for await (const event of withSignatureFallback()) { if (incoming.abortSignal?.aborted) { // Emitting nothing here left the bridge to synthesize adapter_eof. // Say what happened instead, the way the other runTurn-only adapter @@ -696,7 +745,7 @@ export function createDevinAdapter( if (event.kind === "reasoning_signature") { // Carried back out so the next turn can replay it in the prompt's // signature field; an unsigned replay is what the service ignores. - emit({ type: "thinking_signature", signature: event.signature }); + emit({ type: "thinking_signature", signature: encodeDevinSignature(event.signature, event.signatureType) }); continue; } if (event.kind === "tool_call_start") { diff --git a/src/adapters/devin/cloud-direct/chat.ts b/src/adapters/devin/cloud-direct/chat.ts index 30a520d3e2a..906ac4deea9 100644 --- a/src/adapters/devin/cloud-direct/chat.ts +++ b/src/adapters/devin/cloud-direct/chat.ts @@ -486,7 +486,7 @@ export type CloudChatEvent = * turn produced. Without decoding it there is nothing to put in the prompt's * #12 on the next turn, so the replay would always be unsigned. */ - | { kind: 'reasoning_signature'; signature: string } + | { kind: 'reasoning_signature'; signature: string; signatureType?: string } | { kind: 'tool_call_start'; id: string; name: string } | { kind: 'tool_call_args'; @@ -792,6 +792,10 @@ export function* decodeChatFrame(proto: Buffer): Generator { } } if (authoritativeUsage) yield authoritativeUsage; + let signatureType: string | undefined; + for (const f of iterFields(proto)) { + if (f.num === 21 && f.wire === 2 && Buffer.isBuffer(f.value)) signatureType = (f.value as Buffer).toString('utf8') || undefined; + } for (const f of iterFields(proto)) { if (f.num === 3 && f.wire === 2 && Buffer.isBuffer(f.value)) { // Visible delta_text — what the user should SEE in the chat. @@ -817,7 +821,8 @@ export function* decodeChatFrame(proto: Buffer): Generator { if (s) yield { kind: 'reasoning', text: s }; } else if (f.num === 10 && f.wire === 2 && Buffer.isBuffer(f.value)) { const s = (f.value as Buffer).toString('utf8'); - if (s) yield { kind: 'reasoning_signature', signature: s }; + // #21 delta_signature_type arrives in the same frame; the prompt replays it as #18. + if (s) yield { kind: 'reasoning_signature', signature: s, ...(signatureType ? { signatureType } : {}) }; } else if (f.num === 6 && f.wire === 2 && Buffer.isBuffer(f.value)) { let id: string | undefined; let name: string | undefined; diff --git a/src/adapters/devin/reasoning-signature.ts b/src/adapters/devin/reasoning-signature.ts new file mode 100644 index 00000000000..1071c87e71c --- /dev/null +++ b/src/adapters/devin/reasoning-signature.ts @@ -0,0 +1,94 @@ +/** + * Devin reasoning signatures across turns. + * + * GetChatMessage returns the turn's reasoning attestation as `delta_signature` + * (#10) together with `delta_signature_type` (#21) in the same frame, and the + * native client replays both on the assistant prompt as #12 and #18. Measured + * live on swe-2-high the pair arrives AFTER the visible answer (reasoning, text, + * then signature), so the Responses layer stores it as its own signature-only + * reasoning item behind the thinking-text item. GPT and Gemini rows stream no + * thinking text at all, only the signature. + * + * The type is carried inside the stored signature because the reasoning + * envelope that round-trips through the client keeps a single signature string. + * A signature stored before this prefix existed replays without a type. + */ +import type { OcxAssistantMessage } from "../../types"; +import { isProviderIssuedThinkingSignature } from "../../responses/reasoning-envelope"; + +const TYPED_SIGNATURE_PREFIX = "devin-sig1:"; + +export function encodeDevinSignature(signature: string, signatureType: string | undefined): string { + return signatureType && !signatureType.includes(":") + ? `${TYPED_SIGNATURE_PREFIX}${signatureType}:${signature}` + : signature; +} + +export function decodeDevinSignature(stored: string): { signature: string; signatureType?: string } { + if (!stored.startsWith(TYPED_SIGNATURE_PREFIX)) return { signature: stored }; + const rest = stored.slice(TYPED_SIGNATURE_PREFIX.length); + const colon = rest.indexOf(":"); + if (colon <= 0) return { signature: stored }; + return { signature: rest.slice(colon + 1), signatureType: rest.slice(0, colon) }; +} + +/** + * The assistant turn's reasoning for ChatMessagePrompt #11/#12/#18. + * + * All of the turn's thinking text is replayed. A signature rides along only + * when it covers exactly that text: + * - one thinking block carrying its own issued signature (a stray + * signature-only block does not displace it); + * - at most one unsigned thinking block plus exactly one signature-only block, + * which is how a single Devin turn arrives once its late #10 frame has been + * split into its own reasoning item. With no thinking block at all this is a + * GPT or Gemini row, where the signature is the only reasoning there is. + * Any other mix (two signed blocks, a signed block beside unsigned text) has no + * single attestation for the joined text, so the turn is replayed unsigned. + * + * `withholdAnthropic` drops an Anthropic signature and keeps the text: the + * fallback for a Claude turn Cognition refused (see hasAnthropicSignature). + */ +export function devinAssistantReasoning( + message: OcxAssistantMessage, + modelId = "", + withholdAnthropic = false, +): { thinking?: string; signature?: string; signature_type?: string } { + const blocks = message.content.filter( + (part): part is Extract => part.type === "thinking", + ); + const textBlocks = blocks.filter(part => Boolean(part.thinking)); + const signatureOnly = blocks.filter(part => !part.thinking && isProviderIssuedThinkingSignature(part.signature)); + const text = textBlocks.map(part => part.thinking).join("\n"); + let stored: string | undefined; + if (textBlocks.length === 1 && isProviderIssuedThinkingSignature(textBlocks[0]!.signature)) { + stored = textBlocks[0]!.signature; + } else if (textBlocks.length <= 1 && signatureOnly.length === 1) { + stored = signatureOnly[0]!.signature; + } + let decoded = stored ? decodeDevinSignature(stored) : undefined; + if (decoded && withholdAnthropic && signatureTypeFor(decoded, modelId) === "anthropic") decoded = undefined; + return { + ...(text ? { thinking: text } : {}), + ...(decoded ? { signature: decoded.signature } : {}), + ...(decoded?.signatureType ? { signature_type: decoded.signatureType } : {}), + }; +} + +/** A stored signature from before its type was recorded falls back to the model being called. */ +function signatureTypeFor(decoded: { signatureType?: string }, modelId: string): string | undefined { + return decoded.signatureType ?? (/claude/i.test(modelId) ? "anthropic" : undefined); +} + +/** + * True when the mapped history replays a Claude signature. Cognition streams + * Claude's thinking as a summary while the signature covers the original, so + * the pair can fail validation: live on claude-opus-5-5 a signed replay of a + * visible-thinking turn was refused with `invalid_argument` in 5 of 6 tries and + * a text-only one in none, while a signed replay that is accepted is what lets + * the model recall its earlier reasoning. The adapter therefore sends the + * signature and retries a refusal once without it. + */ +export function hasAnthropicSignature(items: ReadonlyArray<{ signature?: string; signature_type?: string }>, modelId: string): boolean { + return items.some(item => Boolean(item.signature) && signatureTypeFor({ signatureType: item.signature_type }, modelId) === "anthropic"); +} diff --git a/structure/providers-and-adapters.md b/structure/providers-and-adapters.md index ad719ac37ec..030ea0789ea 100644 --- a/structure/providers-and-adapters.md +++ b/structure/providers-and-adapters.md @@ -127,7 +127,7 @@ rewrite rules and the routed-id settlement. | `src/adapters/declaration-carrier.ts`, `src/adapters/input-media-guard.ts` | Default-deny allowlists for constraints the normalized request carries but a wire may not be able to express: `tools[*].allowed_callers`, which fences a tool off from callers, and inline document bytes. Both are refused with a 400 at the single guard every registered adapter passes through, rather than left to each adapter, because an adapter that never learned about the carrier rebuilds without it and answers normally. `allowed_callers` reaches the `anthropic` wire; document bytes reach `anthropic`, `openai-chat` and `google`; the `openai-responses` wire is exempt from the whole guard because it forwards the original body. Adding an `AdapterWire` member makes the omission visible in these lists instead of at a customer's upstream. The unrestricted `["direct"]` caller default is not a restriction. | | `src/adapters/azure.ts` | Azure OpenAI bridge. | | `src/adapters/cursor.ts`, `src/adapters/cursor/` | Cursor protobuf transport: discovery, request builder, event decoding, MCP, thread continuity, native-exec policy. | -| `src/adapters/devin.ts`, `src/adapters/devin/cloud-direct/` | Devin runTurn transport over Cognition Connect-RPC. `GetChatMessage` uses the Responses provider executor and shared physical-send budget; catalog, JWT, and `src/web-search/devin-executor.ts` native search support RPCs remain outside inference-send accounting. Provider-stated pre-output 429 reset delays are surfaced immediately by default, releasing shared active-turn capacity. A positive `OPENCODEX_DEVIN_STATED_RESET_WAIT_MS` explicitly enables bounded waiting and up to two replays on standalone turns, which hold that capacity until completion or cancellation. Combo children bypass that wait and surface a pre-output 429 so the next target can run. During an opted-in standalone wait, safe heartbeats commit the response preflight and keep the stream's stall watchdog fed. Invalid values fail closed to the immediate-refusal behavior. A recorded tenant host is used only for the stored account whose credential owns the transmitted key, searched in the configured provider id and then its deprecated alias; a configured, forwarded, or unmatched key uses the configured base URL or the US default. Native search previews the current route by effective adapter without mutating combo selection state, pins one admitted active-account snapshot for the request, and calls `GetWebSearchResults`, so it starts no CLI or second model. | +| `src/adapters/devin.ts`, `src/adapters/devin/cloud-direct/` | Devin runTurn transport over Cognition Connect-RPC. Assistant reasoning replays as ChatMessagePrompt #11 thinking, #12 signature and #18 signature type (`src/adapters/devin/reasoning-signature.ts`): the #10/#21 pair arrives after the visible answer and becomes its own signature-only reasoning item, so a single unsigned thinking block plus exactly one signature-only block is replayed as one signed prompt, and a signature-only turn (GPT, Gemini) is replayed rather than dropped. An Anthropic signature is replayed, but because the streamed thinking is a summary the signature may not cover, a turn Cognition refuses with `invalid_argument` before any visible output (reasoning alone does not count) is retried once with Anthropic signatures withheld and the thinking text kept. `GetChatMessage` uses the Responses provider executor and shared physical-send budget; catalog, JWT, and `src/web-search/devin-executor.ts` native search support RPCs remain outside inference-send accounting. Provider-stated pre-output 429 reset delays are surfaced immediately by default, releasing shared active-turn capacity. A positive `OPENCODEX_DEVIN_STATED_RESET_WAIT_MS` explicitly enables bounded waiting and up to two replays on standalone turns, which hold that capacity until completion or cancellation. Combo children bypass that wait and surface a pre-output 429 so the next target can run. During an opted-in standalone wait, safe heartbeats commit the response preflight and keep the stream's stall watchdog fed. Invalid values fail closed to the immediate-refusal behavior. A recorded tenant host is used only for the stored account whose credential owns the transmitted key, searched in the configured provider id and then its deprecated alias; a configured, forwarded, or unmatched key uses the configured base URL or the US default. Native search previews the current route by effective adapter without mutating combo selection state, pins one admitted active-account snapshot for the request, and calls `GetWebSearchResults`, so it starts no CLI or second model. | | `src/adapters/kiro.ts` and `src/adapters/kiro/` | Kiro event/tool/thinking/truncation/retry handling, including an egress-aware completion fallback, fixed public HTTP 5xx text, and closed-set status/code diagnostics. The original path is a facade over leaves for wire identity, reasoning, conversation state, token estimation, payload assembly, streaming, and the adapter. | | `src/adapters/mimo-free.ts` | Mimo Free transport (client identity + JWT). Concurrent requests share one JWT bootstrap bound only to its timeout; each request stops waiting on its own abort without cancelling the others. | | `src/adapters/command-code.ts`, `src/adapters/command-code-tool-text.ts`, `src/adapters/command-code-restored-schema.ts` | Command Code OAuth NDJSON translation. For every `xiaomi/mimo-` model, text, native calls, reasoning, and terminal decisions share one byte-bounded queue with linear queue visits. Markup is deduplicated against matching native calls; text-only restoration requires one contiguous text run, a clean finish, a declared tool, and arguments validated against supported schema constraints. A parameter-free (freeform) block may omit `` but must end with ``; parameter blocks keep the canonical close. Markup appended after prose in the same delta is split off at the marker and held like a block that opens with ``; a marker split across deltas after prose is still released as text. Native, reasoning, and other intervening events interrupt a still-probing block but leave a held block held in arrival order, and the queued byte bound still flushes an unresolved envelope as text. An envelope the strict parser rejects but that opens with ``, closes with ``, and names a declared function is dropped when a native call for that same function arrives and on a clean finish; markup that parses but fits no supported schema is still released as text. Regex patterns, other unsupported constraints, and abnormal finishes fail closed. `tests/providers/command-code-tool-text-prose-split.test.ts` covers the split, the interleaved-event hold, and both drop paths. | diff --git a/tests/fixtures/test-layout-expected.json b/tests/fixtures/test-layout-expected.json index 01768b541df..bb42b6245e5 100644 --- a/tests/fixtures/test-layout-expected.json +++ b/tests/fixtures/test-layout-expected.json @@ -759,6 +759,8 @@ "devin-provider-merge-migration.test.ts": "providers", "devin-stated-reset-hardening.test.ts": "providers", "devin-stated-reset-retry.test.ts": "providers", + "devin-anthropic-signature-fallback.test.ts": "providers", + "devin-reasoning-continuation.test.ts": "providers", "devin-stream-deadline.test.ts": "providers", "digitalocean-scaleway-provider.test.ts": "providers", "docs-429-failover-claims.test.ts": "ci-workflows", diff --git a/tests/providers/devin-anthropic-signature-fallback.test.ts b/tests/providers/devin-anthropic-signature-fallback.test.ts new file mode 100644 index 00000000000..4d9f2e54f7c --- /dev/null +++ b/tests/providers/devin-anthropic-signature-fallback.test.ts @@ -0,0 +1,184 @@ +import { mkdtempSync } from "node:fs"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import { afterEach, beforeEach, describe, expect, spyOn, test } from "bun:test"; +import { createDevinAdapter } from "../../src/adapters/devin"; +import { setCachedCatalogForTests } from "../../src/adapters/devin/cloud-direct/catalog"; +import { devinCacheIdentity, invalidateSessionIdentity } from "../../src/adapters/devin/cloud-direct/chat"; +import { encodeMessage, encodeString, encodeVarintField, iterFields } from "../../src/adapters/devin/cloud-direct/wire"; +import { encodeDevinSignature } from "../../src/adapters/devin/reasoning-signature"; +import { createTranslatorBudget } from "../../src/lib/translator-budget"; +import { encodeReasoningEnvelope } from "../../src/responses/reasoning-envelope"; +import { parseRequest } from "../../src/responses/parser"; +import type { AdapterEvent } from "../../src/types"; +import { removeTreeWithRetry } from "../helpers/remove-tree"; + +// Cognition streams Claude's thinking as a summary while the signature covers the +// original, so a signed replay can be refused with an opaque invalid_argument before +// any output. The adapter sends the signature (it is what carries the reasoning) and +// retries a refusal once without it. +describe("Devin Anthropic signature fallback", () => { + const apiKey = "ocx-devin-signature-fixture"; + const host = "https://server.codeium.com"; + const previousHome = process.env.OPENCODEX_HOME; + const previousFetch = globalThis.fetch; + let home = ""; + let requests: Buffer[] = []; + let responses: Array<"refuse" | "ok" | "text-then-refuse" | "reasoning-then-refuse" | "reasoning-then-ok" | "usage-reasoning-then-refuse" | "usage-ok"> = []; + + const frame = (body: Buffer, flags = 0) => { + const header = Buffer.alloc(5); + header[0] = flags; + header.writeUInt32BE(body.length, 1); + return Buffer.concat([header, body]); + }; + const refusal = frame(Buffer.from(JSON.stringify({ error: { code: "invalid_argument", message: "an internal error occurred" } })), 2); + const ok = Buffer.concat([frame(Buffer.concat([encodeString(3, "ok"), encodeVarintField(5, 2)])), frame(Buffer.from("{}"), 2)]); + + function assistantSignature(request: Buffer): { thinking?: string; signature?: string } { + const prompts = [...iterFields(request)].filter(f => f.num === 3).map(f => f.value as Buffer); + const assistant = prompts.find(p => [...iterFields(p)].some(f => f.num === 2 && f.value === 2n))!; + const byNum = new Map([...iterFields(assistant)].filter(f => f.wire === 2).map(f => [f.num, (f.value as Buffer).toString("utf8")])); + return { thinking: byNum.get(11), signature: byNum.get(12) }; + } + + async function run(signature: string, modelId: string): Promise { + const parsed = parseRequest({ + model: `devin/${modelId}`, + input: [ + { role: "user", content: [{ type: "input_text", text: "go" }] }, + { type: "reasoning", id: "rs", summary: [], encrypted_content: encodeReasoningEnvelope({ txt: "summarised thought", sig: signature }) }, + { type: "function_call", call_id: "call_1", name: "get_time", arguments: "{}" }, + { type: "function_call_output", call_id: "call_1", output: "12:00" }, + ], + }); + parsed.modelId = modelId; + const adapter = createDevinAdapter({ adapter: "devin", apiKey, baseUrl: host }); + const events: AdapterEvent[] = []; + await adapter.runTurn!(parsed, { headers: new Headers(), translatorBudget: createTranslatorBudget() }, event => { events.push(event); }); + return events; + } + + beforeEach(() => { + home = mkdtempSync(join(tmpdir(), "ocx-devin-sigfallback-")); + process.env.OPENCODEX_HOME = home; + requests = []; + responses = []; + setCachedCatalogForTests(null); + globalThis.fetch = (async (input: RequestInfo | URL, init?: RequestInit) => { + if (!String(input).endsWith("/GetChatMessage")) return new Response("unavailable", { status: 503 }); + requests.push(Buffer.from(await (init!.body as Blob).arrayBuffer()).subarray(5)); + const next = responses.shift() ?? "ok"; + const body = next === "refuse" ? refusal + : next === "text-then-refuse" ? Buffer.concat([frame(encodeString(3, "partial")), refusal]) + // The live shape: reasoning, its signature and a finish frame, then the refusal trailer. + : next === "reasoning-then-refuse" ? Buffer.concat([frame(Buffer.concat([encodeString(9, "thinking"), encodeString(10, "EpcBNew"), encodeString(21, "anthropic"), encodeVarintField(5, 2)])), refusal]) + : next === "reasoning-then-ok" ? Buffer.concat([frame(Buffer.concat([encodeString(9, "thinking"), encodeString(10, "EpcBNew"), encodeString(21, "anthropic")])), ok]) + // ModelUsageStats (#7) arrives with the reasoning, before the refusal trailer. + : next === "usage-reasoning-then-refuse" ? Buffer.concat([frame(Buffer.concat([encodeMessage(7, Buffer.concat([encodeVarintField(2, 1000), encodeVarintField(3, 40)])), encodeString(9, "thinking")])), frame(encodeString(9, " more")), refusal]) + : next === "usage-ok" ? Buffer.concat([frame(Buffer.concat([encodeMessage(7, Buffer.concat([encodeVarintField(2, 1100), encodeVarintField(3, 20)])), encodeString(3, "ok"), encodeVarintField(5, 2)])), frame(Buffer.from("{}"), 2)]) + : ok; + return new Response(body, { headers: { "content-type": "application/connect+proto" } }); + }) as typeof fetch; + }); + afterEach(() => { + globalThis.fetch = previousFetch; + setCachedCatalogForTests(null); + if (previousHome === undefined) delete process.env.OPENCODEX_HOME; + else process.env.OPENCODEX_HOME = previousHome; + invalidateSessionIdentity(devinCacheIdentity(apiKey, host)); + removeTreeWithRetry(home); + }); + + test("a refused signed Claude turn is retried once with the signature withheld", async () => { + responses = ["refuse", "ok"]; + const events = await run(encodeDevinSignature("EpcBClaude", "anthropic"), "claude-opus-5-5-medium"); + expect(requests).toHaveLength(2); + expect(assistantSignature(requests[0]!)).toEqual({ thinking: "summarised thought", signature: "EpcBClaude" }); + expect(assistantSignature(requests[1]!)).toEqual({ thinking: "summarised thought", signature: undefined }); + expect(events.some(e => e.type === "error")).toBe(false); + expect(events).toContainEqual({ type: "text_delta", text: "ok" }); + }); + + test("a refusal after reasoning alone is still retried, and the refused attempt's reasoning never reaches the client", async () => { + responses = ["reasoning-then-refuse", "ok"]; + const events = await run(encodeDevinSignature("EpcBClaude", "anthropic"), "claude-opus-5-5-medium"); + expect(requests).toHaveLength(2); + expect(assistantSignature(requests[1]!).signature).toBeUndefined(); + expect(events.some(e => e.type === "error")).toBe(false); + expect(events).toContainEqual({ type: "text_delta", text: "ok" }); + // The refused attempt streamed "thinking" and signature EpcBNew; neither may leak into the turn. + expect(events.some(e => e.type === "thinking_delta")).toBe(false); + expect(events.some(e => e.type === "thinking_signature")).toBe(false); + }); + + test("an accepted signed turn still delivers its held reasoning", async () => { + responses = ["reasoning-then-ok"]; + const events = await run(encodeDevinSignature("EpcBClaude", "anthropic"), "claude-opus-5-5-medium"); + expect(requests).toHaveLength(1); + const kinds = events.map(e => e.type); + expect(kinds).toContain("thinking_delta"); + expect(events).toContainEqual({ type: "thinking_signature", signature: encodeDevinSignature("EpcBNew", "anthropic") }); + expect(kinds.indexOf("thinking_delta")).toBeLessThan(kinds.indexOf("text_delta")); + }); + + test("the refused attempt's usage is added to the retry's", async () => { + responses = ["usage-reasoning-then-refuse", "usage-ok"]; + const events = await run(encodeDevinSignature("EpcBClaude", "anthropic"), "claude-opus-5-5-medium"); + expect(requests).toHaveLength(2); + const done = events.find(e => e.type === "done") as { usage?: { inputTokens?: number; outputTokens?: number } } | undefined; + expect(done?.usage?.inputTokens).toBe(2100); + expect(done?.usage?.outputTokens).toBe(60); + }); + + test("the refused attempt's usage survives a retry with no usage frame or an early failure", async () => { + responses = ["usage-reasoning-then-refuse", "ok"]; + const done = (await run(encodeDevinSignature("EpcBClaude", "anthropic"), "claude-opus-5-5-medium")) + .find(e => e.type === "done") as { usage?: { inputTokens?: number; outputTokens?: number } } | undefined; + expect(done?.usage?.inputTokens).toBe(1000); + expect(done?.usage?.outputTokens).toBe(40); + + requests = []; + responses = ["usage-reasoning-then-refuse", "refuse"]; + const failed = (await run(encodeDevinSignature("EpcBClaude", "anthropic"), "claude-opus-5-5-medium")) + .find(e => e.type === "error") as { usage?: { inputTokens?: number } } | undefined; + expect(requests).toHaveLength(2); + expect(failed?.usage?.inputTokens).toBe(1000); + }); + + test("held reasoning emits heartbeats, never the held events", async () => { + // Each clock read advances 20s, so every held frame is past the heartbeat interval. + let clock = Date.now(); + const now = spyOn(Date, "now").mockImplementation(() => (clock += 20_000)); + try { + responses = ["usage-reasoning-then-refuse", "ok"]; + const events = await run(encodeDevinSignature("EpcBClaude", "anthropic"), "claude-opus-5-5-medium"); + const kinds = events.map(e => e.type); + expect(kinds.filter(k => k === "heartbeat").length).toBeGreaterThan(0); + expect(kinds).not.toContain("thinking_delta"); + expect(kinds.indexOf("heartbeat")).toBeLessThan(kinds.indexOf("text_delta")); + } finally { + now.mockRestore(); + } + }); + + test("an accepted signed Claude turn is sent once, signature included", async () => { + responses = ["ok"]; + await run(encodeDevinSignature("EpcBClaude", "anthropic"), "claude-opus-5-5-medium"); + expect(requests).toHaveLength(1); + expect(assistantSignature(requests[0]!).signature).toBe("EpcBClaude"); + }); + + test("a refusal is not retried for a non-Anthropic signature or after output", async () => { + responses = ["refuse", "ok"]; + const sealed = await run(encodeDevinSignature("sealed.v1.x", "sealed"), "swe-2-high"); + expect(requests).toHaveLength(1); + expect(sealed.some(e => e.type === "error")).toBe(true); + + requests = []; + responses = ["text-then-refuse", "ok"]; + const partial = await run(encodeDevinSignature("EpcBClaude", "anthropic"), "claude-opus-5-5-medium"); + expect(requests).toHaveLength(1); + expect(partial.some(e => e.type === "error")).toBe(true); + }); +}); diff --git a/tests/providers/devin-reasoning-continuation.test.ts b/tests/providers/devin-reasoning-continuation.test.ts new file mode 100644 index 00000000000..58a76d9411b --- /dev/null +++ b/tests/providers/devin-reasoning-continuation.test.ts @@ -0,0 +1,136 @@ +import { describe, expect, test } from "bun:test"; +import { mapOcxMessagesToDevin } from "../../src/adapters/devin"; +import { buildGetChatMessageRequestForTests, decodeChatFrame } from "../../src/adapters/devin/cloud-direct/chat"; +import { encodeString, iterFields } from "../../src/adapters/devin/cloud-direct/wire"; +import { decodeDevinSignature, encodeDevinSignature, hasAnthropicSignature } from "../../src/adapters/devin/reasoning-signature"; +import { encodeReasoningEnvelope } from "../../src/responses/reasoning-envelope"; +import { parseRequest } from "../../src/responses/parser"; + +// Shapes measured live on GetChatMessage: swe-2-high streams reasoning, then the +// visible answer, then one frame carrying #10 delta_signature and #21 +// delta_signature_type ("sealed"); gpt-6-sol streams no thinking text and a +// signature of type "openai". +const SEALED = "sealed.v1.opaque-attestation"; + +function assistantPrompt(history: ReturnType): Map { + const request = buildGetChatMessageRequestForTests({ + apiKey: "devin-session-token$x", modelUid: "swe-2-high", messages: history, cascadeId: "c", + } as never); + const prompts = [...iterFields(request)].filter(f => f.num === 3).map(f => f.value as Buffer); + const assistant = prompts.find(p => [...iterFields(p)].some(f => f.num === 2 && f.value === 2n))!; + return new Map([...iterFields(assistant)].filter(f => f.wire === 2).map(f => [f.num, (f.value as Buffer).toString("utf8")])); +} + +describe("Devin reasoning continuation across turns", () => { + test("the signature frame yields its type", () => { + const frame = Buffer.concat([encodeString(10, SEALED), encodeString(21, "sealed")]); + expect([...decodeChatFrame(frame)]).toContainEqual({ kind: "reasoning_signature", signature: SEALED, signatureType: "sealed" }); + expect([...decodeChatFrame(encodeString(10, SEALED))]).toContainEqual({ kind: "reasoning_signature", signature: SEALED }); + }); + + test("the stored signature carries its type and an older stored signature still replays", () => { + const stored = encodeDevinSignature(SEALED, "sealed"); + expect(decodeDevinSignature(stored)).toEqual({ signature: SEALED, signatureType: "sealed" }); + expect(decodeDevinSignature(SEALED)).toEqual({ signature: SEALED }); + expect(encodeDevinSignature(SEALED, undefined)).toBe(SEALED); + }); + + test("a SWE-2 turn split into a thinking item and a late signature item replays as one signed prompt", () => { + // What the client sends back: the thinking summary item (no envelope) and the + // signature-only item the late #10 frame became, then the tool loop. + const history = mapOcxMessagesToDevin(parseRequest({ + model: "devin/swe-2", + input: [ + { role: "user", content: [{ type: "input_text", text: "go" }] }, + { type: "reasoning", id: "rs_text", summary: [{ type: "summary_text", text: "pick 482916, then call the tool" }] }, + { type: "reasoning", id: "rs_sig", summary: [], encrypted_content: encodeReasoningEnvelope({ sig: encodeDevinSignature(SEALED, "sealed") }) }, + { type: "function_call", call_id: "call_1", name: "get_time", arguments: "{}" }, + { type: "function_call_output", call_id: "call_1", output: "12:00" }, + ], + })); + const assistant = history.find(m => m.role === "assistant"); + expect(assistant?.thinking).toBe("pick 482916, then call the tool"); + expect(assistant?.signature).toBe(SEALED); + expect(assistant?.signature_type).toBe("sealed"); + const wire = assistantPrompt(history); + expect(wire.get(11)).toBe("pick 482916, then call the tool"); + expect(wire.get(12)).toBe(SEALED); + expect(wire.get(18)).toBe("sealed"); + }); + + test("a signature-only turn is replayed instead of dropped", () => { + const openaiSig = '[{"id":"rs_1","encrypted_content":"opaque"}]'; + const history = mapOcxMessagesToDevin(parseRequest({ + model: "devin/gpt-6-sol", + input: [ + { role: "user", content: [{ type: "input_text", text: "go" }] }, + { type: "reasoning", id: "rs_sig", summary: [], encrypted_content: encodeReasoningEnvelope({ sig: encodeDevinSignature(openaiSig, "openai") }) }, + { type: "function_call", call_id: "call_1", name: "get_time", arguments: "{}" }, + { type: "function_call_output", call_id: "call_1", output: "12:00" }, + ], + })); + const assistant = history.find(m => m.role === "assistant"); + expect(assistant?.thinking).toBeUndefined(); + expect(assistant?.signature).toBe(openaiSig); + expect(assistant?.signature_type).toBe("openai"); + const wire = assistantPrompt(history); + expect(wire.has(11)).toBe(false); + expect(wire.get(12)).toBe(openaiSig); + expect(wire.get(18)).toBe("openai"); + + // No text, no tool call: the signature alone still keeps the assistant turn. + const bare = mapOcxMessagesToDevin(parseRequest({ + model: "devin/gemini-3-8-flash", + input: [ + { role: "user", content: [{ type: "input_text", text: "go" }] }, + { type: "reasoning", id: "rs_sig", summary: [], encrypted_content: encodeReasoningEnvelope({ sig: encodeDevinSignature("AY89gemini", "gemini") }) }, + { role: "user", content: [{ type: "input_text", text: "and then?" }] }, + ], + })); + const kept = bare.find(m => m.role === "assistant"); + expect(kept?.signature).toBe("AY89gemini"); + expect(kept?.signature_type).toBe("gemini"); + expect(assistantPrompt(bare).get(12)).toBe("AY89gemini"); + }); + + test("two late signatures beside one thinking block cannot be paired", () => { + const history = mapOcxMessagesToDevin(parseRequest({ + model: "devin/swe-2", + input: [ + { role: "user", content: [{ type: "input_text", text: "go" }] }, + { type: "reasoning", id: "rs_text", summary: [{ type: "summary_text", text: "thought" }] }, + { type: "reasoning", id: "rs_a", summary: [], encrypted_content: encodeReasoningEnvelope({ sig: "sealed.v1.a" }) }, + { type: "reasoning", id: "rs_b", summary: [], encrypted_content: encodeReasoningEnvelope({ sig: "sealed.v1.b" }) }, + { type: "message", role: "assistant", content: [{ type: "output_text", text: "answer" }] }, + ], + })); + const assistant = history.find(m => m.role === "assistant"); + expect(assistant?.thinking).toBe("thought"); + expect(assistant?.signature).toBeUndefined(); + }); + + test("an Anthropic signature is replayed by default and withheld only for the fallback", () => { + const parsed = (sig: string, model: string) => parseRequest({ + model, + input: [ + { role: "user", content: [{ type: "input_text", text: "go" }] }, + { type: "reasoning", id: "rs", summary: [], encrypted_content: encodeReasoningEnvelope({ txt: "summarised thought", sig }) }, + { type: "function_call", call_id: "call_1", name: "get_time", arguments: "{}" }, + { type: "function_call_output", call_id: "call_1", output: "12:00" }, + ], + }); + const typed = parsed(encodeDevinSignature("EpcBClaude", "anthropic"), "devin/claude-opus-5-5"); + const signed = mapOcxMessagesToDevin(typed); + expect(signed.find(m => m.role === "assistant")?.signature).toBe("EpcBClaude"); + expect(hasAnthropicSignature(signed, typed.modelId)).toBe(true); + const unsigned = mapOcxMessagesToDevin(typed, { withholdAnthropicSignatures: true }).find(m => m.role === "assistant"); + expect(unsigned?.thinking).toBe("summarised thought"); + expect(unsigned?.signature).toBeUndefined(); + // A signature stored before its type was recorded falls back to the model being called. + const legacy = parsed("EpcBClaude", "devin/claude-opus-5-5"); + expect(hasAnthropicSignature(mapOcxMessagesToDevin(legacy), legacy.modelId)).toBe(true); + const sealed = parsed(SEALED, "devin/swe-2"); + expect(hasAnthropicSignature(mapOcxMessagesToDevin(sealed), sealed.modelId)).toBe(false); + expect(mapOcxMessagesToDevin(sealed, { withholdAnthropicSignatures: true }).find(m => m.role === "assistant")?.signature).toBe(SEALED); + }); +});