Files
Peter Steinberger 9d77dad313 refactor(providers): deslop model provider plugins (#158403)
* refactor(providers): deslop model provider plugins

* refactor(providers): finish validation cleanup

* fix(openai): preserve cold provider catalog imports

Keep promise construction local on cold catalog paths to avoid loading logging redaction through extension-shared. Withdraw the unrelated ElevenLabs header-only edit.
2026-09-26 00:32:42 -07:00

578 lines
20 KiB
TypeScript

import { mkdtemp, rm } from "node:fs/promises";
import path from "node:path";
import process from "node:process";
import { pruneMapToMaxSize } from "openclaw/plugin-sdk/collection-runtime";
import { createSubsystemLogger } from "openclaw/plugin-sdk/logging-core";
import { runCommandBuffered } from "openclaw/plugin-sdk/process-runtime";
import type {
SessionCatalogSession,
SessionCatalogTranscriptItem,
SessionsCatalogReadResult,
} from "openclaw/plugin-sdk/session-catalog";
import { sessionCatalogPaging } from "openclaw/plugin-sdk/session-catalog";
import {
isRecord,
normalizeBoundedOptionalString as optionalOpenCodeString,
} from "openclaw/plugin-sdk/string-coerce-runtime";
import { resolvePreferredOpenClawTmpDir } from "openclaw/plugin-sdk/temp-path";
import { truncateUtf16Safe } from "openclaw/plugin-sdk/text-utility-runtime";
import {
materializeWindowsSpawnProgram,
resolveWindowsSpawnProgram,
} from "openclaw/plugin-sdk/windows-spawn";
import {
OPENCODE_SESSION_CATALOG_MAX_PAGE_LIMIT,
OPENCODE_SESSION_ID_PATTERN,
} from "./session-catalog-shared.js";
const LOCAL_HOST_ID = "gateway";
const MAX_SEARCH_LENGTH = 500;
const MAX_CLI_LIST_SESSIONS = 10_000;
const MAX_CLI_OUTPUT_BYTES = 32 * 1024 * 1024;
const CLI_TIMEOUT_MS = 30_000;
const OPENCODE_QUERY_CACHE_TTL_MS = 32_000;
const OPENCODE_QUERY_CACHE_MAX_ENTRIES = 32;
const SAFE_ENV_KEYS = [
"APPDATA",
"COMSPEC",
"HOME",
"LANG",
"LC_ALL",
"LOCALAPPDATA",
"OPENCODE_DB",
"PATH",
"Path",
"PATHEXT",
"SYSTEMROOT",
"TEMP",
"TMP",
"TMPDIR",
"USERPROFILE",
"WINDIR",
"XDG_CACHE_HOME",
"XDG_CONFIG_HOME",
"XDG_DATA_HOME",
"XDG_STATE_HOME",
] as const;
type OpenCodeQueryCacheEntry = {
expiresAt: number;
result: Promise<unknown>;
resolved?: true;
};
type OpenCodeQueryCacheOptions = {
configIdentity?: object;
forceRefresh?: boolean;
};
const openCodeConfigIdentities = new WeakMap<object, number>();
// Query results are valid for one immutable OpenClaw config identity, CLI environment, and SQL text.
// Config/env changes or 32s expiry invalidate them; failures are removed so recovery retries at once.
// The bounded map prevents pagination variants from growing while avoiding a subprocess every poll.
const openCodeQueryCache = new Map<string, OpenCodeQueryCacheEntry>();
let nextOpenCodeConfigIdentity = 1;
const log = createSubsystemLogger("opencode/session-catalog");
let cliVersion: { key: string; expiresAt: number; result: Promise<boolean> } | undefined;
export async function isOpenCodeV2(): Promise<boolean> {
const key = SAFE_ENV_KEYS.map((name) => process.env[name] ?? "").join("\0");
if (!cliVersion || cliVersion.key !== key || cliVersion.expiresAt <= Date.now()) {
const result = runOpenCode(["--version"]).then((output) => {
const major = /^(?:opencode v)?(\d+)\./.exec(output.trim())?.[1];
if (major !== "1" && major !== "2") {
throw new Error("Unsupported OpenCode CLI version");
}
return major === "2";
});
cliVersion = { key, expiresAt: Date.now() + OPENCODE_QUERY_CACHE_TTL_MS, result };
result.catch(() => {
if (cliVersion?.result === result) {
cliVersion = undefined;
}
});
}
return cliVersion.result;
}
function openCodeQueryCacheKey(query: string, configIdentity: object): string {
let identity = openCodeConfigIdentities.get(configIdentity);
if (identity === undefined) {
identity = nextOpenCodeConfigIdentity++;
openCodeConfigIdentities.set(configIdentity, identity);
}
const environment = SAFE_ENV_KEYS.map((key) => `${key}=${process.env[key] ?? ""}`).join("\0");
return `${String(identity)}\0${environment}\0${query}`;
}
type OpenCodeSessionPage = {
sessions: SessionCatalogSession[];
nextCursor?: string;
};
export const isExactOpenCodeSessionCursor = sessionCatalogPaging.isExactCursor;
const OPENCODE_PARAMETER_MESSAGES = {
listNotObject: "OpenCode session list parameters must be an object",
unknownListParameter: (key: string) => `unknown OpenCode session list parameter: ${key}`,
invalidSearchTerm: "searchTerm is invalid",
readNotObject: "OpenCode session read parameters must be an object",
unknownReadParameter: (key: string) => `unknown OpenCode session read parameter: ${key}`,
invalidThreadId: "threadId is invalid",
};
async function runOpenCode(args: string[], timeoutMs = CLI_TIMEOUT_MS): Promise<string> {
const configDirectory = args.includes("--standalone")
? await mkdtemp(path.join(resolvePreferredOpenClawTmpDir(), "opencode-config-"))
: undefined;
try {
if (timeoutMs <= 0) {
throw new Error("OpenCode session scan exceeded the time limit");
}
return await executeOpenCode(args, timeoutMs, configDirectory);
} catch (error) {
log.warn(`OpenCode catalog CLI failed: ${String(error)}`);
throw error;
} finally {
if (configDirectory) {
await rm(configDirectory, { recursive: true, force: true });
}
}
}
async function executeOpenCode(
args: string[],
timeoutMs: number,
configDirectory?: string,
): Promise<string> {
const invocation = materializeWindowsSpawnProgram(
resolveWindowsSpawnProgram({
command: "opencode",
platform: process.platform,
env: process.env,
execPath: process.execPath,
packageName: "opencode-ai",
}),
args,
);
const env: NodeJS.ProcessEnv = { OPENCODE_PURE: "1", NO_COLOR: "1" };
for (const key of SAFE_ENV_KEYS) {
if (process.env[key] !== undefined) {
env[key] = process.env[key];
}
}
if (configDirectory) {
// v2 removed --pure. A private server and empty config root avoid executing
// user plugins; project configuration must also be disabled when exporting.
env.OPENCODE_CONFIG_DIR = configDirectory;
env.OPENCODE_CONFIG_PROJECT_DISABLE = "1";
}
const result = await runCommandBuffered([invocation.command, ...invocation.argv], {
baseEnv: {},
env,
input: "",
maxCombinedOutputBytes: MAX_CLI_OUTPUT_BYTES,
maxOutputBytes: MAX_CLI_OUTPUT_BYTES,
terminateOnOutputError: true,
timeoutMs,
});
if (result.termination === "output-limit") {
throw new Error("OpenCode session output exceeded the safety limit");
}
if (result.errorStream) {
const message = result.error?.message ?? "unknown error";
throw new Error(`OpenCode ${result.errorStream} stream failed: ${message}`, {
cause: result.error,
});
}
if (result.termination === "error" && result.error) {
throw result.error;
}
if (result.code !== 0) {
const detail = result.stderr.toString("utf8").trim();
throw new Error(detail || `OpenCode exited with code ${String(result.code)}`, {
cause: result.stdout.toString("utf8"),
});
}
return result.stdout.toString("utf8");
}
export async function queryOpenCodeDatabase(query: string): Promise<unknown> {
const output = await runOpenCode(["--pure", "db", query, "--format", "json"]);
return output.trim() ? (JSON.parse(output) as unknown) : [];
}
export async function runOpenCodeApi(
operation: string,
params: string[],
timeoutMs = CLI_TIMEOUT_MS,
): Promise<string> {
return runOpenCode(
["api", "--standalone", operation, ...params.flatMap((param) => ["--param", param])],
timeoutMs,
);
}
async function queryOpenCodeSessions(query: string, count: number): Promise<unknown> {
if (!(await isOpenCodeV2())) {
return queryOpenCodeDatabase(query);
}
const sessions: unknown[] = [];
let cursor: string | undefined;
const deadline = Date.now() + CLI_TIMEOUT_MS;
while (sessions.length < count) {
const limit = Math.max(OPENCODE_SESSION_CATALOG_MAX_PAGE_LIMIT, count - sessions.length);
const parsed: unknown = JSON.parse(
await runOpenCodeApi(
"session.list",
[`limit=${limit}`, ...(cursor ? [`cursor=${cursor}`] : ["order=desc", "parentID=null"])],
deadline - Date.now(),
),
);
if (!isRecord(parsed) || !Array.isArray(parsed.data) || parsed.data.length > limit) {
throw new Error("OpenCode returned an invalid session list");
}
for (const session of parsed.data) {
if (!isRecord(session) || !isRecord(session.time) || !isRecord(session.location)) {
throw new Error("OpenCode returned an invalid session");
}
if (session.time.archived === undefined) {
sessions.push({
id: session.id,
title: session.title,
directory: session.location.directory,
created: session.time.created,
updated: session.time.updated,
});
if (sessions.length === count) {
break;
}
}
}
// The API applies its limit before archived rows are removed.
if (parsed.data.length < limit || sessions.length >= count) {
break;
}
const next = isRecord(parsed.cursor) ? parsed.cursor.next : undefined;
if (typeof next !== "string" || !next || next === cursor) {
throw new Error("OpenCode returned an invalid session cursor");
}
cursor = next;
}
return sessions;
}
async function queryCachedOpenCodeSessions(
query: string,
count: number,
options: OpenCodeQueryCacheOptions,
): Promise<unknown> {
const key = openCodeQueryCacheKey(query, options.configIdentity ?? process.env);
const cached = openCodeQueryCache.get(key);
if (options.forceRefresh !== true && cached && cached.expiresAt > Date.now()) {
openCodeQueryCache.delete(key);
openCodeQueryCache.set(key, cached);
return await cached.result;
}
if (cached) {
openCodeQueryCache.delete(key);
}
const result = queryOpenCodeSessions(query, count);
const entry: OpenCodeQueryCacheEntry = {
expiresAt: Date.now() + OPENCODE_QUERY_CACHE_TTL_MS,
result,
};
openCodeQueryCache.set(key, entry);
pruneMapToMaxSize(openCodeQueryCache, OPENCODE_QUERY_CACHE_MAX_ENTRIES);
try {
const value = await result;
entry.resolved = true;
return value;
} catch (error) {
if (openCodeQueryCache.get(key) === entry) {
if (cached?.resolved) {
openCodeQueryCache.set(key, cached);
} else {
openCodeQueryCache.delete(key);
}
}
throw error;
}
}
export async function exportOpenCodeSession(threadId: string): Promise<unknown> {
if (!(await isOpenCodeV2())) {
return JSON.parse(await runOpenCode(["--pure", "export", threadId])) as unknown;
}
const exported: unknown = JSON.parse(
await runOpenCode(["session", "export", "--standalone", threadId]),
);
if (!isRecord(exported) || !Array.isArray(exported.messages)) {
throw new Error("OpenCode returned an invalid session export");
}
return {
...exported,
messages: exported.messages.flatMap((message) => {
if (!isRecord(message)) {
return [];
}
const model = isRecord(message.model) ? message.model : {};
const parts =
message.type === "assistant" && Array.isArray(message.content)
? message.content.map((part) => {
if (!isRecord(part) || part.type !== "tool" || !isRecord(part.state)) {
return part;
}
return {
...part,
tool: part.name,
state: {
...part.state,
output: Array.isArray(part.state.content)
? part.state.content
.flatMap((content) =>
isRecord(content) &&
content.type === "text" &&
typeof content.text === "string"
? [content.text]
: [],
)
.join("\n")
: undefined,
error: isRecord(part.state.error) ? jsonText(part.state.error) : part.state.error,
},
};
})
: typeof message.text === "string"
? [{ type: "text", text: message.text, synthetic: message.type !== "user" }]
: [];
if (message.type === "user" && Array.isArray(message.files)) {
parts.push(
...message.files.flatMap((file) => (isRecord(file) ? [{ ...file, type: "file" }] : [])),
);
}
return [
{
info: {
id: message.id,
role: message.type === "user" ? "user" : "assistant",
time: message.time,
model,
...model,
},
parts,
},
];
}),
};
}
function parseOpenCodeSession(value: unknown): SessionCatalogSession | undefined {
if (!isRecord(value)) {
return undefined;
}
const threadId = optionalOpenCodeString(value.id, 256);
if (!threadId || !OPENCODE_SESSION_ID_PATTERN.test(threadId)) {
return undefined;
}
const name = optionalOpenCodeString(value.title, 1_000);
const cwd = optionalOpenCodeString(value.directory, 4_096);
const createdAt =
typeof value.created === "number" && Number.isFinite(value.created) ? value.created : undefined;
const updatedAt =
typeof value.updated === "number" && Number.isFinite(value.updated) ? value.updated : undefined;
return {
threadId,
...(name ? { name } : {}),
...(cwd ? { cwd } : {}),
status: "stored",
...(createdAt !== undefined ? { createdAt } : {}),
...(updatedAt !== undefined ? { updatedAt, recencyAt: updatedAt } : {}),
source: "opencode-cli",
modelProvider: "opencode",
archived: false,
canContinue: true,
canArchive: false,
};
}
export async function listLocalOpenCodeSessionPage(
value?: unknown,
options: OpenCodeQueryCacheOptions = {},
): Promise<OpenCodeSessionPage> {
const params = sessionCatalogPaging.parseListParams(value, {
searchMaxLength: MAX_SEARCH_LENGTH,
messages: OPENCODE_PARAMETER_MESSAGES,
});
const offset = sessionCatalogPaging.decodeCursor(params.cursor);
const requestedCount = params.searchTerm
? MAX_CLI_LIST_SESSIONS
: Math.min(MAX_CLI_LIST_SESSIONS, offset + params.limit + 1);
const query = [
"SELECT id, title, time_created AS created, time_updated AS updated,",
"project_id AS projectId, directory FROM session",
"WHERE parent_id IS NULL AND time_archived IS NULL",
`ORDER BY time_updated DESC, id DESC LIMIT ${String(requestedCount)}`,
].join(" ");
const parsed = await queryCachedOpenCodeSessions(query, requestedCount, options);
if (!Array.isArray(parsed) || parsed.length > MAX_CLI_LIST_SESSIONS) {
throw new Error("OpenCode returned an invalid session list");
}
const needle = params.searchTerm?.toLocaleLowerCase();
const sessions = parsed
.flatMap((entry) => {
const session = parseOpenCodeSession(entry);
return session ? [session] : [];
})
.filter((session) => {
if (!needle) {
return true;
}
return [session.threadId, session.name, session.cwd].some((field) =>
field?.toLocaleLowerCase().includes(needle),
);
});
const page = sessions.slice(offset, offset + params.limit);
return {
sessions: page,
...(offset + page.length < sessions.length
? { nextCursor: sessionCatalogPaging.encodeCursor(offset + page.length) }
: {}),
};
}
export async function requireLocalOpenCodeSession(
threadId: string,
): Promise<SessionCatalogSession> {
const page = await listLocalOpenCodeSessionPage({
searchTerm: threadId,
limit: OPENCODE_SESSION_CATALOG_MAX_PAGE_LIMIT,
});
const session = page.sessions.find((candidate) => candidate.threadId === threadId);
if (!session) {
throw new Error("OpenCode session is unavailable");
}
return session;
}
function jsonText(value: unknown, maxLength = 20_000): string | undefined {
try {
const text = JSON.stringify(value);
return text.length > maxLength ? `${truncateUtf16Safe(text, maxLength)}…` : text;
} catch {
return undefined;
}
}
function timestampFromInfo(info: Record<string, unknown>): string | undefined {
if (!isRecord(info.time) || typeof info.time.created !== "number") {
return undefined;
}
const date = new Date(info.time.created);
return Number.isNaN(date.getTime()) ? undefined : date.toISOString();
}
function openCodeTranscriptItems(value: unknown): SessionCatalogTranscriptItem[] {
if (!isRecord(value) || !Array.isArray(value.messages)) {
throw new Error("OpenCode returned an invalid session export");
}
return value.messages.flatMap((message): SessionCatalogTranscriptItem[] => {
if (!isRecord(message) || !isRecord(message.info) || !Array.isArray(message.parts)) {
return [];
}
const info = message.info;
const role = info.role;
const messageId = optionalOpenCodeString(info.id, 256);
const timestamp = timestampFromInfo(info);
const modelId =
role === "assistant"
? optionalOpenCodeString(info.modelID, 256)
: isRecord(info.model)
? optionalOpenCodeString(info.model.modelID, 256)
: undefined;
const providerId =
role === "assistant"
? optionalOpenCodeString(info.providerID, 256)
: isRecord(info.model)
? optionalOpenCodeString(info.model.providerID, 256)
: undefined;
const model = providerId && modelId ? `${providerId}/${modelId}` : modelId;
return message.parts.flatMap((part, partIndex): SessionCatalogTranscriptItem[] => {
if (!isRecord(part)) {
return [];
}
const id =
optionalOpenCodeString(part.id, 256) ??
(messageId ? `${messageId}:${String(partIndex)}` : undefined);
const common = {
...(id ? { id } : {}),
...(timestamp ? { timestamp } : {}),
...(model ? { model } : {}),
};
if (part.type === "text" && typeof part.text === "string") {
return [
{ ...common, type: role === "user" ? "userMessage" : "agentMessage", text: part.text },
];
}
if (part.type === "reasoning" && typeof part.text === "string") {
return [{ ...common, type: "reasoning", text: part.text }];
}
if (part.type === "tool") {
const tool = optionalOpenCodeString(part.tool, 256) ?? "tool";
const state = isRecord(part.state) ? part.state : undefined;
const callText = state && "input" in state ? jsonText(state.input) : undefined;
const resultText =
state?.status === "completed" && typeof state.output === "string"
? state.output
: state?.status === "error" && typeof state.error === "string"
? state.error
: undefined;
return [
{ ...common, type: "toolCall", text: callText ? `${tool}\n${callText}` : tool },
...(resultText
? [
{
...common,
...(id ? { id: `${id}:result` } : {}),
type: "toolResult" as const,
text: resultText,
},
]
: []),
];
}
if (part.type === "file") {
const filename = optionalOpenCodeString(part.filename, 1_000);
const mime = optionalOpenCodeString(part.mime, 256);
return [
{
...common,
type: "other",
text: `[Attachment${filename ? `: ${filename}` : ""}${mime ? ` (${mime})` : ""}]`,
},
];
}
return [];
});
});
}
export async function readLocalOpenCodeTranscriptPage(
value: unknown,
): Promise<SessionsCatalogReadResult> {
const params = sessionCatalogPaging.parseReadParams(value, {
threadIdMaxLength: 256,
threadIdPattern: OPENCODE_SESSION_ID_PATTERN,
messages: OPENCODE_PARAMETER_MESSAGES,
});
const offset = sessionCatalogPaging.decodeCursor(params.cursor);
const items = openCodeTranscriptItems(await exportOpenCodeSession(params.threadId));
const page = sessionCatalogPaging.boundTranscriptPage(items, params.limit, offset);
return {
hostId: LOCAL_HOST_ID,
label: "Local OpenCode",
threadId: params.threadId,
...page,
};
}