// Openrouter plugin module implements stream behavior. import type { StreamFn } from "openclaw/plugin-sdk/agent-core"; import type { ProviderWrapStreamFnContext } from "openclaw/plugin-sdk/plugin-entry"; import { buildProviderStreamFamilyHooks } from "openclaw/plugin-sdk/provider-stream-family"; import { composeProviderStreamWrappers, createPayloadPatchStreamWrapper, normalizeOpenAICompatibleReasoningReplay, } from "openclaw/plugin-sdk/provider-stream-shared"; import { createSubsystemLogger } from "openclaw/plugin-sdk/runtime-env"; import { asNonArrayRecord, readStringValue } from "openclaw/plugin-sdk/string-coerce-runtime"; import { isOpenRouterDeepSeekV4ModelId, normalizeOpenRouterModelFamilyId } from "./models.js"; import { isOpenRouterProxyReasoningUnsupportedModel, normalizeOpenRouterBaseUrl, OPENROUTER_BASE_URL, } from "./provider-catalog.js"; const log = createSubsystemLogger("openrouter-stream"); const openRouterThinkingStreamHooks = buildProviderStreamFamilyHooks("openrouter-thinking"); function normalizeOpenRouterStringPreservingEmpty(value: unknown): string | undefined { return readStringValue(value)?.trim(); } function isVerifiedOpenRouterRoute(model: Parameters[0]): boolean { const provider = normalizeOpenRouterStringPreservingEmpty(model.provider)?.toLowerCase(); const baseUrl = normalizeOpenRouterStringPreservingEmpty(model.baseUrl); if (baseUrl) { return normalizeOpenRouterBaseUrl(baseUrl) === OPENROUTER_BASE_URL; } return provider === "openrouter"; } function shouldPatchAnthropicOpenRouterPayload(model: Parameters[0]): boolean { const api = normalizeOpenRouterStringPreservingEmpty(model.api); return ( (api === undefined || api === "openai-completions") && normalizeOpenRouterModelFamilyId(model.id)?.startsWith("anthropic/") === true && isVerifiedOpenRouterRoute(model) ); } function shouldPatchDeepSeekV4OpenRouterPayload(model: Parameters[0]): boolean { const api = normalizeOpenRouterStringPreservingEmpty(model.api); return ( (api === undefined || api === "openai-completions") && isOpenRouterDeepSeekV4ModelId(model.id) && isVerifiedOpenRouterRoute(model) ); } function shouldPatchOpenRouterRoutingPayload(model: Parameters[0]): boolean { const api = normalizeOpenRouterStringPreservingEmpty(model.api); return (api === undefined || api === "openai-completions") && isVerifiedOpenRouterRoute(model); } function mergeOpenRouterAuthHeaders(options: Parameters[2]): Parameters[2] { const apiKey = normalizeOpenRouterStringPreservingEmpty(options?.apiKey); if (!apiKey) { return options; } const headers = new Headers((options as { headers?: HeadersInit } | undefined)?.headers); if (!headers.has("authorization")) { headers.set("Authorization", `Bearer ${apiKey}`); } if (!headers.has("http-referer")) { headers.set("HTTP-Referer", "https://openclaw.ai"); } if (!headers.has("x-openrouter-title")) { headers.set("X-OpenRouter-Title", "OpenClaw"); } return { ...options, headers: Object.fromEntries(headers.entries()), } as Parameters[2]; } function createOpenRouterAuthHeaderWrapper( baseStreamFn: StreamFn | undefined, ): StreamFn | undefined { if (!baseStreamFn) { return baseStreamFn; } return (model, context, options) => baseStreamFn( model, context, isVerifiedOpenRouterRoute(model) ? mergeOpenRouterAuthHeaders(options) : options, ); } function assistantMessageHasOpenAIToolCalls(message: Record): boolean { return Array.isArray(message.tool_calls) && message.tool_calls.length > 0; } function isAnthropicToolCallContentBlock(value: unknown): boolean { return ( value !== null && typeof value === "object" && ((value as { type?: unknown }).type === "tool_use" || (value as { type?: unknown }).type === "toolCall") ); } function assistantMessageHasAnthropicToolUse(message: Record): boolean { const content = message.content; return Array.isArray(content) && content.some(isAnthropicToolCallContentBlock); } function shouldStripOpenRouterTrailingMessage(value: unknown): boolean { if (!value || typeof value !== "object") { return false; } const message = value as Record; return ( message.role === "assistant" && !assistantMessageHasOpenAIToolCalls(message) && !assistantMessageHasAnthropicToolUse(message) ); } function stripTrailingOpenRouterAssistantPrefillMessages(payload: Record): number { const messages = payload.messages; if (!Array.isArray(messages)) { return 0; } let keep = messages.length; while (keep > 0 && shouldStripOpenRouterTrailingMessage(messages[keep - 1])) { keep -= 1; } if (keep === messages.length) { return 0; } const stripped = messages.length - keep; messages.splice(keep); return stripped; } function isEnabledReasoningValue(value: unknown): boolean { if (value === undefined || value === null || value === false) { return false; } if (typeof value === "string") { const normalized = value.trim().toLowerCase(); return normalized !== "" && normalized !== "off" && normalized !== "none"; } if (typeof value === "object" && !Array.isArray(value)) { const reasoning = value as Record; if (reasoning.enabled === false) { return false; } const effort = reasoning.effort; if (typeof effort === "string") { const normalized = effort.trim().toLowerCase(); return normalized !== "" && normalized !== "off" && normalized !== "none"; } } return true; } function isOpenRouterReasoningPayloadEnabled(payload: Record): boolean { return ( isEnabledReasoningValue(payload.reasoning) || isEnabledReasoningValue(payload.reasoning_effort) ); } function injectOpenRouterRouting( baseStreamFn: StreamFn | undefined, providerRouting?: Record, ): StreamFn | undefined { if (!providerRouting) { return baseStreamFn; } const routedStreamFn: StreamFn = (model, context, options) => ( baseStreamFn ?? ((nextModel) => { throw new Error( `OpenRouter routing wrapper requires an underlying streamFn for ${nextModel.id}.`, ); }) )( { ...model, compat: { ...model.compat, openRouterRouting: providerRouting }, } as typeof model, context, options, ); return createPayloadPatchStreamWrapper( routedStreamFn, ({ payload }) => { if (payload.provider === undefined) { payload.provider = providerRouting; } }, { shouldPatch: ({ model }) => shouldPatchOpenRouterRoutingPayload(model), }, ); } function createOpenRouterAnthropicPrefillWrapper(baseStreamFn: StreamFn | undefined): StreamFn { return createPayloadPatchStreamWrapper( baseStreamFn, ({ payload }) => { if (!isOpenRouterReasoningPayloadEnabled(payload)) { return; } const stripped = stripTrailingOpenRouterAssistantPrefillMessages(payload); if (stripped > 0) { log.warn( `removed ${stripped} trailing assistant prefill message${stripped === 1 ? "" : "s"} because OpenRouter-routed Anthropic reasoning requires conversations to end with a user turn`, ); } }, { shouldPatch: ({ model }) => shouldPatchAnthropicOpenRouterPayload(model), }, ); } function resolveOpenRouterDeepSeekV4ReasoningEffort( thinkingLevel: ProviderWrapStreamFnContext["thinkingLevel"], ): "high" | "xhigh" | undefined { if (thinkingLevel === "off") { return undefined; } if (thinkingLevel === "xhigh" || thinkingLevel === "max") { return "xhigh"; } return "high"; } function applyOpenRouterDeepSeekV4ReasoningEffort( payload: Record, thinkingLevel: ProviderWrapStreamFnContext["thinkingLevel"], ): boolean { const effort = resolveOpenRouterDeepSeekV4ReasoningEffort(thinkingLevel); if (!effort) { delete payload.reasoning; return false; } const reasoning = asNonArrayRecord(payload.reasoning); reasoning.effort = effort; payload.reasoning = reasoning; return true; } function createOpenRouterDeepSeekV4ReplayWrapper( baseStreamFn: StreamFn | undefined, thinkingLevel: ProviderWrapStreamFnContext["thinkingLevel"], ): StreamFn { return createPayloadPatchStreamWrapper( baseStreamFn, ({ payload }) => { delete payload.thinking; delete payload.reasoning_effort; normalizeOpenAICompatibleReasoningReplay(payload, { thinkingEnabled: applyOpenRouterDeepSeekV4ReasoningEffort(payload, thinkingLevel), shouldBackfillAssistantMessage: (message) => !assistantMessageHasOpenAIToolCalls(message), }); }, { shouldPatch: ({ model }) => shouldPatchDeepSeekV4OpenRouterPayload(model), }, ); } export function wrapOpenRouterProviderStream( ctx: ProviderWrapStreamFnContext, ): StreamFn | null | undefined { const providerRouting = ctx.extraParams?.provider != null && typeof ctx.extraParams.provider === "object" ? (ctx.extraParams.provider as Record) : undefined; const routedStreamFn = providerRouting ? injectOpenRouterRouting(ctx.streamFn, providerRouting) : ctx.streamFn; const wrapStreamFn = openRouterThinkingStreamHooks.wrapStreamFn ?? undefined; return composeProviderStreamWrappers( routedStreamFn, wrapStreamFn && ((streamFn) => wrapStreamFn({ ...ctx, streamFn, thinkingLevel: isOpenRouterProxyReasoningUnsupportedModel(ctx.modelId) ? undefined : ctx.thinkingLevel, }) ?? undefined), (streamFn) => createOpenRouterDeepSeekV4ReplayWrapper(streamFn, ctx.thinkingLevel), createOpenRouterAuthHeaderWrapper, createOpenRouterAnthropicPrefillWrapper, ); }