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 <amp@ampcode.com>
This commit is contained in:
Peter Steinberger
2026-08-12 10:53:53 -07:00
committed by GitHub
parent 2e86f7cc95
commit 39c733ff73
2 changed files with 67 additions and 79 deletions
+6
View File
@@ -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) => {
+61 -79
View File
@@ -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,