diff --git a/README.md b/README.md index 7dcc7ad4..35da7b89 100644 --- a/README.md +++ b/README.md @@ -2,6 +2,44 @@ This repository is fully consolidated under `oclaw/`. +## Quickstart (Open Source) + +### Prerequisites +- Python 3.11+ +- Node.js 22+ (required by the official Weixin plugin) + +### 1) Bootstrap venv (Windows) + +```powershell +powershell -ExecutionPolicy Bypass -File .\scripts\bootstrap_venv.ps1 +``` + +### 2) Start gateway (background) + +```powershell +powershell -ExecutionPolicy Bypass -File .\scripts\start_gateway.ps1 -SkipInstall -Background +``` + +Open: +- Admin: `http://127.0.0.1:8787/admin` +- Chat: `http://127.0.0.1:8787/chat` + +If port 8787 is stuck/occupied, force-stop: + +```powershell +powershell -ExecutionPolicy Bypass -File .\scripts\stop_gateway.ps1 -Force +``` + +### 3) Weixin (Personal WeChat): install → login → start + +```powershell +powershell -ExecutionPolicy Bypass -File .\runtime\operations\scripts\weixin_install.ps1 -UseOpenclawCli +powershell -ExecutionPolicy Bypass -File .\runtime\operations\scripts\weixin_login.ps1 +powershell -ExecutionPolicy Bypass -File .\runtime\operations\scripts\weixin_start.ps1 +``` + +Full runbook (recommended): see `docs/RUNBOOK.md` → “开源快速安装(从零到跑起来)”. + ## Layers - `runtime/`: core execution loop, routing, skill runtime, hook runtime - `interfaces/`: transport adapters (HTTP/WS) diff --git a/docs/RUNBOOK.md b/docs/RUNBOOK.md index a1508762..6f09452e 100644 --- a/docs/RUNBOOK.md +++ b/docs/RUNBOOK.md @@ -30,6 +30,90 @@ --- +## 1.1 开源快速安装(从零到跑起来) + +本节面向“第一次拿到开源仓库的用户”,目标是 **15 分钟内跑通**: + +- 网关(Admin + Chat) +- 微信(Personal WeChat)收消息、回消息(官方插件 + 本地原生宿主 `/weixin/native/reply`) + +### 1.1.1 前置依赖 + +- **Windows 10/11** +- **Python**:建议 `3.11+`(仓库脚本默认会创建 `.venv/`) +- **Node.js**:要求 `>=22`(官方微信插件声明 `engines.node >=22`) +- **npm**:随 Node 安装 +- (可选)Git:用于拉取仓库 + +> 注意:Node 版本过低会导致微信插件/sidecar 无法运行。 + +### 1.1.2 初始化 Python venv(Windows) + +在仓库根目录执行: + +```powershell +powershell -ExecutionPolicy Bypass -File .\scripts\bootstrap_venv.ps1 +``` + +### 1.1.3 启动网关(建议后台) + +```powershell +powershell -ExecutionPolicy Bypass -File .\scripts\start_gateway.ps1 -SkipInstall -Background +``` + +访问: + +- Admin:`http://127.0.0.1:8787/admin` +- Chat:`http://127.0.0.1:8787/chat` + +停止网关(如果遇到“端口被占用 / 重启不生效”,务必加 `-Force`): + +```powershell +powershell -ExecutionPolicy Bypass -File .\scripts\stop_gateway.ps1 -Force +``` + +### 1.1.4 安装官方微信插件(Personal WeChat) + +> 这一步会安装/更新官方插件到:`%USERPROFILE%\.openclaw\extensions\openclaw-weixin\` + +```powershell +powershell -ExecutionPolicy Bypass -File .\runtime\operations\scripts\weixin_install.ps1 -UseOpenclawCli +``` + +### 1.1.5 扫码登录(绑定账号) + +```powershell +powershell -ExecutionPolicy Bypass -File .\runtime\operations\scripts\weixin_login.ps1 +``` + +按提示扫码完成绑定(会写入账号 ID / token 等状态到 `%USERPROFILE%\.openclaw\openclaw-weixin\`)。 + +### 1.1.6 启动微信 sidecar(原生模式) + +```powershell +powershell -ExecutionPolicy Bypass -File .\runtime\operations\scripts\weixin_start.ps1 +``` + +检查状态: + +```powershell +powershell -ExecutionPolicy Bypass -File .\runtime\operations\scripts\weixin_status.ps1 +``` + +### 1.1.7 常见问题 + +- **发消息无回复** + - 先确认网关活着:访问 `http://127.0.0.1:8787/health` 应返回 `{"ok":"1"}` + - 再看微信 sidecar 日志: + - `data/channel_sidecar/oclaw-weixin/logs/weixin_sidecar.log` + - `data/channel_sidecar/oclaw-weixin/logs/weixin_sidecar.err.log` +- **启动网关提示端口占用 / 你以为重启了但没生效** + - 用 `stop_gateway.ps1 -Force` 强制按端口清理旧监听进程,然后再启动。 +- **Node 版本不对** + - 官方插件要求 `node >=22`;请升级 Node 后重新执行 `weixin_install.ps1 -UseOpenclawCli`。 + +--- + ## 2. 首次初始化 Windows: @@ -96,6 +180,36 @@ Linux/macOS: - 聊天页面统一使用 `http://127.0.0.1:8787/chat` - `--with-ui` 为历史参数,不再生效 +### 4.1 微信(Personal WeChat)当前接入模式 + +当前默认链路已经切到“官方插件优先”: + +- 官方插件模块负责: + - 扫码登录 + - 持久化账号 ID / bot token / context token + - 直接调用微信云端 `ilink/bot/getupdates`、`ilink/bot/sendmessage` + - 复用官方媒体下载/上传实现 +- 本仓库宿主适配负责: + - 把官方入站消息转换成 `oclaw` 可消费的本地 payload + - 通过本地 `/weixin/native/reply` 同步生成回复 + - 保留一个历史 `runner.ts` fallback,便于短期回滚 + +对应脚本行为: + +- `runtime/operations/scripts/weixin_install.ps1 -UseOpenclawCli` + - 安装官方 `openclaw-weixin` 插件 + - 安装运行官方模块所需的本地 Node 依赖 + - 同步 `runtime/operations/weixin_bridge/official_runner.ts` / `runner.ts` / `login.ts` +- `runtime/operations/scripts/weixin_start.ps1` + - 默认启动 `official_runner.ts` + - 设置 `AIA_WEIXIN_RUNNER_MODE=legacy` 时,临时回退到历史 `runner.ts` + +现阶段主路径建议: + +- 官方登录态 +- 官方收发与媒体模块 +- 本仓库本地 reply 宿主适配 + --- ## 5. 仅启动网关 diff --git a/interfaces/http/weixin_ilink_api.py b/interfaces/http/weixin_ilink_api.py index 8b1629d8..90755037 100644 --- a/interfaces/http/weixin_ilink_api.py +++ b/interfaces/http/weixin_ilink_api.py @@ -9,11 +9,19 @@ from typing import Any from fastapi import APIRouter, Header, HTTPException, Request -from oclaw.runtime.application.gateway import process_inbound_payload_usecase +def _process_inbound_payload_usecase(payload: dict[str, Any]) -> dict[str, Any]: + # Lazy import to avoid circular imports during FastAPI app bootstrap. + from oclaw.runtime.application.gateway import process_inbound_payload_usecase + + return process_inbound_payload_usecase(payload) router = APIRouter() +# Historical iLink-compatible fallback endpoints. The primary path is now +# `/weixin/native/reply` + `official_runner.ts`, but these routes remain for +# short-term rollback compatibility. + def _require_ilink_auth( *, @@ -123,6 +131,48 @@ def _extract_inbound_identity(body: dict[str, Any]) -> tuple[str, str]: return user_id, chat_id +def _build_native_reply_payload(body: dict[str, Any]) -> dict[str, Any]: + channel = _normalize_channel(body.get("channel")) + account_id = _resolve_account_id(body) + ctx = body.get("ctx") if isinstance(body.get("ctx"), dict) else {} + metadata = body.get("metadata") if isinstance(body.get("metadata"), dict) else {} + attachments = body.get("attachments") if isinstance(body.get("attachments"), list) else [] + + user_id = str( + body.get("user_id") + or ctx.get("From") + or ctx.get("To") + or body.get("external_user_id") + or "" + ).strip() + chat_id = str( + body.get("chat_id") + or ctx.get("To") + or ctx.get("From") + or body.get("external_chat_id") + or user_id + or "" + ).strip() + text = str(body.get("text") or ctx.get("Body") or ctx.get("CommandBody") or "").strip() + if not metadata: + metadata = {"source": "weixin_official_native"} + else: + metadata = dict(metadata) + metadata.setdefault("source", "weixin_official_native") + if ctx: + metadata["weixin_ctx"] = ctx + + return { + "channel": channel, + "account_id": account_id, + "user_id": user_id, + "chat_id": chat_id or user_id, + "text": text, + "attachments": [a for a in attachments if isinstance(a, dict)], + "metadata": metadata, + } + + class _IlinkBridge: def __init__(self) -> None: self._lock = threading.Lock() @@ -236,7 +286,7 @@ async def _process_inbound_and_enqueue( return "" try: - out = await asyncio.wait_for(asyncio.to_thread(process_inbound_payload_usecase, payload), timeout=60.0) + out = await asyncio.wait_for(asyncio.to_thread(_process_inbound_payload_usecase, payload), timeout=60.0) except asyncio.TimeoutError: ctx_token = _extract_context_token(payload) _BRIDGE.enqueue_reply( @@ -406,5 +456,27 @@ async def ilink_sendtyping( return {"ret": 0} +@router.post("/weixin/native/reply") +async def weixin_native_reply( + body: dict[str, Any], + authorizationtype: str | None = Header(default=None, alias="AuthorizationType"), + authorization: str | None = Header(default=None, alias="Authorization"), +) -> dict[str, Any]: + _require_ilink_auth(authorization_type=authorizationtype, authorization=authorization) + payload = _build_native_reply_payload(body) + if not str(payload.get("user_id") or "").strip(): + return {"ok": False, "error": "missing user_id", "replies": []} + if not str(payload.get("account_id") or "").strip(): + return {"ok": False, "error": "missing account_id", "replies": []} + out = await asyncio.to_thread(_process_inbound_payload_usecase, payload) + replies = out.get("replies") if isinstance(out, dict) else [] + if not isinstance(replies, list): + replies = [] + return { + "ok": bool((out or {}).get("ok", True)) if isinstance(out, dict) else True, + "replies": [r for r in replies if isinstance(r, dict)], + } + + __all__ = ["router"] diff --git a/runtime/operations/scripts/stop_gateway.ps1 b/runtime/operations/scripts/stop_gateway.ps1 index 316f3807..4f7f39f5 100644 --- a/runtime/operations/scripts/stop_gateway.ps1 +++ b/runtime/operations/scripts/stop_gateway.ps1 @@ -37,9 +37,15 @@ if (Test-Path $pidFile) { $ok = Kill-ProcId -procId $procId if ($ok) { Remove-Item $pidFile -Force -ErrorAction SilentlyContinue - exit 0 + # Keep going to also clean up any orphan listeners on $Port. + } else { + # Stale pid file is common after crashes; fall back to stop-by-port. + if ($Force) { + Remove-Item $pidFile -Force -ErrorAction SilentlyContinue + } else { + Warn "PID file kill failed; falling back to stop-by-port. Re-run with -Force to also clear pid file." + } } - if (-not $Force) { exit 1 } } } diff --git a/runtime/operations/scripts/weixin_install.ps1 b/runtime/operations/scripts/weixin_install.ps1 index c29270b4..a95ed6c2 100644 --- a/runtime/operations/scripts/weixin_install.ps1 +++ b/runtime/operations/scripts/weixin_install.ps1 @@ -15,11 +15,27 @@ 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" 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 + } +} + if ($UseOpenclawCli) { $openclawCmd = Get-Command openclaw -ErrorAction SilentlyContinue if (-not $openclawCmd) { @@ -29,6 +45,7 @@ if ($UseOpenclawCli) { if ($LASTEXITCODE -ne 0) { throw "openclaw-weixin-cli install failed with exit code $LASTEXITCODE" } + Ensure-OfficialPluginRuntimeDeps Push-Location $sidecarRoot try { if (-not (Test-Path (Join-Path $sidecarRoot "package.json"))) { @@ -37,17 +54,18 @@ if ($UseOpenclawCli) { throw "npm init failed with exit code $LASTEXITCODE" } } - npm.cmd install --save-exact tsx@4.21.0 typescript@6.0.3 + npm.cmd install openclaw@latest --save tsx@4.21.0 typescript@6.0.3 if ($LASTEXITCODE -ne 0) { throw "npm install bridge runtime deps failed with exit code $LASTEXITCODE" } $bridgeSrc = Join-Path $oclawRoot "runtime\\operations\\weixin_bridge" Copy-Item -Path (Join-Path $bridgeSrc "runner.ts") -Destination (Join-Path $sidecarRoot "runner.ts") -Force + Copy-Item -Path (Join-Path $bridgeSrc "official_runner.ts") -Destination (Join-Path $sidecarRoot "official_runner.ts") -Force Copy-Item -Path (Join-Path $bridgeSrc "login.ts") -Destination (Join-Path $sidecarRoot "login.ts") -Force } finally { Pop-Location } - Write-Host "[ok] installed official openclaw-weixin plugin + local bridge runtime" + Write-Host "[ok] installed official openclaw-weixin plugin + native/fallback sidecar runtime" exit 0 } diff --git a/runtime/operations/scripts/weixin_start.ps1 b/runtime/operations/scripts/weixin_start.ps1 index 492bb7f1..75decd41 100644 --- a/runtime/operations/scripts/weixin_start.ps1 +++ b/runtime/operations/scripts/weixin_start.ps1 @@ -16,12 +16,18 @@ $sidecarRoot = Join-Path $oclawRoot "data\\channel_sidecar\\$ChannelId" $stateDir = Join-Path $sidecarRoot "state" $logDir = Join-Path $sidecarRoot "logs" $pidFile = Join-Path $sidecarRoot "pid.txt" +$bridgeSrc = Join-Path $oclawRoot "runtime\\operations\\weixin_bridge" +$pluginRoot = Join-Path $env:USERPROFILE ".openclaw\\extensions\\openclaw-weixin" +$runnerMode = [string]($env:AIA_WEIXIN_RUNNER_MODE) +if (-not $runnerMode) { $runnerMode = "official" } +$runnerMode = $runnerMode.Trim().ToLowerInvariant() function Get-SidecarProcesses { $escapedSidecarRoot = $sidecarRoot.Replace("\", "\\") $patterns = @( "*$ChannelId*", "*runner.ts*", + "*official_runner.ts*", "*$escapedSidecarRoot*" ) Get-CimInstance Win32_Process | Where-Object { @@ -46,38 +52,21 @@ function Stop-SidecarProcesses { return $procs.Count } -function Set-OfficialWeixinBaseUrl([string]$BaseUrl) { - $weixinRoot = Join-Path $env:USERPROFILE ".openclaw\\openclaw-weixin" - $accountsListPath = Join-Path $weixinRoot "accounts.json" - if (-not (Test-Path $accountsListPath)) { - Write-Host "[warn] official mode: accounts.json not found, skip baseUrl rewrite" +function Ensure-OfficialPluginRuntimeDeps { + if (-not (Test-Path (Join-Path $pluginRoot "package.json"))) { + throw "official plugin root not found at $pluginRoot" + } + if (Test-Path (Join-Path $pluginRoot "node_modules\\openclaw\\package.json")) { return } - $ids = @() + Push-Location $pluginRoot try { - $parsed = Get-Content -Path $accountsListPath -Raw | ConvertFrom-Json - if ($parsed -is [System.Array]) { - $ids = @($parsed) - } - } catch { - Write-Host "[warn] official mode: failed to parse accounts.json" - return - } - foreach ($aid in $ids) { - $idText = [string]$aid - if (-not $idText) { continue } - $accPath = Join-Path (Join-Path $weixinRoot "accounts") "$idText.json" - if (-not (Test-Path $accPath)) { continue } - try { - $obj = Get-Content -Path $accPath -Raw | ConvertFrom-Json - $obj.baseUrl = $BaseUrl - $json = $obj | ConvertTo-Json -Depth 8 - $utf8NoBom = New-Object System.Text.UTF8Encoding($false) - [System.IO.File]::WriteAllText($accPath, $json + "`n", $utf8NoBom) - Write-Host "[ok] official mode: set baseUrl for $idText -> $BaseUrl" - } catch { - Write-Host "[warn] official mode: failed to rewrite $accPath" + 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 } } @@ -90,17 +79,36 @@ New-Item -ItemType Directory -Force -Path $stateDir | Out-Null $cleaned = Stop-SidecarProcesses Remove-Item -Force $pidFile -ErrorAction SilentlyContinue +if (Test-Path $bridgeSrc) { + foreach ($name in @("runner.ts", "official_runner.ts", "login.ts")) { + $srcPath = Join-Path $bridgeSrc $name + if (Test-Path $srcPath) { + Copy-Item -Path $srcPath -Destination (Join-Path $sidecarRoot $name) -Force + } + } +} + $logPath = Join-Path $logDir "weixin_sidecar.log" $errPath = Join-Path $logDir "weixin_sidecar.err.log" -if (Test-Path (Join-Path $sidecarRoot "runner.ts")) { +if ((Test-Path (Join-Path $sidecarRoot "runner.ts")) -or (Test-Path (Join-Path $sidecarRoot "official_runner.ts"))) { + $runnerFile = "official_runner.ts" + if ($runnerMode -eq "legacy") { + $runnerFile = "runner.ts" + } + if ($runnerMode -eq "official") { + Ensure-OfficialPluginRuntimeDeps + } + if (-not (Test-Path (Join-Path $sidecarRoot $runnerFile))) { + throw "selected runner file missing: $runnerFile" + } $cmd = "cmd.exe" $args = @( "/c", - "cd /d $sidecarRoot && set OCLAW_STATE_DIR=$stateDir&& set AIA_GATEWAY_BASE_URL=$GatewayBaseUrl&& npm.cmd exec -- tsx runner.ts" + "cd /d $sidecarRoot && set OCLAW_STATE_DIR=$stateDir&& set AIA_GATEWAY_BASE_URL=$GatewayBaseUrl&& set NODE_PATH=$sidecarRoot\node_modules&& npm.cmd exec -- tsx $runnerFile" ) $p = Start-Process -FilePath $cmd -ArgumentList $args -WorkingDirectory $sidecarRoot -PassThru -WindowStyle Hidden -RedirectStandardOutput $logPath -RedirectStandardError $errPath Set-Content -Path $pidFile -Value $p.Id - Write-Host "[ok] started weixin sidecar pid=$($p.Id) cleaned=$cleaned out=$logPath err=$errPath" + Write-Host "[ok] started weixin sidecar pid=$($p.Id) mode=$runnerMode cleaned=$cleaned out=$logPath err=$errPath" exit 0 } @@ -113,9 +121,8 @@ $systemNodeDir = "C:\\Program Files\\nodejs" if (Test-Path (Join-Path $systemNodeDir "node.exe")) { $env:PATH = "$systemNodeDir;$env:PATH" } -Set-OfficialWeixinBaseUrl -BaseUrl $GatewayBaseUrl $args = @("/c", "openclaw gateway --allow-unconfigured") $p = Start-Process -FilePath "cmd.exe" -ArgumentList $args -WorkingDirectory $oclawRoot -PassThru -WindowStyle Hidden -RedirectStandardOutput $logPath -RedirectStandardError $errPath Set-Content -Path $pidFile -Value $p.Id -Write-Host "[ok] started openclaw gateway bridge pid=$($p.Id) cleaned=$cleaned out=$logPath err=$errPath" +Write-Host "[ok] started openclaw gateway pid=$($p.Id) cleaned=$cleaned out=$logPath err=$errPath" diff --git a/runtime/operations/scripts/weixin_status.ps1 b/runtime/operations/scripts/weixin_status.ps1 index 82f59d24..a992ef83 100644 --- a/runtime/operations/scripts/weixin_status.ps1 +++ b/runtime/operations/scripts/weixin_status.ps1 @@ -23,6 +23,7 @@ function Get-SidecarProcesses { $patterns = @( "*$ChannelId*", "*runner.ts*", + "*official_runner.ts*", "*$escapedSidecarRoot*" ) Get-CimInstance Win32_Process | Where-Object { diff --git a/runtime/operations/scripts/weixin_stop.ps1 b/runtime/operations/scripts/weixin_stop.ps1 index a998b899..67636303 100644 --- a/runtime/operations/scripts/weixin_stop.ps1 +++ b/runtime/operations/scripts/weixin_stop.ps1 @@ -24,6 +24,7 @@ function Get-SidecarProcesses { $patterns = @( "*$ChannelId*", "*runner.ts*", + "*official_runner.ts*", "*$escapedSidecarRoot*" ) Get-CimInstance Win32_Process | Where-Object { diff --git a/runtime/operations/weixin_bridge/official_runner.ts b/runtime/operations/weixin_bridge/official_runner.ts new file mode 100644 index 00000000..569d8e1b --- /dev/null +++ b/runtime/operations/weixin_bridge/official_runner.ts @@ -0,0 +1,430 @@ +import fs from "node:fs"; +import os from "node:os"; +import path from "node:path"; +import { pathToFileURL } from "node:url"; + +type Json = Record; +type TokenMap = Record; + +const LOCAL_BASE_URL = (process.env.AIA_GATEWAY_BASE_URL || "http://127.0.0.1:8787").trim(); +const STATE_DIR = (process.env.OCLAW_STATE_DIR || path.resolve(process.cwd(), "state")).trim(); +const STATE_FILE = path.join(STATE_DIR, "official_bridge_state.json"); +const POLL_TIMEOUT_MS = 35_000; +const DEFAULT_CDN_BASE_URL = "https://novac2c.cdn.weixin.qq.com/c2c"; + +type OfficialModules = { + getUpdates: (params: { + baseUrl: string; + token?: string; + get_updates_buf?: string; + timeoutMs?: number; + }) => Promise; + setContextToken: (accountId: string, userId: string, token: string) => void; + getContextToken: (accountId: string, userId: string) => string | undefined; + restoreContextTokens: (accountId: string) => void; + weixinMessageToMsgContext: (msg: Json, accountId: string, opts?: Json) => Json; + downloadMediaFromItem: ( + item: Json, + deps: { + cdnBaseUrl: string; + saveMedia: ( + buffer: Buffer, + contentType?: string, + subdir?: string, + maxBytes?: number, + originalFilename?: string, + ) => Promise<{ path: string }>; + log: (msg: string) => void; + errLog: (msg: string) => void; + label: string; + }, + ) => Promise; + sendMessageWeixin: (params: { + to: string; + text: string; + opts: { baseUrl: string; token?: string; contextToken?: string }; + }) => Promise<{ messageId: string }>; + sendWeixinMediaFile: (params: { + filePath: string; + to: string; + text: string; + opts: { baseUrl: string; token?: string; contextToken?: string }; + cdnBaseUrl: string; + }) => Promise<{ messageId: string }>; +}; + +let officialModulesPromise: Promise | null = null; + +class HttpStatusError extends Error { + status: number; + + constructor(status: number, message: string) { + super(message); + this.name = "HttpStatusError"; + this.status = status; + } +} + +function log(msg: string): void { + process.stdout.write(`${new Date().toISOString()} [official-weixin] ${msg}\n`); +} + +function ensureDir(dir: string): void { + fs.mkdirSync(dir, { recursive: true }); +} + +function readJsonFile(p: string): T | null { + try { + return JSON.parse(fs.readFileSync(p, "utf8")) as T; + } catch { + return null; + } +} + +function writeJsonFile(p: string, obj: unknown): void { + fs.writeFileSync(p, JSON.stringify(obj, null, 2) + "\n", "utf8"); +} + +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"); +} + +async function loadOfficialModules(): Promise { + if (officialModulesPromise) { + return officialModulesPromise; + } + officialModulesPromise = (async () => { + 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 [apiMod, inboundMod, mediaMod, sendMod, sendMediaMod] = await Promise.all([ + importTs(path.join("api", "api.ts")), + importTs(path.join("messaging", "inbound.ts")), + importTs(path.join("media", "media-download.ts")), + importTs(path.join("messaging", "send.ts")), + importTs(path.join("messaging", "send-media.ts")), + ]); + return { + getUpdates: apiMod.getUpdates, + setContextToken: inboundMod.setContextToken, + getContextToken: inboundMod.getContextToken, + restoreContextTokens: inboundMod.restoreContextTokens, + weixinMessageToMsgContext: inboundMod.weixinMessageToMsgContext, + downloadMediaFromItem: mediaMod.downloadMediaFromItem, + sendMessageWeixin: sendMod.sendMessageWeixin, + sendWeixinMediaFile: sendMediaMod.sendWeixinMediaFile, + } satisfies OfficialModules; + })(); + return officialModulesPromise; +} + +function resolveAccount(): { accountId: string; token: string; cloudBaseUrl: string } { + const ids = readJsonFile(homeOpenclawPath("openclaw-weixin", "accounts.json")) || []; + 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(); + 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 cloudBaseUrl = (envCloud || cfgCloud || "https://ilinkai.weixin.qq.com").trim(); + return { accountId, token, cloudBaseUrl }; +} + +function resolveHeaders(token: string): Record { + return { + "Content-Type": "application/json", + AuthorizationType: "ilink_bot_token", + Authorization: `Bearer ${token}`, + }; +} + +async function postNativeReply(token: string, body: Json): Promise { + const url = `${LOCAL_BASE_URL.replace(/\/+$/, "")}/weixin/native/reply`; + const res = await fetch(url, { + method: "POST", + headers: resolveHeaders(token), + body: JSON.stringify(body), + }); + const text = await res.text(); + if (!res.ok) { + throw new HttpStatusError(res.status, `native reply ${res.status}: ${text.slice(0, 300)}`); + } + return text ? (JSON.parse(text) as Json) : {}; +} + +async function safeSendNativeFailureNotice(modules: OfficialModules, params: { + to: string; + cloudBaseUrl: string; + token: string; + contextToken?: string; + err: unknown; +}): Promise { + const msg = `[weixin] native reply failed: ${String(params.err)}`.slice(0, 3500); + try { + await modules.sendMessageWeixin({ + to: params.to, + text: msg, + opts: { + baseUrl: params.cloudBaseUrl, + token: params.token, + contextToken: params.contextToken, + }, + }); + } catch { + // best-effort + } +} + +function inferExtension(contentType: string, originalFilename: string): string { + const explicit = path.extname(originalFilename || "").trim(); + if (explicit) { + return explicit; + } + const low = (contentType || "").toLowerCase(); + if (low.includes("jpeg")) return ".jpg"; + if (low.includes("png")) return ".png"; + if (low.includes("gif")) return ".gif"; + if (low.includes("webp")) return ".webp"; + if (low.includes("mp4")) return ".mp4"; + if (low.includes("wav")) return ".wav"; + if (low.includes("pdf")) return ".pdf"; + if (low.includes("text/plain")) return ".txt"; + return ".bin"; +} + +async function saveMediaBuffer( + buffer: Buffer, + contentType?: string, + subdir?: string, + maxBytes?: number, + originalFilename?: string, +): Promise<{ path: string }> { + if (typeof maxBytes === "number" && maxBytes > 0 && buffer.length > maxBytes) { + throw new Error(`media exceeds max bytes: ${buffer.length} > ${maxBytes}`); + } + const dir = path.join(STATE_DIR, "media", subdir || "inbound"); + ensureDir(dir); + const filePath = path.join( + dir, + `${Date.now()}-${Math.random().toString(16).slice(2, 10)}${inferExtension(contentType || "", originalFilename || "")}`, + ); + await fs.promises.writeFile(filePath, buffer); + return { path: filePath }; +} + +function pickDownloadableMedia(msg: Json): Json | null { + const items = Array.isArray(msg.item_list) ? (msg.item_list as Json[]) : []; + const hasDownloadableMedia = (item: Json, kind: string): boolean => { + const media = item[kind] && typeof item[kind] === "object" ? (item[kind] as Json).media : null; + if (!media || typeof media !== "object") return false; + return Boolean((media as Json).encrypt_query_param || (media as Json).full_url); + }; + const direct = + items.find((item) => Number(item.type || 0) === 3 && hasDownloadableMedia(item, "image_item")) || + items.find((item) => Number(item.type || 0) === 4 && hasDownloadableMedia(item, "video_item")) || + items.find((item) => Number(item.type || 0) === 5 && hasDownloadableMedia(item, "file_item")) || + items.find((item) => Number(item.type || 0) === 2 && hasDownloadableMedia(item, "voice_item")); + if (direct) { + return direct; + } + for (const item of items) { + if (Number(item.type || 0) !== 1) continue; + const refMsg = item.ref_msg; + if (!refMsg || typeof refMsg !== "object") continue; + const messageItem = (refMsg as Json).message_item; + if (!messageItem || typeof messageItem !== "object") continue; + const ref = messageItem as Json; + if ( + (Number(ref.type || 0) === 3 && hasDownloadableMedia(ref, "image_item")) || + (Number(ref.type || 0) === 4 && hasDownloadableMedia(ref, "video_item")) || + (Number(ref.type || 0) === 5 && hasDownloadableMedia(ref, "file_item")) || + (Number(ref.type || 0) === 2 && hasDownloadableMedia(ref, "voice_item")) + ) { + return ref; + } + } + return null; +} + +function buildAttachmentsFromMedia(mediaOpts: Json): Json[] { + const out: Json[] = []; + const pushIf = (key: string, mimeKey: string, kind: string): void => { + const filePath = String(mediaOpts[key] || "").trim(); + if (!filePath) return; + out.push({ + kind, + local_path: filePath, + media_type: String(mediaOpts[mimeKey] || "").trim(), + }); + }; + pushIf("decryptedPicPath", "", "image"); + pushIf("decryptedVideoPath", "", "video"); + pushIf("decryptedFilePath", "fileMediaType", "file"); + pushIf("decryptedVoicePath", "voiceMediaType", "voice"); + return out; +} + +async function handleInboundMessage( + modules: OfficialModules, + params: { + token: string; + accountId: string; + cloudBaseUrl: string; + full: Json; + userContextTokens: TokenMap; + localCursor: string; + }, +): Promise { + const fromUser = String(params.full.from_user_id || "").trim(); + const toUser = String(params.full.to_user_id || "").trim(); + if (!fromUser) return params.localCursor; + if (toUser && toUser === fromUser) return params.localCursor; + if (Number(params.full.message_type || 1) !== 1) return params.localCursor; + + const contextToken = String(params.full.context_token || "").trim(); + if (contextToken) { + params.userContextTokens[fromUser] = contextToken; + modules.setContextToken(params.accountId, fromUser, contextToken); + } + + const mediaItem = pickDownloadableMedia(params.full); + const mediaOpts = mediaItem + ? await modules.downloadMediaFromItem(mediaItem, { + cdnBaseUrl: DEFAULT_CDN_BASE_URL, + saveMedia: saveMediaBuffer, + log: (msg: string) => log(`media ${msg}`), + errLog: (msg: string) => log(`media-error ${msg}`), + label: "inbound", + }) + : {}; + const ctx = modules.weixinMessageToMsgContext(params.full, params.accountId, mediaOpts); + let replies: Json[] = []; + try { + const native = await postNativeReply(params.token, { + channel: "wechat", + account_id: params.accountId, + ctx, + attachments: buildAttachmentsFromMedia(mediaOpts), + metadata: { + source: "weixin_official_native", + raw: { + msg: params.full, + }, + }, + }); + replies = Array.isArray(native.replies) ? (native.replies as Json[]) : []; + } catch (err) { + log(`native reply failed; no fallback enabled err=${String(err)}`); + await safeSendNativeFailureNotice(modules, { + to: fromUser, + cloudBaseUrl: params.cloudBaseUrl, + token: params.token, + contextToken: contextToken || undefined, + err, + }); + return params.localCursor; + } + for (const reply of replies) { + const text = String(reply.text || "").trim(); + const mediaPath = String(reply.media_path || reply.mediaPath || "").trim(); + const mediaUrl = String(reply.media_url || reply.mediaUrl || "").trim(); + const deliverTo = String(reply.chat_id || reply.to || fromUser).trim() || fromUser; + const replyContextToken = String( + reply.context_token || modules.getContextToken(params.accountId, deliverTo) || contextToken || "", + ).trim(); + if (mediaPath || mediaUrl) { + const filePath = mediaPath || mediaUrl; + await modules.sendWeixinMediaFile({ + filePath, + to: deliverTo, + text, + opts: { + baseUrl: params.cloudBaseUrl, + token: params.token, + contextToken: replyContextToken || undefined, + }, + cdnBaseUrl: DEFAULT_CDN_BASE_URL, + }); + log(`reply media sent to=${deliverTo} textLen=${text.length}`); + continue; + } + if (!text) continue; + await modules.sendMessageWeixin({ + to: deliverTo, + text, + opts: { + baseUrl: params.cloudBaseUrl, + token: params.token, + contextToken: replyContextToken || undefined, + }, + }); + log(`reply text sent to=${deliverTo} textLen=${text.length}`); + } + return params.localCursor; +} + +async function main(): Promise { + ensureDir(STATE_DIR); + const state = (readJsonFile(STATE_FILE) || {}) as Json; + let cloudCursor = String(state.cloud_cursor || "").trim(); + const userContextTokens: TokenMap = + state.user_context_tokens && typeof state.user_context_tokens === "object" + ? (state.user_context_tokens as TokenMap) + : {}; + const { accountId, token, cloudBaseUrl } = resolveAccount(); + const modules = await loadOfficialModules(); + modules.restoreContextTokens(accountId); + log(`official runner started account=${accountId} cloud=${cloudBaseUrl} local=${LOCAL_BASE_URL}`); + while (true) { + try { + const out = await modules.getUpdates({ + baseUrl: cloudBaseUrl, + token, + get_updates_buf: cloudCursor, + timeoutMs: POLL_TIMEOUT_MS + 5000, + }); + const msgs = Array.isArray(out.msgs) ? (out.msgs as Json[]) : []; + const nextCloudCursor = String(out.get_updates_buf || cloudCursor || "").trim(); + if (nextCloudCursor) { + cloudCursor = nextCloudCursor; + } + for (const full of msgs) { + await handleInboundMessage(modules, { + token, + accountId, + cloudBaseUrl, + full, + userContextTokens, + localCursor: "0", + }); + } + writeJsonFile(STATE_FILE, { + cloud_cursor: cloudCursor, + user_context_tokens: userContextTokens, + updated_at: new Date().toISOString(), + }); + } catch (err) { + log(`loop error: ${String(err)}`); + await sleep(1500); + } + } +} + +void main(); diff --git a/runtime/operations/weixin_bridge/runner.ts b/runtime/operations/weixin_bridge/runner.ts index f02ab158..bfbf7707 100644 --- a/runtime/operations/weixin_bridge/runner.ts +++ b/runtime/operations/weixin_bridge/runner.ts @@ -1,3 +1,5 @@ +// Legacy fallback bridge. The default startup path now uses `official_runner.ts`, +// which drives the official openclaw-weixin modules directly. import fs from "node:fs"; import crypto from "node:crypto"; import os from "node:os"; diff --git a/tests/test_weixin_ilink_api.py b/tests/test_weixin_ilink_api.py index 1cdd2183..109775bb 100644 --- a/tests/test_weixin_ilink_api.py +++ b/tests/test_weixin_ilink_api.py @@ -30,7 +30,7 @@ class WeixinIlinkApiTests(unittest.TestCase): self.assertEqual((r.json() or {}).get("ret"), 400) def test_sendmessage_enqueues_reply_for_getupdates(self) -> None: - old_usecase = weixin_ilink_api.process_inbound_payload_usecase + old_usecase = weixin_ilink_api._process_inbound_payload_usecase def _fake_usecase(payload: dict[str, object]) -> dict[str, object]: text = str(payload.get("text") or "") @@ -45,7 +45,7 @@ class WeixinIlinkApiTests(unittest.TestCase): } try: - weixin_ilink_api.process_inbound_payload_usecase = _fake_usecase # type: ignore[assignment] + weixin_ilink_api._process_inbound_payload_usecase = _fake_usecase # type: ignore[assignment] s = self.client.post( "/ilink/bot/sendmessage", headers=self.headers, @@ -77,7 +77,49 @@ class WeixinIlinkApiTests(unittest.TestCase): self.assertTrue(msgs, data) self.assertEqual(str((msgs[0] or {}).get("text") or ""), "echo:ping") finally: - weixin_ilink_api.process_inbound_payload_usecase = old_usecase # type: ignore[assignment] + weixin_ilink_api._process_inbound_payload_usecase = old_usecase # type: ignore[assignment] + + def test_native_reply_returns_sync_replies(self) -> None: + old_usecase = weixin_ilink_api._process_inbound_payload_usecase + + def _fake_usecase(payload: dict[str, object]) -> dict[str, object]: + text = str(payload.get("text") or "") + self.assertEqual(str(payload.get("channel") or ""), "wechat") + self.assertEqual(str(payload.get("user_id") or ""), "wxid_u3") + return { + "ok": True, + "replies": [ + { + "chat_id": str(payload.get("chat_id") or ""), + "text": f"native:{text}", + } + ], + } + + try: + weixin_ilink_api._process_inbound_payload_usecase = _fake_usecase # type: ignore[assignment] + r = self.client.post( + "/weixin/native/reply", + headers=self.headers, + json={ + "channel": "wechat", + "account_id": "bot-1", + "ctx": { + "From": "wxid_u3", + "To": "wxid_u3", + "Body": "hello native", + "context_token": "ctx-1", + }, + }, + ) + self.assertEqual(r.status_code, 200, r.text) + data = r.json() or {} + self.assertTrue(data.get("ok"), data) + replies = data.get("replies") if isinstance(data.get("replies"), list) else [] + self.assertEqual(len(replies), 1, data) + self.assertEqual(str((replies[0] or {}).get("text") or ""), "native:hello native") + finally: + weixin_ilink_api._process_inbound_payload_usecase = old_usecase # type: ignore[assignment] if __name__ == "__main__": diff --git a/tests/test_weixin_official_scripts.py b/tests/test_weixin_official_scripts.py new file mode 100644 index 00000000..f3d6e6f2 --- /dev/null +++ b/tests/test_weixin_official_scripts.py @@ -0,0 +1,50 @@ +from __future__ import annotations + +from pathlib import Path + + +REPO_ROOT = Path(__file__).resolve().parents[1] + + +def _read(rel: str) -> str: + return (REPO_ROOT / rel).read_text(encoding="utf-8") + + +def test_weixin_start_defaults_to_official_runner() -> None: + text = _read("runtime/operations/scripts/weixin_start.ps1") + assert '$runnerMode = "official"' in text + assert '$runnerFile = "official_runner.ts"' in text + assert 'AIA_WEIXIN_RUNNER_MODE' in text + assert 'official_runner.ts' in text + + +def test_weixin_status_and_stop_track_official_runner() -> None: + status_text = _read("runtime/operations/scripts/weixin_status.ps1") + stop_text = _read("runtime/operations/scripts/weixin_stop.ps1") + assert "official_runner.ts" in status_text + assert "official_runner.ts" in stop_text + + +def test_weixin_install_copies_official_runner() -> None: + text = _read("runtime/operations/scripts/weixin_install.ps1") + assert 'Copy-Item -Path (Join-Path $bridgeSrc "official_runner.ts")' in text + assert "Ensure-OfficialPluginRuntimeDeps" in text + + +def test_weixin_start_ensures_plugin_runtime_deps() -> None: + text = _read("runtime/operations/scripts/weixin_start.ps1") + assert "Ensure-OfficialPluginRuntimeDeps" in text + assert "npm.cmd install openclaw@latest --no-save" in text + + +def test_official_runner_prefers_direct_inbound_before_other_fallbacks() -> None: + text = _read("runtime/operations/weixin_bridge/official_runner.ts") + assert "const native = await postNativeReply" in text + assert "falling back to direct /inbound bridge" not in text + assert "legacy local ilink bridge" not in text + + +def test_official_runner_logs_active_bridge_path() -> None: + text = _read("runtime/operations/weixin_bridge/official_runner.ts") + assert "official runner started account=" in text + assert "native reply failed; no fallback enabled" in text