From 39c733ff731e171606b5178264cda86a35e8d9d2 Mon Sep 17 00:00:00 2001 From: Peter Steinberger Date: Wed, 12 Aug 2026 10:53:53 -0700 Subject: [PATCH] refactor(opencode-go): compose provider stream wrappers (#122749) Amp-Thread-ID: https://ampcode.com/threads/T-019ff438-3b93-77b8-9828-1d3c586cb127 Co-authored-by: Amp --- extensions/opencode-go/index.test.ts | 6 ++ extensions/opencode-go/stream.ts | 140 ++++++++++++--------------- 2 files changed, 67 insertions(+), 79 deletions(-) diff --git a/extensions/opencode-go/index.test.ts b/extensions/opencode-go/index.test.ts index a5374b5eccb3..e0200af56285 100644 --- a/extensions/opencode-go/index.test.ts +++ b/extensions/opencode-go/index.test.ts @@ -586,6 +586,12 @@ describe("opencode-go provider plugin", () => { ); }); + it("does not synthesize a stream when the runtime provides none", async () => { + const provider = await registerSingleProviderPlugin(plugin); + + expect(provider.wrapStreamFn?.({ streamFn: undefined } as never)).toBeUndefined(); + }); + it.each(["deepseek-v4-pro", "deepseek-v4-flash"] as const)( "disables invalid DeepSeek V4 reasoning_effort off payloads on OpenCode Go for %s", async (modelId) => { diff --git a/extensions/opencode-go/stream.ts b/extensions/opencode-go/stream.ts index 19033fcbbfc6..e6ab51203dc5 100644 --- a/extensions/opencode-go/stream.ts +++ b/extensions/opencode-go/stream.ts @@ -1,6 +1,7 @@ // Opencode Go plugin module implements stream behavior. import type { ProviderWrapStreamFnContext } from "openclaw/plugin-sdk/plugin-entry"; import { + composeProviderStreamWrappers, createDeepSeekV4OpenAICompatibleThinkingWrapper, createOpenAICompatibleCompletionsThinkingOffWrapper, createPayloadPatchStreamWrapper, @@ -14,76 +15,6 @@ import { OPENCODE_GO_STREAM_IDLE_TIMEOUT_MS_DEFAULT, } from "./stream-termination.js"; -function createOpencodeGoDeepSeekV4Wrapper( - baseStreamFn: ProviderWrapStreamFnContext["streamFn"], - thinkingLevel: ProviderWrapStreamFnContext["thinkingLevel"], -): ProviderWrapStreamFnContext["streamFn"] { - const flashWrapped = createDeepSeekV4OpenAICompatibleThinkingWrapper({ - baseStreamFn, - thinkingLevel, - shouldPatchModel: (model) => - model.provider === "opencode-go" && model.id === "deepseek-v4-flash", - resolveReasoningEffort: (level) => (level === "low" ? "low" : level === "max" ? "max" : "high"), - }); - return createDeepSeekV4OpenAICompatibleThinkingWrapper({ - baseStreamFn: flashWrapped, - thinkingLevel, - shouldPatchModel: (model) => model.provider === "opencode-go" && model.id === "deepseek-v4-pro", - }); -} - -function createOpencodeGoKimiNoReasoningWrapper( - baseStreamFn: ProviderWrapStreamFnContext["streamFn"], -): ProviderWrapStreamFnContext["streamFn"] { - if (!baseStreamFn) { - return undefined; - } - return createPayloadPatchStreamWrapper( - baseStreamFn, - ({ payload }) => stripOpencodeGoKimiReasoningPayload(payload), - { - shouldPatch: ({ model }) => - model.provider === "opencode-go" && isOpencodeGoKimiNoReasoningModelId(model.id), - }, - ); -} - -function createOpencodeGoFixedAnthropicReasoningWrapper( - baseStreamFn: ProviderWrapStreamFnContext["streamFn"], -): ProviderWrapStreamFnContext["streamFn"] { - if (!baseStreamFn) { - return undefined; - } - return createPayloadPatchStreamWrapper( - baseStreamFn, - ({ payload }) => { - delete payload.thinking; - delete payload.output_config; - }, - { - shouldPatch: ({ model }) => - model.provider === "opencode-go" && isOpencodeGoFixedAnthropicReasoningModelId(model.id), - }, - ); -} - -function createOpencodeGoKimiK3ThinkingOffWrapper( - baseStreamFn: ProviderWrapStreamFnContext["streamFn"], - thinkingLevel: ProviderWrapStreamFnContext["thinkingLevel"], -): ProviderWrapStreamFnContext["streamFn"] { - if (!baseStreamFn) { - return undefined; - } - const thinkingOff = createOpenAICompatibleCompletionsThinkingOffWrapper( - baseStreamFn, - thinkingLevel, - ); - return (model, context, options) => - model.provider === "opencode-go" && model.id === "kimi-k3" - ? thinkingOff(model, context, options) - : baseStreamFn(model, context, options); -} - export function createOpencodeGoWrapper( baseStreamFn: ProviderWrapStreamFnContext["streamFn"], thinkingLevel: ProviderWrapStreamFnContext["thinkingLevel"], @@ -91,18 +22,69 @@ export function createOpencodeGoWrapper( if (!baseStreamFn) { return undefined; } - const kimiWrapped = createOpencodeGoKimiNoReasoningWrapper(baseStreamFn) ?? baseStreamFn; - const kimiK3Wrapped = - createOpencodeGoKimiK3ThinkingOffWrapper(kimiWrapped, thinkingLevel) ?? kimiWrapped; - const fixedAnthropicWrapped = - createOpencodeGoFixedAnthropicReasoningWrapper(kimiK3Wrapped) ?? kimiK3Wrapped; - const deepSeekWrapped = - createOpencodeGoDeepSeekV4Wrapper(fixedAnthropicWrapped, thinkingLevel) ?? - fixedAnthropicWrapped; + const wrapped = + composeProviderStreamWrappers( + baseStreamFn, + (streamFn) => + streamFn + ? createPayloadPatchStreamWrapper( + streamFn, + ({ payload }) => stripOpencodeGoKimiReasoningPayload(payload), + { + shouldPatch: ({ model }) => + model.provider === "opencode-go" && isOpencodeGoKimiNoReasoningModelId(model.id), + }, + ) + : undefined, + (streamFn) => { + if (!streamFn) { + return undefined; + } + const thinkingOff = createOpenAICompatibleCompletionsThinkingOffWrapper( + streamFn, + thinkingLevel, + ); + return (model, context, options) => + model.provider === "opencode-go" && model.id === "kimi-k3" + ? thinkingOff(model, context, options) + : streamFn(model, context, options); + }, + (streamFn) => + streamFn + ? createPayloadPatchStreamWrapper( + streamFn, + ({ payload }) => { + delete payload.thinking; + delete payload.output_config; + }, + { + shouldPatch: ({ model }) => + model.provider === "opencode-go" && + isOpencodeGoFixedAnthropicReasoningModelId(model.id), + }, + ) + : undefined, + (streamFn) => + createDeepSeekV4OpenAICompatibleThinkingWrapper({ + baseStreamFn: streamFn, + thinkingLevel, + shouldPatchModel: (model) => + model.provider === "opencode-go" && model.id === "deepseek-v4-flash", + resolveReasoningEffort: (level) => + level === "low" ? "low" : level === "max" ? "max" : "high", + }) ?? streamFn, + (streamFn) => + createDeepSeekV4OpenAICompatibleThinkingWrapper({ + baseStreamFn: streamFn, + thinkingLevel, + shouldPatchModel: (model) => + model.provider === "opencode-go" && model.id === "deepseek-v4-pro", + }) ?? streamFn, + ) ?? baseStreamFn; // Outermost layer: provider-owned stalled SSE termination so the underlying // OpenAI SDK request is aborted at the raw opencode-go boundary instead of // waiting for the shared runtime stuck-session recovery. - return createOpencodeGoStalledStreamWrapper(deepSeekWrapped, { + return createOpencodeGoStalledStreamWrapper(wrapped, { provider: "opencode-go", idleTimeoutMs: OPENCODE_GO_STREAM_IDLE_TIMEOUT_MS_DEFAULT, firstEventTimeoutMs: OPENCODE_GO_STREAM_FIRST_EVENT_TIMEOUT_MS_DEFAULT,