From df91434936b512716b72d34109fa5e58b2af1140 Mon Sep 17 00:00:00 2001 From: oliver Date: Fri, 1 May 2026 17:40:52 +0800 Subject: [PATCH] =?UTF-8?q?=E9=80=82=E9=85=8D=20Bailian=20WebParser=20?= =?UTF-8?q?=E5=85=BC=E5=AE=B9=E6=8E=A5=E5=85=A5=E5=B9=B6=E5=AE=8C=E5=96=84?= =?UTF-8?q?=20sidecar=20=E6=96=87=E6=A1=A3=E3=80=82?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 新增 WebParser 兼容工具与 MCP health/tools-sync 兼容模式,让 WebParser/sse 可注册到 agent 并暴露明确 schema;同时修复 mcp argv 环境变量展开并补齐微信 sidecar 纯本地化安装与文档说明。 Made-with: Cursor --- docs/LOCAL_PUBLIC_TOOLS.md | 1 + docs/OPEN_SOURCE_QUICKSTART.md | 2 + docs/RUNBOOK.md | 8 +- interfaces/admin/routes.py | 45 +++ runtime/operations/scripts/weixin_install.ps1 | 31 +- runtime/operations/scripts/weixin_login.ps1 | 2 + runtime/operations/weixin_bridge/login.ts | 84 +++++- .../weixin_bridge/official_runner.ts | 34 ++- runtime/tools/mcp/adapter.py | 29 +- runtime/tools/mcp/filesystem_argv.py | 12 +- .../tools/public/bailian_webparser_tool.py | 281 ++++++++++++++++++ tests/test_bailian_webparser_tool.py | 81 +++++ tests/test_local_public_tools.py | 15 + tests/test_mcp_adapter.py | 32 ++ tests/test_mcp_admin_api.py | 27 ++ tests/test_mcp_filesystem_argv.py | 12 + 16 files changed, 640 insertions(+), 56 deletions(-) create mode 100644 runtime/tools/public/bailian_webparser_tool.py create mode 100644 tests/test_bailian_webparser_tool.py diff --git a/docs/LOCAL_PUBLIC_TOOLS.md b/docs/LOCAL_PUBLIC_TOOLS.md index b49e117a..666d7d98 100644 --- a/docs/LOCAL_PUBLIC_TOOLS.md +++ b/docs/LOCAL_PUBLIC_TOOLS.md @@ -22,6 +22,7 @@ This project exposes local atomic capabilities as shared `public` tools for all - `set_env` - `list_processes` - `kill_process` +- `bailian_webparser` (DashScope WebParser compatibility tool) ## Loading path diff --git a/docs/OPEN_SOURCE_QUICKSTART.md b/docs/OPEN_SOURCE_QUICKSTART.md index dd97834a..e51cc6d3 100644 --- a/docs/OPEN_SOURCE_QUICKSTART.md +++ b/docs/OPEN_SOURCE_QUICKSTART.md @@ -44,6 +44,8 @@ Notes: - No global `openclaw` CLI installation is required. - Runtime dependencies are installed locally in sidecar workspace. +- Login/account state is persisted under `data/channel_sidecar/oclaw-weixin/state/`. +- The Weixin flow does not rely on `%USERPROFILE%\.openclaw`. ## Path C: Enable WhatsApp (experimental) diff --git a/docs/RUNBOOK.md b/docs/RUNBOOK.md index c1f0137e..323d3e0a 100644 --- a/docs/RUNBOOK.md +++ b/docs/RUNBOOK.md @@ -75,7 +75,7 @@ powershell -ExecutionPolicy Bypass -File .\scripts\stop_gateway.ps1 -Force ### 1.1.4 安装官方微信插件(Personal WeChat) -> 这一步会安装/更新官方插件到:`%USERPROFILE%\.openclaw\extensions\openclaw-weixin\` +> 这一步会安装/更新微信 sidecar 运行时到:`data/channel_sidecar/oclaw-weixin/` ```powershell powershell -ExecutionPolicy Bypass -File .\runtime\operations\scripts\weixin_install.ps1 @@ -87,10 +87,10 @@ powershell -ExecutionPolicy Bypass -File .\runtime\operations\scripts\weixin_ins powershell -ExecutionPolicy Bypass -File .\runtime\operations\scripts\weixin_login.ps1 ``` -按提示扫码完成绑定(会写入账号 ID / token 等状态到 `%USERPROFILE%\.openclaw\openclaw-weixin\`)。 +按提示扫码完成绑定(会写入账号 ID / token 等状态到 `data/channel_sidecar/oclaw-weixin/state/`)。 > 说明:`weixin_install.ps1` **不要求你全局安装 openclaw**(不需要 `npm install -g openclaw`)。 -> 脚本会在 `data/channel_sidecar/oclaw-weixin/` 下安装本地运行时依赖,并完成官方插件安装。 +> 脚本会在 `data/channel_sidecar/oclaw-weixin/` 下安装本地运行时依赖与官方插件包(`node_modules`),全流程不依赖 `%USERPROFILE%\.openclaw`。 ### 1.1.6 启动微信 sidecar(原生模式) @@ -548,6 +548,7 @@ powershell -ExecutionPolicy Bypass -File .\scripts\weixin_install.ps1 ``` 注意:当前仅保留官方插件单链路,直接执行 `weixin_install.ps1` 即可。 +且为纯 sidecar 方案:不会写入 `%USERPROFILE%\.openclaw\`。 安装目录: @@ -563,6 +564,7 @@ powershell -ExecutionPolicy Bypass -File .\scripts\weixin_login.ps1 - `data/channel_sidecar/oclaw-weixin/state/oclaw-weixin/accounts/*.json` - `data/channel_sidecar/oclaw-weixin/state/oclaw-weixin/accounts.json` +- 不会写入 `%USERPROFILE%\.openclaw\openclaw-weixin\` ### 11.3 启动微信 sidecar diff --git a/interfaces/admin/routes.py b/interfaces/admin/routes.py index 441e7473..94b52514 100644 --- a/interfaces/admin/routes.py +++ b/interfaces/admin/routes.py @@ -120,6 +120,12 @@ def _mcp_health_and_sync_one(store: SqliteStore, row: dict[str, Any]) -> dict[st } store.set_mcp_server_health(server_id=sid, status="error", detail=item["health"]) return item + if _is_bailian_webparser_remote(entry_command=cmd, entry_args=args): + tools = _bailian_webparser_virtual_tools() + store.replace_mcp_server_tools(server_id=sid, tools=tools) + detail = {"synced_tools": len(tools), "compat_mode": "bailian_webparser"} + store.set_mcp_server_health(server_id=sid, status="ok", detail=detail) + return {"server_id": sid, "ok": True, "health": detail, "tools_synced": len(tools)} rt = McpProcessRuntime( build_mcp_process_command(cmd, args, store=store), timeout_s=float(row.get("timeout_s") or 30.0), @@ -158,6 +164,32 @@ def _mcp_health_and_sync_one(store: SqliteStore, row: dict[str, Any]) -> dict[st rt.stop() +def _is_bailian_webparser_remote(*, entry_command: str, entry_args: list[str]) -> bool: + cmd = str(entry_command or "").strip().lower() + if cmd not in {"npx", "npx.cmd", "node"}: + return False + joined = " ".join(str(x or "").strip().lower() for x in (entry_args or [])) + return "mcp-remote" in joined and "/api/v1/mcps/webparser/sse" in joined + + +def _bailian_webparser_virtual_tools() -> list[dict[str, Any]]: + return [ + { + "tool_name": "bailian_webparser_parse", + "description": "Parse webpage via DashScope WebParser compatibility mode. Requires `url` (http/https).", + "parameters": { + "type": "object", + "properties": { + "url": {"type": "string", "description": "Target webpage URL (required). Example: https://example.com"}, + "timeout": {"type": "integer", "default": 35, "minimum": 8, "maximum": 90}, + }, + "required": ["url"], + "additionalProperties": False, + }, + } + ] + + def _http_get_json(url: str, *, timeout: float = 8.0) -> dict[str, Any]: req = urllib_request.Request( url, @@ -3178,6 +3210,10 @@ def build_admin_router() -> APIRouter: args = [str(x) for x in (row.get("entry_args") or []) if str(x).strip()] if not cmd: return {"ok": False, "error_code": "mcp_entry_missing", "error": "entry_command_missing"} + if _is_bailian_webparser_remote(entry_command=cmd, entry_args=args): + detail = {"ok": True, "status": "ok", "compat_mode": "bailian_webparser", "tools_count": 1} + store.set_mcp_server_health(server_id=server_id, status="ok", detail=detail) + return {"ok": True, "response": detail} rt = McpProcessRuntime( build_mcp_process_command(cmd, args, store=store), timeout_s=float(row.get("timeout_s") or 30.0), @@ -3225,6 +3261,15 @@ def build_admin_router() -> APIRouter: args = [str(x) for x in (row.get("entry_args") or []) if str(x).strip()] if not cmd: return {"ok": False, "error_code": "mcp_entry_missing", "error": "entry_command_missing"} + if _is_bailian_webparser_remote(entry_command=cmd, entry_args=args): + norm = _bailian_webparser_virtual_tools() + store.replace_mcp_server_tools(server_id=server_id, tools=norm) + store.set_mcp_server_health( + server_id=server_id, + status="ok", + detail={"synced_tools": len(norm), "compat_mode": "bailian_webparser"}, + ) + return {"ok": True, "server_id": server_id, "tools": norm, "compat_mode": "bailian_webparser"} rt = McpProcessRuntime( build_mcp_process_command(cmd, args, store=store), timeout_s=float(row.get("timeout_s") or 30.0), diff --git a/runtime/operations/scripts/weixin_install.ps1 b/runtime/operations/scripts/weixin_install.ps1 index f1d74023..ebad1039 100644 --- a/runtime/operations/scripts/weixin_install.ps1 +++ b/runtime/operations/scripts/weixin_install.ps1 @@ -13,7 +13,6 @@ function Resolve-RepoRoot { $oclawRoot = Resolve-RepoRoot $sidecarRoot = Join-Path $oclawRoot "data\\channel_sidecar\\$ChannelId" $stateDir = Join-Path $sidecarRoot "state" -$pluginRoot = Join-Path $env:USERPROFILE ".openclaw\\extensions\\openclaw-weixin" function Sync-WeixinBridgeRunners { $bridgeSrc = Join-Path $oclawRoot "runtime\\operations\\weixin_bridge" foreach ($name in @("official_runner.ts", "login.ts")) { @@ -29,41 +28,29 @@ New-Item -ItemType Directory -Force -Path $sidecarRoot | Out-Null New-Item -ItemType Directory -Force -Path (Join-Path $sidecarRoot "logs") | Out-Null New-Item -ItemType Directory -Force -Path $stateDir | Out-Null -function Ensure-OfficialPluginRuntimeDeps { - if (-not (Test-Path (Join-Path $pluginRoot "package.json"))) { - throw "official plugin root not found at $pluginRoot" - } - Push-Location $pluginRoot - try { - npm.cmd install openclaw@latest --no-save - if ($LASTEXITCODE -ne 0) { - throw "npm install official plugin runtime deps failed with exit code $LASTEXITCODE" - } - } finally { - Pop-Location - } -} - Push-Location $sidecarRoot try { + $npmRegistry = ($env:OCLAW_NPM_REGISTRY).Trim() + $npmRegistryArgs = @() + if ($npmRegistry) { + $npmRegistryArgs = @("--registry", $npmRegistry) + } if (-not (Test-Path (Join-Path $sidecarRoot "package.json"))) { npm.cmd init -y | Out-Null if ($LASTEXITCODE -ne 0) { throw "npm init failed with exit code $LASTEXITCODE" } } - npx.cmd -y @tencent-weixin/openclaw-weixin-cli@latest install - if ($LASTEXITCODE -ne 0) { - throw "openclaw-weixin-cli install failed with exit code $LASTEXITCODE" - } - npm.cmd install openclaw@latest --save tsx@4.21.0 typescript@6.0.3 + # Pure sidecar mode: + # - No ~/.openclaw usage + # - Install the official plugin package into node_modules and point it to $stateDir via OPENCLAW_STATE_DIR at runtime + npm.cmd install --no-audit --no-fund --save @tencent-weixin/openclaw-weixin@latest openclaw@latest tsx@4.21.0 typescript@6.0.3 @npmRegistryArgs if ($LASTEXITCODE -ne 0) { throw "npm install bridge runtime deps failed with exit code $LASTEXITCODE" } } finally { Pop-Location } -Ensure-OfficialPluginRuntimeDeps Sync-WeixinBridgeRunners Write-Host "[ok] installed official openclaw-weixin plugin sidecar runtime into $sidecarRoot" diff --git a/runtime/operations/scripts/weixin_login.ps1 b/runtime/operations/scripts/weixin_login.ps1 index 4e25e3c5..07c22591 100644 --- a/runtime/operations/scripts/weixin_login.ps1 +++ b/runtime/operations/scripts/weixin_login.ps1 @@ -23,6 +23,8 @@ New-Item -ItemType Directory -Force -Path $stateDir | Out-Null Push-Location $sidecarRoot try { $env:OCLAW_STATE_DIR = $stateDir + # Ensure the official plugin stores state under the sidecar. + $env:OPENCLAW_STATE_DIR = $stateDir if (Test-Path (Join-Path $sidecarRoot "login.ts")) { npm.cmd exec -- tsx login.ts exit 0 diff --git a/runtime/operations/weixin_bridge/login.ts b/runtime/operations/weixin_bridge/login.ts index bb343959..d3c87879 100644 --- a/runtime/operations/weixin_bridge/login.ts +++ b/runtime/operations/weixin_bridge/login.ts @@ -1,24 +1,76 @@ -import { spawn } from "node:child_process"; import path from "node:path"; +import { pathToFileURL } from "node:url"; -function resolveLocalOpenclawBin(): string { - // Prefer local openclaw installed in the sidecar runtime (no global CLI required). - // Windows: node_modules/.bin/openclaw.cmd - return path.join(process.cwd(), "node_modules", ".bin", process.platform === "win32" ? "openclaw.cmd" : "openclaw"); +const STATE_DIR = (process.env.OCLAW_STATE_DIR || path.resolve(process.cwd(), "state")).trim(); + +function resolvePluginRoot(): string { + const configured = String(process.env.OCLAW_WEIXIN_PLUGIN_ROOT || "").trim(); + return configured || path.join(process.cwd(), "node_modules", "@tencent-weixin", "openclaw-weixin"); } -function run(): Promise { - return new Promise((resolve, reject) => { - const openclawBin = resolveLocalOpenclawBin(); - const child = spawn(openclawBin, ["channels", "login", "--channel", "openclaw-weixin"], { - stdio: "inherit", - shell: true, - }); - child.on("error", reject); - child.on("exit", (code) => resolve(code ?? 1)); +async function importFromPluginSrc(relativePath: string): Promise { + const pluginRoot = resolvePluginRoot(); + const srcRoot = path.join(pluginRoot, "src"); + const fullPath = path.join(srcRoot, relativePath); + return import(pathToFileURL(fullPath).href); +} + +async function run(): Promise { + // Force the plugin's state-dir resolver away from ~/.openclaw + process.env.OPENCLAW_STATE_DIR = STATE_DIR; + + const loginQr = await importFromPluginSrc(path.join("auth", "login-qr.ts")); + const accounts = await importFromPluginSrc(path.join("auth", "accounts.ts")); + + const startWeixinLoginWithQr = loginQr.startWeixinLoginWithQr as ((opts: any) => Promise) | undefined; + const waitForWeixinLogin = loginQr.waitForWeixinLogin as ((opts: any) => Promise) | undefined; + const displayQRCode = loginQr.displayQRCode as ((qrcodeUrl: string) => Promise) | undefined; + const registerWeixinAccountId = accounts.registerWeixinAccountId as ((accountId: string) => void) | undefined; + const saveWeixinAccount = accounts.saveWeixinAccount as ((accountId: string, update: any) => void) | undefined; + + if (!startWeixinLoginWithQr || !waitForWeixinLogin || !displayQRCode) { + throw new Error("weixin plugin login-qr module missing exports"); + } + if (!registerWeixinAccountId || !saveWeixinAccount) { + throw new Error("weixin plugin accounts module missing exports"); + } + + const apiBaseUrl = "https://ilinkai.weixin.qq.com"; + const botType = process.env.OCLAW_WEIXIN_BOT_TYPE?.trim(); + + const start = await startWeixinLoginWithQr({ + apiBaseUrl, + botType, }); + if (!start.qrcodeUrl) { + throw new Error(start.message || "failed to start login"); + } + + process.stdout.write(`${start.message}\n`); + await displayQRCode(start.qrcodeUrl); + + const result = await waitForWeixinLogin({ + sessionKey: start.sessionKey, + apiBaseUrl, + botType, + }); + + process.stdout.write(`${result.message}\n`); + if (!result.connected || !result.accountId || !result.botToken) { + process.exitCode = 1; + return; + } + + registerWeixinAccountId(result.accountId); + saveWeixinAccount(result.accountId, { + token: result.botToken, + baseUrl: result.baseUrl, + userId: result.userId, + }); + process.stdout.write(`saved account=${result.accountId} into ${STATE_DIR}\n`); } -void run().then((code) => { - process.exitCode = code; +void run().catch((err) => { + process.stderr.write(`weixin login failed: ${String(err)}\n`); + process.exitCode = 1; }); diff --git a/runtime/operations/weixin_bridge/official_runner.ts b/runtime/operations/weixin_bridge/official_runner.ts index aa956d57..5736e7c0 100644 --- a/runtime/operations/weixin_bridge/official_runner.ts +++ b/runtime/operations/weixin_bridge/official_runner.ts @@ -1,5 +1,4 @@ import fs from "node:fs"; -import os from "node:os"; import path from "node:path"; import { pathToFileURL } from "node:url"; @@ -89,13 +88,9 @@ function sleep(ms: number): Promise { return new Promise((resolve) => setTimeout(resolve, ms)); } -function homeOpenclawPath(...parts: string[]): string { - return path.join(os.homedir(), ".openclaw", ...parts); -} - function resolvePluginRoot(): string { const configured = String(process.env.OCLAW_WEIXIN_PLUGIN_ROOT || "").trim(); - return configured || homeOpenclawPath("extensions", "openclaw-weixin"); + return configured || path.join(process.cwd(), "node_modules", "@tencent-weixin", "openclaw-weixin"); } async function loadOfficialModules(): Promise { @@ -130,19 +125,34 @@ async function loadOfficialModules(): Promise { return officialModulesPromise; } -function resolveAccount(): { accountId: string; token: string; cloudBaseUrl: string } { - const ids = readJsonFile(homeOpenclawPath("openclaw-weixin", "accounts.json")) || []; +async function resolveAccount(): Promise<{ accountId: string; token: string; cloudBaseUrl: string }> { + if (!String(process.env.OPENCLAW_STATE_DIR || "").trim()) { + process.env.OPENCLAW_STATE_DIR = STATE_DIR; + } + const pluginRoot = resolvePluginRoot(); + const srcRoot = path.join(pluginRoot, "src"); + const importTs = async (relativePath: string): Promise => { + const fullPath = path.join(srcRoot, relativePath); + return import(pathToFileURL(fullPath).href); + }; + const accountsMod = await importTs(path.join("auth", "accounts.ts")); + const listIds = accountsMod.listIndexedWeixinAccountIds as (() => string[]) | undefined; + const loadAcc = accountsMod.loadWeixinAccount as ((accountId: string) => any) | undefined; + if (!listIds || !loadAcc) { + throw new Error("weixin plugin accounts module missing exports"); + } + const ids = listIds() || []; const accountId = String(ids[0] || "").trim(); if (!accountId) { throw new Error("no weixin account id found; run login first"); } - const account = readJsonFile(homeOpenclawPath("openclaw-weixin", "accounts", `${accountId}.json`)) || {}; - const token = String(account.token || "").trim(); + const data = loadAcc(accountId) || {}; + const token = String(data.token || "").trim(); if (!token) { throw new Error(`missing token for account ${accountId}; run login again`); } const envCloud = String(process.env.OCLAW_WEIXIN_CLOUD_BASE_URL || "").trim(); - const cfgCloud = String(account.baseUrl || "").trim(); + const cfgCloud = String(data.baseUrl || "").trim(); const cloudBaseUrl = (envCloud || cfgCloud || "https://ilinkai.weixin.qq.com").trim(); return { accountId, token, cloudBaseUrl }; } @@ -424,7 +434,7 @@ async function main(): Promise { state.user_context_tokens && typeof state.user_context_tokens === "object" ? (state.user_context_tokens as TokenMap) : {}; - const { accountId, token, cloudBaseUrl } = resolveAccount(); + const { accountId, token, cloudBaseUrl } = await resolveAccount(); const modules = await loadOfficialModules(); modules.restoreContextTokens(accountId); log(`official runner started account=${accountId} cloud=${cloudBaseUrl} local=${LOCAL_BASE_URL}`); diff --git a/runtime/tools/mcp/adapter.py b/runtime/tools/mcp/adapter.py index 69b27907..0b66377e 100644 --- a/runtime/tools/mcp/adapter.py +++ b/runtime/tools/mcp/adapter.py @@ -10,6 +10,7 @@ from oclaw.runtime.skills import SkillSpec, materialize_skills_from_tool_specs from oclaw.runtime.tools.base import ToolSpec from oclaw.runtime.tools.mcp.filesystem_argv import build_mcp_process_command from oclaw.runtime.tools.mcp.runtime import McpProcessRuntime +from oclaw.runtime.tools.public.bailian_webparser_tool import bailian_webparser_tool @dataclass @@ -64,6 +65,14 @@ def materialize_mcp_tools_for_specialist( path_policy_tenant_id: str | None = None, path_policy_user_id: str | None = None, ) -> list[ToolSpec]: + def _is_bailian_webparser_remote_row(r: dict[str, Any]) -> bool: + cmd2 = str(r.get("entry_command") or "").strip().lower() + if cmd2 not in {"npx", "npx.cmd", "node"}: + return False + argv = [str(x or "").strip().lower() for x in (r.get("entry_args") or [])] + joined = " ".join(argv) + return "mcp-remote" in joined and "/api/v1/mcps/webparser/sse" in joined + sp = str(specialist or "").strip().lower() if sp == "manager": # Manager is a first-class binding role in admin UI/config. @@ -125,9 +134,27 @@ def materialize_mcp_tools_for_specialist( except Exception: tools = [] for t in tools: + tname = str(t.get("tool_name") or "") + if _is_bailian_webparser_remote_row(row) and tname == "bailian_webparser_parse": + compat = bailian_webparser_tool() + out.append( + ToolSpec( + name=f"mcp__{server_id}__{tname}", + description=str(t.get("description") or compat.description), + parameters=t.get("parameters") if isinstance(t.get("parameters"), dict) else compat.parameters, + handler=compat.handler, + tags=frozenset({"mcp", "plugin", "compat"}), + version="v1", + risk_level="high", + timeout_s=float(row.get("timeout_s") or 30.0), + required_permissions=frozenset(str(x) for x in (row.get("required_permissions") or [])), + execution_mode="subprocess", + ) + ) + continue spec = _McpBoundTool( server_id=server_id, - tool_name=str(t.get("tool_name") or ""), + tool_name=tname, description=str(t.get("description") or f"MCP tool {t.get('tool_name') or ''}"), parameters=t.get("parameters") if isinstance(t.get("parameters"), dict) else {}, command=command, diff --git a/runtime/tools/mcp/filesystem_argv.py b/runtime/tools/mcp/filesystem_argv.py index a4a7bce6..a6bead5c 100644 --- a/runtime/tools/mcp/filesystem_argv.py +++ b/runtime/tools/mcp/filesystem_argv.py @@ -212,9 +212,17 @@ def build_mcp_process_command( path_policy_tenant_id: str | None = None, path_policy_user_id: str | None = None, ) -> list[str]: - """``[cmd] + args`` after filesystem argv augmentation.""" + """``[cmd] + args`` after env placeholder expansion + filesystem argv augmentation.""" + base = [str(cmd)] + for x in args: + s = str(x).strip() + if not s: + continue + # Allow mcp entry args like: "Authorization: Bearer ${DASHSCOPE_API_KEY}". + # Unknown vars are kept unchanged by expandvars. + base.append(os.path.expandvars(s)) return augment_filesystem_mcp_argv( - [cmd] + [x for x in args if str(x).strip()], + base, store=store, policy_session_id=policy_session_id, path_policy_tenant_id=path_policy_tenant_id, diff --git a/runtime/tools/public/bailian_webparser_tool.py b/runtime/tools/public/bailian_webparser_tool.py new file mode 100644 index 00000000..85848c42 --- /dev/null +++ b/runtime/tools/public/bailian_webparser_tool.py @@ -0,0 +1,281 @@ +from __future__ import annotations + +import json +import os +import time +import urllib.parse +import urllib.request +from typing import Any + +from oclaw.runtime.tools.base import ToolSpec + +_WEBPARSER_SSE_URL = "https://dashscope.aliyuncs.com/api/v1/mcps/WebParser/sse" +_UA = ( + "Mozilla/5.0 (Windows NT 10.0; Win64; x64) " + "AppleWebKit/537.36 (KHTML, like Gecko) Chrome/124.0 Safari/537.36" +) + + +def _open_request( + url: str, + *, + timeout_s: int, + headers: dict[str, str], + body: dict[str, Any] | None = None, +): + data = None + method = "GET" + if body is not None: + data = json.dumps(body, ensure_ascii=False).encode("utf-8") + method = "POST" + req = urllib.request.Request(str(url), data=data, headers=headers, method=method) + return urllib.request.urlopen(req, timeout=max(5, min(int(timeout_s), 90))) + + +def _read_sse_until( + sse_resp, + *, + timeout_s: int, + message_url: str | None = None, + wait_id: int | None = None, + auth_header: str = "", +) -> dict[str, Any] | None: + deadline = time.time() + max(3, int(timeout_s)) + event_name = "" + while time.time() < deadline: + raw = sse_resp.readline() + if not raw: + continue + line = raw.decode("utf-8", errors="replace").strip() + if not line: + continue + if line.startswith("event:"): + event_name = line[6:].strip() + continue + if not line.startswith("data:"): + continue + data = line[5:].strip() + if event_name == "endpoint" and wait_id is None: + return {"endpoint": data} + try: + obj = json.loads(data) + except Exception: + continue + if not isinstance(obj, dict): + continue + # Keep-alive ping from server; respond with empty result. + if ( + message_url + and isinstance(obj.get("id"), (int, str)) + and str(obj.get("method") or "").strip() == "ping" + ): + try: + _open_request( + message_url, + timeout_s=10, + headers={ + "Authorization": auth_header, + "Content-Type": "application/json", + "User-Agent": _UA, + }, + body={"jsonrpc": "2.0", "id": obj.get("id"), "result": {}}, + ).close() + except Exception: + pass + continue + if wait_id is not None and obj.get("id") == wait_id: + return obj + return None + + +def _call_webparser(url: str, *, api_key: str, timeout_s: int) -> dict[str, Any]: + auth_header = f"Bearer {api_key}" + sse_headers = { + "Authorization": auth_header, + "Accept": "text/event-stream", + "User-Agent": _UA, + } + with _open_request(_WEBPARSER_SSE_URL, timeout_s=timeout_s, headers=sse_headers) as sse_resp: + endpoint_evt = _read_sse_until(sse_resp, timeout_s=timeout_s) + endpoint = str((endpoint_evt or {}).get("endpoint") or "").strip() + if not endpoint or "/message?" not in endpoint: + return {"ok": False, "error_code": "webparser_no_endpoint", "error": "webparser_no_endpoint"} + message_url = urllib.parse.urljoin(_WEBPARSER_SSE_URL, endpoint) + + def _post(payload: dict[str, Any], post_timeout: int = 30) -> None: + with _open_request( + message_url, + timeout_s=post_timeout, + headers={ + "Authorization": auth_header, + "Content-Type": "application/json", + "User-Agent": _UA, + }, + body=payload, + ): + return + + _post( + { + "jsonrpc": "2.0", + "id": 1, + "method": "initialize", + "params": { + "protocolVersion": "2024-11-05", + "capabilities": {}, + "clientInfo": {"name": "oclaw-bailian-webparser", "version": "0.1.0"}, + }, + } + ) + init_res = _read_sse_until( + sse_resp, + timeout_s=timeout_s, + message_url=message_url, + wait_id=1, + auth_header=auth_header, + ) + if not init_res or not isinstance(init_res.get("result"), dict): + return {"ok": False, "error_code": "webparser_initialize_failed", "error": "initialize_no_result"} + + _post({"jsonrpc": "2.0", "method": "notifications/initialized", "params": {}}, post_timeout=12) + _post({"jsonrpc": "2.0", "id": 2, "method": "tools/list", "params": {}}) + list_res = _read_sse_until( + sse_resp, + timeout_s=timeout_s, + message_url=message_url, + wait_id=2, + auth_header=auth_header, + ) + tools = [] + try: + tools = list((list_res or {}).get("result", {}).get("tools") or []) + except Exception: + tools = [] + if not tools: + return { + "ok": False, + "error_code": "webparser_tools_empty", + "error": "tools_list_empty", + "initialize": init_res, + "tools_list": list_res, + } + tool_name = str((tools[0] or {}).get("name") or "").strip() + if not tool_name: + return {"ok": False, "error_code": "webparser_tool_name_missing", "error": "tool_name_missing"} + + candidate_args = ( + {"url": url}, + {"urls": [url]}, + {"query": url}, + ) + call_res: dict[str, Any] | None = None + used_args: dict[str, Any] | None = None + for args in candidate_args: + _post({"jsonrpc": "2.0", "id": 3, "method": "tools/call", "params": {"name": tool_name, "arguments": args}}) + call_res = _read_sse_until( + sse_resp, + timeout_s=timeout_s, + message_url=message_url, + wait_id=3, + auth_header=auth_header, + ) + if call_res and not isinstance(call_res.get("error"), dict): + used_args = args + break + if not call_res: + return {"ok": False, "error_code": "webparser_call_timeout", "error": "call_timeout"} + if isinstance(call_res.get("error"), dict): + return { + "ok": False, + "error_code": "webparser_call_failed", + "error": str((call_res.get("error") or {}).get("message") or "webparser_call_failed"), + "rpc_error": call_res.get("error"), + "tool_name": tool_name, + "arguments": used_args or {"url": url}, + } + result = call_res.get("result") + text_parts: list[str] = [] + if isinstance(result, dict): + for it in list(result.get("content") or []): + if isinstance(it, dict) and str(it.get("type") or "") == "text": + txt = str(it.get("text") or "").strip() + if txt: + text_parts.append(txt) + return { + "ok": True, + "tool_name": tool_name, + "arguments": used_args or {"url": url}, + "text": "\n\n".join(text_parts).strip(), + "result": result, + } + + +def bailian_webparser_tool() -> ToolSpec: + def _handler(args: dict[str, Any]) -> dict[str, Any]: + payload = args if isinstance(args, dict) else {} + url = str(payload.get("url") or "").strip() + if not url: + return { + "ok": False, + "error_code": "url_required", + "error": 'url_required: provide {"url":"https://example.com/article"}', + } + timeout_s = max(8, min(int(payload.get("timeout") or 35), 90)) + api_key = str(payload.get("api_key") or os.getenv("DASHSCOPE_API_KEY") or "").strip() + if not api_key: + return {"ok": False, "error_code": "api_key_missing", "error": "DASHSCOPE_API_KEY required"} + try: + p = urllib.parse.urlparse(url) + except Exception: + return { + "ok": False, + "error_code": "invalid_url", + "error": 'invalid_url: provide full http/https URL, e.g. {"url":"https://example.com/article"}', + } + if p.scheme not in {"http", "https"}: + return { + "ok": False, + "error_code": "unsupported_scheme", + "error": "unsupported_scheme: url must start with http:// or https://", + } + try: + out = _call_webparser(url, api_key=api_key, timeout_s=timeout_s) + return out + except Exception as exc: + return {"ok": False, "error_code": "webparser_runtime_failed", "error": f"{type(exc).__name__}: {exc}"} + + return ToolSpec( + name="bailian_webparser", + description="Parse webpage content via DashScope WebParser (SSE-compatible fallback path). `url` is required.", + parameters={ + "type": "object", + "properties": { + "url": { + "type": "string", + "description": "Target webpage URL to parse (required). Example: https://example.com/article", + }, + "timeout": { + "type": "integer", + "default": 35, + "minimum": 8, + "maximum": 90, + "description": "Total timeout seconds [8..90].", + }, + "api_key": { + "type": "string", + "description": "Optional DashScope API key override; defaults to DASHSCOPE_API_KEY.", + }, + }, + "required": ["url"], + "additionalProperties": False, + }, + handler=_handler, + tags=frozenset({"public", "web", "parser", "read"}), + risk_level="low", + read_only=True, + timeout_s=95.0, + ) + + +__all__ = ["bailian_webparser_tool"] + diff --git a/tests/test_bailian_webparser_tool.py b/tests/test_bailian_webparser_tool.py new file mode 100644 index 00000000..a00291b2 --- /dev/null +++ b/tests/test_bailian_webparser_tool.py @@ -0,0 +1,81 @@ +from __future__ import annotations + +import json + +from oclaw.runtime.tools.public.bailian_webparser_tool import bailian_webparser_tool + + +class _SseResp: + def __init__(self, lines: list[str]): + self._lines = [x.encode("utf-8") for x in lines] + + def readline(self) -> bytes: + if not self._lines: + return b"" + return self._lines.pop(0) + + def __enter__(self): + return self + + def __exit__(self, exc_type, exc, tb): + return False + + +class _PostResp: + def __enter__(self): + return self + + def __exit__(self, exc_type, exc, tb): + return False + + +def test_bailian_webparser_sse_flow(monkeypatch) -> None: + sse_lines = [ + "event:endpoint\n", + "data:/api/v1/mcps/WebParser/message?sessionId=abc\n", + "\n", + "event:message\n", + 'data:{"jsonrpc":"2.0","id":1,"result":{"protocolVersion":"2024-11-05","capabilities":{"tools":{"listChanged":true}},"serverInfo":{"name":"x","version":"1"}}}\n', + "\n", + "event:message\n", + 'data:{"jsonrpc":"2.0","id":2,"result":{"tools":[{"name":"webparser","description":"d","inputSchema":{"type":"object"}}]}}\n', + "\n", + "event:message\n", + 'data:{"jsonrpc":"2.0","id":3,"result":{"content":[{"type":"text","text":"parsed ok"}]}}\n', + "\n", + ] + + def _fake_open(req, timeout=0): # noqa: ARG001 + method = str(getattr(req, "method", "") or "GET").upper() + url = str(getattr(req, "full_url", "") or "") + if method == "GET" and url.endswith("/WebParser/sse"): + return _SseResp(sse_lines) + if method == "POST" and "/message?sessionId=abc" in url: + # Ensure request is valid JSON-RPC object. + raw = getattr(req, "data", b"") or b"" + _ = json.loads(raw.decode("utf-8")) + return _PostResp() + raise AssertionError(f"unexpected request method={method} url={url}") + + monkeypatch.setenv("DASHSCOPE_API_KEY", "sk-test") + monkeypatch.setattr("urllib.request.urlopen", _fake_open) + + out = bailian_webparser_tool().handler({"url": "https://example.com/x"}) + assert out.get("ok") is True, out + assert out.get("tool_name") == "webparser" + assert "parsed ok" in str(out.get("text") or "") + + +def test_bailian_webparser_requires_key(monkeypatch) -> None: + monkeypatch.delenv("DASHSCOPE_API_KEY", raising=False) + out = bailian_webparser_tool().handler({"url": "https://example.com/x"}) + assert out.get("ok") is False + assert out.get("error_code") == "api_key_missing" + + +def test_bailian_webparser_missing_url_hint() -> None: + out = bailian_webparser_tool().handler({}) + assert out.get("ok") is False + assert out.get("error_code") == "url_required" + assert "https://example.com/article" in str(out.get("error") or "") + diff --git a/tests/test_local_public_tools.py b/tests/test_local_public_tools.py index 0a945771..499ea72e 100644 --- a/tests/test_local_public_tools.py +++ b/tests/test_local_public_tools.py @@ -197,3 +197,18 @@ def test_run_command_does_not_follow_cd_state(tmp_path: Path, monkeypatch) -> No # If run_command follows cd state this would be "subdir"; default should be data/workspace. assert str(out_run.get("cwd") or "").replace("\\", "/").rstrip("/").endswith("/data/workspace") + +def test_run_command_reads_db_toggle_without_ops_db_env(tmp_path: Path, monkeypatch) -> None: + monkeypatch.setenv("OPS_WORKSPACE_ROOT", str(tmp_path)) + monkeypatch.delenv("OPS_ASSISTANT_DB_PATH", raising=False) + monkeypatch.setenv("AIA_ENABLE_RUN_COMMAND", "0") + + db_file = tmp_path / "ops.sqlite" + from oclaw.platform.persistence.sqlite_store import SqliteStore + + SqliteStore(str(db_file)).set_setting("AIA_ENABLE_RUN_COMMAND", "1") + monkeypatch.setattr("oclaw.platform.config.paths.db_path", lambda: str(db_file)) + + out = LocalAdapter().run_command(command="echo hi", timeout=10) + assert out.get("ok") is True, out + diff --git a/tests/test_mcp_adapter.py b/tests/test_mcp_adapter.py index 2572be5a..e45288b0 100644 --- a/tests/test_mcp_adapter.py +++ b/tests/test_mcp_adapter.py @@ -166,10 +166,42 @@ class McpAdapterTests(unittest.TestCase): self.assertIn("GOOGLE_CALENDAR_MCP_TOKEN_PATH", keys) self.assertIn("GITHUB_PERSONAL_ACCESS_TOKEN", keys) self.assertIn("CONTEXT7_API_KEY", keys) + self.assertIn("DASHSCOPE_API_KEY", keys) finally: if old is not None: os.environ["OPS_MCP_ENV_ALLOWLIST"] = old + def test_materialize_bailian_webparser_compat_tool(self) -> None: + with tempfile.TemporaryDirectory(ignore_cleanup_errors=True) as td: + store = SqliteStore(str(Path(td) / "ops.sqlite")) + store.upsert_mcp_server( + server_id="webparser-compat", + source_type="npm", + source_ref="mcp-remote", + entry_command="npx", + entry_args=[ + "-y", + "mcp-remote", + "https://dashscope.aliyuncs.com/api/v1/mcps/WebParser/sse", + "--header", + "Authorization: Bearer ${DASHSCOPE_API_KEY}", + ], + enabled=True, + ) + store.replace_mcp_server_tools( + server_id="webparser-compat", + tools=[ + { + "tool_name": "bailian_webparser_parse", + "description": "compat", + "parameters": {"type": "object", "properties": {"url": {"type": "string"}}}, + } + ], + ) + specs = materialize_mcp_tools(store) + names = {x.name for x in specs} + self.assertIn("mcp__webparser-compat__bailian_webparser_parse", names) + if __name__ == "__main__": unittest.main() diff --git a/tests/test_mcp_admin_api.py b/tests/test_mcp_admin_api.py index 8c411b0b..9b0d87ee 100644 --- a/tests/test_mcp_admin_api.py +++ b/tests/test_mcp_admin_api.py @@ -140,6 +140,33 @@ class McpAdminApiTests(unittest.TestCase): tools = sync.json().get("tools") or [] self.assertTrue(any(str(t.get("tool_name") or "") == "ping" for t in tools)) + def test_healthcheck_and_tools_sync_bailian_webparser_compat(self) -> None: + store = SqliteStore(db_path()) + store.upsert_mcp_server( + server_id="webparser-compat", + source_type="npm", + source_ref="mcp-remote", + entry_command="npx", + entry_args=[ + "-y", + "mcp-remote", + "https://dashscope.aliyuncs.com/api/v1/mcps/WebParser/sse", + "--header", + "Authorization: Bearer ${DASHSCOPE_API_KEY}", + ], + enabled=True, + ) + health = self.client.post("/admin/api/mcp/healthcheck", json={"server_id": "webparser-compat"}, headers=self._headers()) + self.assertEqual(health.status_code, 200) + self.assertTrue(health.json().get("ok"), health.json()) + self.assertEqual(str((health.json().get("response") or {}).get("compat_mode") or ""), "bailian_webparser") + + sync = self.client.post("/admin/api/mcp/tools/sync", json={"server_id": "webparser-compat"}, headers=self._headers()) + self.assertEqual(sync.status_code, 200) + self.assertTrue(sync.json().get("ok"), sync.json()) + tools = sync.json().get("tools") or [] + self.assertTrue(any(str(t.get("tool_name") or "") == "bailian_webparser_parse" for t in tools)) + def test_reinstall_from_saved_manifest(self) -> None: script = self._write_mcp_server() store = SqliteStore(db_path()) diff --git a/tests/test_mcp_filesystem_argv.py b/tests/test_mcp_filesystem_argv.py index 52a19c94..162f0e83 100644 --- a/tests/test_mcp_filesystem_argv.py +++ b/tests/test_mcp_filesystem_argv.py @@ -38,6 +38,18 @@ class McpFilesystemArgvTests(unittest.TestCase): out = build_mcp_process_command("echo", ["a"], store=None) self.assertEqual(out, ["echo", "a"]) + def test_build_mcp_process_command_expands_env_placeholders(self) -> None: + with mock.patch.dict(os.environ, {"DASHSCOPE_API_KEY": "sk-test-123"}, clear=False): + out = build_mcp_process_command( + "npx", + ["-y", "mcp-remote", "https://example.com/sse", "--header", "Authorization: Bearer ${DASHSCOPE_API_KEY}"], + store=None, + ) + self.assertEqual( + out, + ["npx", "-y", "mcp-remote", "https://example.com/sse", "--header", "Authorization: Bearer sk-test-123"], + ) + def test_policy_session_only_that_users_db_roots(self) -> None: with tempfile.TemporaryDirectory(ignore_cleanup_errors=True) as td: dbf = Path(td) / "ops.sqlite"