From b086c5e35899784f0018a70ec37ac9b75f459bdd Mon Sep 17 00:00:00 2001 From: Zhuyuan Operations Date: Wed, 29 Jul 2026 18:49:26 +0800 Subject: [PATCH] =?UTF-8?q?feat:=20=E6=8E=A5=E5=85=A5=E8=8B=8D=E8=80=B3?= =?UTF-8?q?=E5=AE=B6=E5=BA=AD=E6=9C=8D=E5=8A=A1=E5=99=A8=E8=BF=90=E7=BB=B4?= =?UTF-8?q?=E4=B8=8E=E5=B9=BF=E6=92=AD=E6=8E=88=E6=9D=83=E5=85=A5=E5=8F=A3?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- CANGER-REMOTE-OPS-ENTRY.hdlp | 49 + remote-node/README.md | 195 ++++ remote-node/config/control.env.example | 20 + remote-node/config/node.example.json | 26 + remote-node/config/pufferfish-persona.json | 14 + remote-node/install/install-node.sh | 107 ++ remote-node/install/preflight-linux.py | 85 ++ remote-node/install/苍耳一键安装.sh | 18 + remote-node/nginx/canger-node.conf.example | 10 + remote-node/node/canger_node.py | 408 ++++++++ .../persona/DEVELOPMENT-RECEIPT-20260729.md | 57 + remote-node/persona/README.md | 69 ++ remote-node/persona/WORKORDER-GUIDE.md | 77 ++ remote-node/server/control_plane.py | 978 ++++++++++++++++++ remote-node/shared.py | 111 ++ remote-node/systemd/canger-control.service | 22 + remote-node/systemd/canger-node.service | 21 + remote-node/tests/test_end_to_end.py | 299 ++++++ remote-node/tools/canger_persona.py | 182 ++++ remote-node/tools/cangerctl.py | 168 +++ 20 files changed, 2916 insertions(+) create mode 100644 CANGER-REMOTE-OPS-ENTRY.hdlp create mode 100644 remote-node/README.md create mode 100644 remote-node/config/control.env.example create mode 100644 remote-node/config/node.example.json create mode 100644 remote-node/config/pufferfish-persona.json create mode 100755 remote-node/install/install-node.sh create mode 100755 remote-node/install/preflight-linux.py create mode 100755 remote-node/install/苍耳一键安装.sh create mode 100644 remote-node/nginx/canger-node.conf.example create mode 100755 remote-node/node/canger_node.py create mode 100644 remote-node/persona/DEVELOPMENT-RECEIPT-20260729.md create mode 100644 remote-node/persona/README.md create mode 100644 remote-node/persona/WORKORDER-GUIDE.md create mode 100755 remote-node/server/control_plane.py create mode 100644 remote-node/shared.py create mode 100644 remote-node/systemd/canger-control.service create mode 100644 remote-node/systemd/canger-node.service create mode 100644 remote-node/tests/test_end_to_end.py create mode 100755 remote-node/tools/canger_persona.py create mode 100755 remote-node/tools/cangerctl.py diff --git a/CANGER-REMOTE-OPS-ENTRY.hdlp b/CANGER-REMOTE-OPS-ENTRY.hdlp new file mode 100644 index 0000000..a36e7d1 --- /dev/null +++ b/CANGER-REMOTE-OPS-ENTRY.hdlp @@ -0,0 +1,49 @@ +# CANGER-REMOTE-OPS-ENTRY · 苍耳服务器运维入口 + +> HLDP://cang-ying/CANGER-REMOTE-OPS-ENTRY +> 状态: ACTIVE_ENTRY · 2026-07-29 +> 仓主与最终批准者: 苍耳 + +```text +trigger: + 人格体需要检查、开发、调试或维护苍耳家庭 Ubuntu + +entry: + remote-node/persona/README.md + +load_order: + 1. remote-node/persona/DEVELOPMENT-RECEIPT-20260729.md + 2. remote-node/persona/WORKORDER-GUIDE.md + 3. remote-node/config/pufferfish-persona.json + +workorder: + requester_name: "谁申请就写谁" + requester_id: "申请者自己的正式编号" + purpose: "本次准备做什么" + duration: "本次希望打开的受控时段" + +approval: + authority: 苍耳 + rule: "申请者不能批准自己;苍耳收到通知后自行决定是否点击" + +route: + 申请者 + → 新加坡语言主控广播塔 + → 第五域语言主控广播塔 + → 苍耳审批频道 + → 苍耳邮箱 + → 苍耳批准 + → 同一申请者获得限时会话 + +current: + home_node_connected: true + local_request_entry_ready: true + singapore_to_fifth_domain_broadcast_bridge: live_verified + canger_email_delivery_live_verified: true + latest_test_workorder: bf72f16b-f81e-4de5-b0aa-aaee3d39e04c + latest_test_state: approved_not_claimed + automatic_session_activation_after_approval: not_yet_implemented +``` + +不得把第五域的 `approved` 邮件回执写成新加坡已经领取会话;只有取得领取回执并在 +家庭节点生成限时授权后,才能写成“本地权限已激活”。 diff --git a/remote-node/README.md b/remote-node/README.md new file mode 100644 index 0000000..6dfe32c --- /dev/null +++ b/remote-node/README.md @@ -0,0 +1,195 @@ +# 苍耳本地 Linux 节点接入 + +这个模块把苍耳家里的 Ubuntu 接到新加坡大脑服务器,但不把家里的电脑暴露到公网。 + +## 工作方式 + +1. 新加坡控制端运行在 `127.0.0.1:8787`,由现有 HTTPS 网关转发。 +2. 苍耳节点首次使用一次性配对码登记,长期节点令牌只保存在本机 `0600` 配置中。 +3. `systemd` 开机自动启动节点;连接只能从家里向新加坡主动发起。 +4. 人格体创建任务后,任务保持 `pending`;苍耳/冰朔在新加坡端批准后才会下发。 +5. 节点执行白名单动作并回传退出码、输出和耗时。 +6. 执行期间持续回传 `started`、分段输出和 `finished` 事件,新加坡端可以 + 实时观察,不必等任务结束才看见错误。 + +## 默认安全边界 + +- 家里路由器不做端口转发,不开放 SSH 到公网。 +- 控制端默认只监听环回地址,公网入口必须使用 HTTPS。 +- 不接受 Shell 字符串,只接受结构化动作。 +- 安装程序自动建立 `/srv/canger` 受管工作区,并用没有登录能力的 + `canger-agent` 系统账户执行任务。 +- 编程命令受可执行文件白名单、受管工作区、超时和授权时段共同限制; + Agent 没有 `sudo` 权限。 +- `repo.fetch` 只更新远端引用,不自动 merge、reset 或覆盖本地修改。 +- 每个任务和结果都写入 SQLite 审计记录。 + +## 新加坡端部署 + +准备 `/etc/canger-control.env`,四个令牌分别代表: + +- `ADMIN`:创建设备配对码; +- `OWNER`:苍耳主控批准、拒绝或撤销副控权限; +- `PERSONA`:苍耳人格体读取任务、结果和执行上下文; +- `OPERATOR`:冰朔与铸渊副控创建运维会话、任务并读取实时回执。 + +实际职责为: + +- 苍耳本人:服务器主控、需求方和最终决策者; +- 冰朔与 `ICE-GL-ZY001` 铸渊:受苍耳授权的运维副控; +- 每张工单的实际申请人格体:如实填写名称和正式编号; +- 本地 Agent:无独立判断权的受限物理执行器; +- 苍耳本人可随时批准、拒绝或撤销数小时副控运维会话。 + +安装 `systemd/canger-control.service`,并把 +`nginx/canger-node.conf.example` 接到已有 TLS 站点。不要直接把 8787 端口开放公网。 + +## 首次配对 + +新加坡端创建十分钟有效的一次性码: + +```bash +python3 tools/cangerctl.py pair-code +``` + +把整个 `remote-node` 目录交给苍耳。他在 Ubuntu 上只执行一次: + +```bash +sudo ./install/install-node.sh \ + https://guanghubingshuo.com/canger-node \ + XXXX-XXXX-XXXX \ + canger-home-ubuntu +``` + +冰朔交付给苍耳的专用包可直接使用简化入口: + +```bash +sudo ./install/苍耳一键安装.sh XXXX-XXXX-XXXX +``` + +配对码十分钟失效,因此安装包可以提前发送,配对码应在苍耳准备执行时再生成。 + +安装脚本会完成登记并启用开机自启动。以后不再需要登录本地终端。 + +安装时还会生成只读环境报告: + +```text +/var/lib/canger-node/preflight.json +``` + +报告会识别 systemd、虚拟化环境、磁盘文件系统和可能存在的 Windows +双启动,但不会修改分区、引导项或另一个操作系统。 + +## 双系统边界 + +真正的双系统切换会关闭当前操作系统。因此: + +- 启动 Ubuntu 时,本节点会自动运行并连接新加坡; +- 切到另一个操作系统后,Ubuntu 服务不可能继续运行; +- 若希望“电脑开机就总有一个节点在线”,另一个系统也需要安装配套节点; +- 若希望“Ubuntu 服务器本身始终在线”,需要让 Ubuntu 成为宿主机并把另一个 + 系统放进虚拟机,或者使用独立服务器。 + +安装程序不会自动改造分区、引导和虚拟机。这类操作可能导致数据丢失,必须在 +确认硬件、另一个系统和磁盘结构后单独处理。 + +## 从新加坡端操作 + +查看节点: + +```bash +python3 tools/cangerctl.py nodes +``` + +人格体创建任务: + +```bash +python3 tools/cangerctl.py task \ + --node node_xxx \ + --action system.status +``` + +正常开发不需要逐任务批准。人格体先申请 1–8 小时的授权时段: + +```bash +python3 tools/cangerctl.py grant \ + --node node_xxx \ + --hours 4 +``` + +系统自动读取节点登记的安全能力,向苍耳主控邮箱发送人话版批准链接。 +主人只看目的和时长后点一次批准,不需要填写命令、目录或权限。授权时段内, +匹配固定安全边界的任务自动批准。所有任务仍保留独立审计记录。 + +查看与撤销授权: + +```bash +python3 tools/cangerctl.py grants +python3 tools/cangerctl.py revoke grant_xxx +``` + +邮件中的链接只用于某一项授权申请,30 分钟失效;打开链接后还需要在确认页 +点击批准,避免邮件安全扫描器误触发。SMTP 密码、主人邮箱和签名密钥只放在 +新加坡服务器私有环境文件中,不写入仓库或安装包。 + +批准任务: + +```bash +python3 tools/cangerctl.py approve task_xxx +``` + +读取回执: + +```bash +python3 tools/cangerctl.py show task_xxx +``` + +实时观察执行过程: + +```bash +python3 tools/cangerctl.py watch task_xxx +``` + +## 苍耳人格体代码仓库入口 + +苍耳人格体不使用副控令牌。它从仓库固定门牌 +`persona/README.md` 进入,并使用独立的 `CANGER_PERSONA_TOKEN`: + +```bash +python3 tools/canger_persona.py request-access \ + --requester-name "申请者名称" \ + --requester-id "申请者正式编号" \ + --hours 4 \ + --reason "本次开发、诊断或修复目的" +``` + +申请回执的 `requester_name` 和 `requested_by` 必须与本次实际申请者一致。人格体可以申请、创建任务和读取结果, +但不能批准、拒绝或撤销授权;批准权继续属于苍耳主控。固定节点和安全边界登记在 +`config/pufferfish-persona.json`,令牌仍只存在新加坡私有环境文件中。 + +完整流程见 `persona/WORKORDER-GUIDE.md`;当天开发回执见 +`persona/DEVELOPMENT-RECEIPT-20260729.md`。 + +## 智能体与模型 API + +本地 Agent 是执行、观察和审计端,不直接保存大模型供应商的 API 密钥。 +运维判断由冰朔与铸渊副控在新加坡大脑服务器完成: + +```text +苍耳提出需求并授权 + ↓ +冰朔与铸渊副控 → 创建运维任务 → 苍耳人格体/本地 Agent 协作执行 + ← 实时事件、错误和回执 ← +``` + +如需调用大模型,只把供应商 API 密钥放在新加坡服务器的密钥环境中。本地节点 +只持有可撤销、仅能代表该节点的身份令牌。禁止把大模型 API 密钥写进安装包、 +仓库、命令行或配对文件。 + +## 动作 + +- `system.status`:主机、系统、在线时间、根分区状态。 +- `repo.status`:只读 Git 状态。 +- `repo.fetch`:获取远端引用,不合并。 +- `service.logs`:读取指定 systemd 服务日志。 +- `command.run`:结构化命令执行;只在受管工作区内、以受限系统账户运行。 diff --git a/remote-node/config/control.env.example b/remote-node/config/control.env.example new file mode 100644 index 0000000..52b03e7 --- /dev/null +++ b/remote-node/config/control.env.example @@ -0,0 +1,20 @@ +CANGER_BIND=127.0.0.1 +CANGER_PORT=8787 +CANGER_CONTROL_DB=/var/lib/canger-control/control.sqlite3 +CANGER_ADMIN_TOKEN=replace-with-at-least-24-random-characters +CANGER_OWNER_TOKEN=replace-with-at-least-24-random-characters +CANGER_PERSONA_TOKEN=replace-with-at-least-24-random-characters +CANGER_OPERATOR_TOKEN=replace-with-at-least-24-random-characters +CANGER_OPERATOR_ID=ICE-GL-ZY001 +CANGER_OWNER_ID=CANGER-HUMAN-OWNER +CANGER_EXECUTOR_PERSONA_ID=ICE-GL-CA001 +CANGER_APPROVAL_SIGNING_KEY=replace-with-at-least-32-random-characters +CANGER_OWNER_EMAIL=owner@example.invalid +CANGER_PUBLIC_URL=https://guanghulab.com/canger-node +# SMTP credentials stay only on the Singapore server. Never put them in Git. +CANGER_SMTP_HOST=smtp.example.invalid +CANGER_SMTP_PORT=465 +CANGER_SMTP_FROM=guanghu-router@example.invalid +CANGER_SMTP_USERNAME=replace-on-server +CANGER_SMTP_PASSWORD=replace-on-server +CANGER_CONTROL_URL=http://127.0.0.1:8787 diff --git a/remote-node/config/node.example.json b/remote-node/config/node.example.json new file mode 100644 index 0000000..3155ee2 --- /dev/null +++ b/remote-node/config/node.example.json @@ -0,0 +1,26 @@ +{ + "server_url": "https://guanghulab.com/canger-node", + "node_name": "canger-home-ubuntu", + "node_id": null, + "node_token": null, + "allowed_roots": [ + "/srv/canger" + ], + "enabled_actions": [ + "system.status", + "repo.status", + "repo.fetch", + "service.logs", + "command.run" + ], + "executable_allowlist": [ + "git", + "python3", + "pytest", + "node", + "npm", + "pnpm" + ], + "poll_seconds": 5, + "verify_tls": true +} diff --git a/remote-node/config/pufferfish-persona.json b/remote-node/config/pufferfish-persona.json new file mode 100644 index 0000000..e9c1fb4 --- /dev/null +++ b/remote-node/config/pufferfish-persona.json @@ -0,0 +1,14 @@ +{ + "profile": "pufferfish-persona-operations/v0.1", + "subsystem_id": "SYS-CE", + "requester_identity_mode": "explicit_name_and_registered_id_per_request", + "human_owner": "苍耳", + "node_id": "node_1cfab77579a3e9018f7a9037", + "node_name": "canger-home-ubuntu-live-20260729", + "control_url": "http://127.0.0.1:8797", + "default_hours": 4, + "minimum_hours": 1, + "maximum_hours": 8, + "approval_authority": "CANGER-HUMAN-OWNER", + "managed_root": "/srv/canger" +} diff --git a/remote-node/install/install-node.sh b/remote-node/install/install-node.sh new file mode 100755 index 0000000..f02679c --- /dev/null +++ b/remote-node/install/install-node.sh @@ -0,0 +1,107 @@ +#!/usr/bin/env bash +set -euo pipefail + +if [[ "${EUID}" -ne 0 ]]; then + echo "请用 sudo 执行此脚本。" >&2 + exit 1 +fi + +if [[ "$#" -ne 3 ]]; then + echo "用法: $0 " >&2 + exit 1 +fi + +server_url="$1" +pair_code="$2" +node_name="$3" + +if ! id canger-agent >/dev/null 2>&1; then + useradd --system --home-dir /var/lib/canger-node --create-home \ + --shell /usr/sbin/nologin canger-agent +fi + +install -d -m 0755 /opt/canger-node +install -d -o canger-agent -g canger-agent -m 0700 /etc/canger-node +install -d -o canger-agent -g canger-agent -m 0750 /var/lib/canger-node +install -d -o canger-agent -g canger-agent -m 0750 /srv/canger +install -m 0755 "$(dirname "$0")/../node/canger_node.py" /opt/canger-node/canger_node.py +install -m 0755 "$(dirname "$0")/preflight-linux.py" /opt/canger-node/preflight-linux.py +install -m 0644 "$(dirname "$0")/../systemd/canger-node.service" /etc/systemd/system/canger-node.service + +already_paired=false +if [[ -f /etc/canger-node/config.json ]] && \ + runuser -u canger-agent -- python3 - /etc/canger-node/config.json <<'PY' +import json +import sys + +with open(sys.argv[1], encoding="utf-8") as handle: + config = json.load(handle) +raise SystemExit(0 if config.get("node_id") and config.get("node_token") else 1) +PY +then + already_paired=true +fi + +if [[ "${already_paired}" == "false" ]]; then + python3 - "$server_url" "$node_name" <<'PY' +import json +import os +import pwd +import sys + +server_url, node_name = sys.argv[1:] +config = { + "server_url": server_url, + "node_name": node_name, + "node_id": None, + "node_token": None, + "allowed_roots": ["/srv/canger"], + "enabled_actions": [ + "system.status", + "repo.status", + "repo.fetch", + "service.logs", + "command.run", + ], + "executable_allowlist": ["git", "python3", "pytest", "node", "npm", "pnpm"], + "poll_seconds": 5, + "verify_tls": True, +} +path = "/etc/canger-node/config.json" +with open(path, "w", encoding="utf-8") as handle: + json.dump(config, handle, ensure_ascii=False, indent=2) + handle.write("\n") +os.chmod(path, 0o600) +account = pwd.getpwnam("canger-agent") +os.chown(path, account.pw_uid, account.pw_gid) +PY +else + echo "检测到现有有效节点身份,保留原身份并跳过重复配对。" +fi + +/usr/bin/python3 /opt/canger-node/preflight-linux.py \ + --output /var/lib/canger-node/preflight.json + +install -d -m 0755 /etc/systemd/system/canger-node.service.d +python3 - <<'PY' +import os + +path = "/etc/systemd/system/canger-node.service.d/paths.conf" +with open(path, "w", encoding="utf-8") as handle: + handle.write("[Service]\n") + handle.write("ReadWritePaths=/etc/canger-node\n") + handle.write("ReadWritePaths=/srv/canger\n") +os.chmod(path, 0o644) +PY + +if [[ "${already_paired}" == "false" ]]; then + runuser -u canger-agent -- /usr/bin/python3 /opt/canger-node/canger_node.py \ + --config /etc/canger-node/config.json pair --code "$pair_code" +fi + +systemctl daemon-reload +systemctl enable --now canger-node.service +systemctl --no-pager --full status canger-node.service + +echo "安装完成。今后开机自动启动,断网后自动重连。" +echo "环境检测报告: /var/lib/canger-node/preflight.json" diff --git a/remote-node/install/preflight-linux.py b/remote-node/install/preflight-linux.py new file mode 100755 index 0000000..bd5bb71 --- /dev/null +++ b/remote-node/install/preflight-linux.py @@ -0,0 +1,85 @@ +#!/usr/bin/env python3 +"""Read-only environment report used before installing the persistent node.""" + +from __future__ import annotations + +import argparse +import json +import os +import platform +import shutil +import subprocess +from pathlib import Path +from typing import Any + + +def command(argv: list[str], timeout: int = 20) -> dict[str, Any]: + if shutil.which(argv[0]) is None: + return {"available": False} + try: + result = subprocess.run( + argv, + stdin=subprocess.DEVNULL, + stdout=subprocess.PIPE, + stderr=subprocess.STDOUT, + text=True, + timeout=timeout, + shell=False, + ) + return { + "available": True, + "exit_code": result.returncode, + "output": result.stdout[:32_000], + } + except subprocess.TimeoutExpired: + return {"available": True, "error": "timeout"} + + +def inspect() -> dict[str, Any]: + efi = command(["efibootmgr", "-v"]) + disks = command(["lsblk", "-J", "-o", "NAME,FSTYPE,LABEL,MOUNTPOINTS"]) + virt = command(["systemd-detect-virt"]) + boot_text = str(efi.get("output", "")).lower() + disk_text = str(disks.get("output", "")).lower() + windows_hint = any( + marker in boot_text or marker in disk_text + for marker in ("windows boot manager", "microsoft", "ntfs", "bitlocker") + ) + return { + "schema": "canger-node-preflight/v1", + "hostname": platform.node(), + "platform": platform.platform(), + "machine": platform.machine(), + "systemd_running": os.path.isdir("/run/systemd/system"), + "virtualization": virt, + "efi_boot_entries": efi, + "block_filesystems": disks, + "possible_windows_dual_boot": windows_hint, + "continuity": { + "while_linux_is_booted": "supported", + "after_reboot_back_into_linux": "supported", + "while_another_dual_boot_os_is_active": "not_possible_from_linux", + }, + "note": ( + "A dual-boot switch powers off Linux. Cross-OS continuity requires " + "a companion agent in the other OS, a virtual machine, or separate hardware." + ), + } + + +def main() -> None: + parser = argparse.ArgumentParser() + parser.add_argument("--output") + args = parser.parse_args() + report = inspect() + text = json.dumps(report, ensure_ascii=False, indent=2) + "\n" + if args.output: + target = Path(args.output) + target.parent.mkdir(parents=True, exist_ok=True) + target.write_text(text, encoding="utf-8") + os.chmod(target, 0o600) + print(text, end="") + + +if __name__ == "__main__": + main() diff --git a/remote-node/install/苍耳一键安装.sh b/remote-node/install/苍耳一键安装.sh new file mode 100755 index 0000000..67e1889 --- /dev/null +++ b/remote-node/install/苍耳一键安装.sh @@ -0,0 +1,18 @@ +#!/usr/bin/env bash +set -euo pipefail + +script_dir="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" + +if [[ "${EUID}" -ne 0 ]]; then + exec sudo "$0" "$@" +fi + +pair_code="${1:-}" +if [[ -z "${pair_code}" ]]; then + read -r -p "请输入冰朔发来的十分钟配对码: " pair_code +fi + +exec bash "${script_dir}/install-node.sh" \ + "https://guanghubingshuo.com/canger-node" \ + "${pair_code}" \ + "canger-home-ubuntu" diff --git a/remote-node/nginx/canger-node.conf.example b/remote-node/nginx/canger-node.conf.example new file mode 100644 index 0000000..e91a0d9 --- /dev/null +++ b/remote-node/nginx/canger-node.conf.example @@ -0,0 +1,10 @@ +# Mount below an existing TLS-enabled server block. +# The Python control plane intentionally listens on loopback only. +location /canger-node/ { + proxy_pass http://127.0.0.1:8787/; + proxy_http_version 1.1; + proxy_set_header Host $host; + proxy_set_header X-Forwarded-Proto https; + proxy_read_timeout 45s; + client_max_body_size 512k; +} diff --git a/remote-node/node/canger_node.py b/remote-node/node/canger_node.py new file mode 100755 index 0000000..49239d2 --- /dev/null +++ b/remote-node/node/canger_node.py @@ -0,0 +1,408 @@ +#!/usr/bin/env python3 +"""Outbound-only Linux node agent for Cang'er's home server.""" + +from __future__ import annotations + +import argparse +import json +import os +import platform +import shutil +import socket +import ssl +import subprocess +import sys +import tempfile +import threading +import time +from queue import Empty, Queue +from pathlib import Path +from typing import Any +from urllib.error import HTTPError, URLError +from urllib.request import Request, urlopen + +AGENT_VERSION = "0.1.0" +MAX_OUTPUT_CHARS = 240_000 + + +def read_config(path: str) -> dict[str, Any]: + with open(path, encoding="utf-8") as handle: + value = json.load(handle) + if not isinstance(value, dict): + raise ValueError("config must be a JSON object") + return value + + +def write_config(path: str, value: dict[str, Any]) -> None: + target = Path(path) + target.parent.mkdir(parents=True, exist_ok=True) + fd, temp_name = tempfile.mkstemp(prefix=".canger-node-", dir=target.parent) + try: + os.fchmod(fd, 0o600) + with os.fdopen(fd, "w", encoding="utf-8") as handle: + json.dump(value, handle, ensure_ascii=False, indent=2) + handle.write("\n") + os.replace(temp_name, target) + finally: + if os.path.exists(temp_name): + os.unlink(temp_name) + + +def ssl_context(config: dict[str, Any]) -> ssl.SSLContext: + context = ssl.create_default_context() + if config.get("verify_tls", True) is False: + if not str(config.get("server_url", "")).startswith( + ("http://127.0.0.1", "http://localhost") + ): + raise ValueError("TLS verification may only be disabled for localhost") + context.check_hostname = False + context.verify_mode = ssl.CERT_NONE + return context + + +def api( + config: dict[str, Any], + method: str, + path: str, + body: dict[str, Any] | None = None, + token: str | None = None, +) -> dict[str, Any]: + base = str(config["server_url"]).rstrip("/") + if not base.startswith(("https://", "http://127.0.0.1", "http://localhost")): + raise ValueError("server_url must use HTTPS (HTTP is only allowed on localhost)") + data = None if body is None else json.dumps(body).encode("utf-8") + headers = {"Accept": "application/json"} + if data is not None: + headers["Content-Type"] = "application/json" + if token: + headers["Authorization"] = f"Bearer {token}" + request = Request(base + path, data=data, headers=headers, method=method) + try: + with urlopen(request, timeout=30, context=ssl_context(config)) as response: + value = json.load(response) + except HTTPError as exc: + detail = exc.read().decode("utf-8", errors="replace") + raise RuntimeError(f"server returned HTTP {exc.code}: {detail}") from exc + if not isinstance(value, dict): + raise RuntimeError("server returned a non-object response") + return value + + +def allowed_path(config: dict[str, Any], raw: str) -> Path: + candidate = Path(raw).expanduser().resolve() + roots = [Path(item).expanduser().resolve() for item in config["allowed_roots"]] + if not any(candidate == root or root in candidate.parents for root in roots): + raise PermissionError(f"path is outside allowed_roots: {candidate}") + return candidate + + +def run_process( + argv: list[str], + cwd: Path | None = None, + timeout: int = 120, + event_cb: Any = None, +) -> dict[str, Any]: + started = time.monotonic() + process = subprocess.Popen( + argv, + cwd=str(cwd) if cwd else None, + stdin=subprocess.DEVNULL, + stdout=subprocess.PIPE, + stderr=subprocess.STDOUT, + text=True, + bufsize=1, + shell=False, + env={ + "HOME": os.environ.get("HOME", "/root"), + "LANG": os.environ.get("LANG", "C.UTF-8"), + "PATH": "/usr/local/sbin:/usr/local/bin:/usr/sbin:/usr/bin:/sbin:/bin", + }, + ) + lines: Queue[str | None] = Queue() + + def reader() -> None: + assert process.stdout is not None + for line in process.stdout: + lines.put(line) + lines.put(None) + + threading.Thread(target=reader, daemon=True).start() + output_parts: list[str] = [] + output_size = 0 + batch: list[str] = [] + last_emit = time.monotonic() + ended = False + try: + while not ended: + if time.monotonic() - started > timeout: + process.terminate() + try: + process.wait(timeout=5) + except subprocess.TimeoutExpired: + process.kill() + process.wait() + raise subprocess.TimeoutExpired(argv, timeout) + try: + line = lines.get(timeout=0.5) + if line is None: + ended = True + else: + if output_size < MAX_OUTPUT_CHARS: + remaining = MAX_OUTPUT_CHARS - output_size + kept = line[:remaining] + output_parts.append(kept) + output_size += len(kept) + batch.append(line) + except Empty: + pass + if batch and ( + sum(map(len, batch)) >= 8192 + or time.monotonic() - last_emit >= 1.0 + or ended + ): + if event_cb: + event_cb( + "output", + { + "text": "".join(batch)[:16_000], + "stream": "combined", + }, + ) + batch.clear() + last_emit = time.monotonic() + return_code = process.wait() + output = "".join(output_parts) + return { + "ok": return_code == 0, + "exit_code": return_code, + "duration_ms": round((time.monotonic() - started) * 1000), + "output": output, + "output_truncated": output_size >= MAX_OUTPUT_CHARS, + } + except subprocess.TimeoutExpired as exc: + output = "".join(output_parts) + return { + "ok": False, + "exit_code": None, + "duration_ms": round((time.monotonic() - started) * 1000), + "output": output, + "output_truncated": output_size >= MAX_OUTPUT_CHARS, + "error": "timeout", + } + + +def execute( + config: dict[str, Any], task: dict[str, Any], event_cb: Any = None +) -> dict[str, Any]: + action = task["action"] + args = task.get("args") or {} + enabled = set(config.get("enabled_actions", [])) + if action not in enabled: + raise PermissionError(f"action is disabled on this node: {action}") + + if action == "system.status": + commands = [ + ["hostname"], + ["uname", "-a"], + ["uptime"], + ["df", "-h", "/"], + ] + results = [ + run_process(command, timeout=30, event_cb=event_cb) for command in commands + ] + return { + "ok": all(item["ok"] for item in results), + "hostname": socket.gethostname(), + "platform": platform.platform(), + "commands": results, + } + + if action == "repo.status": + repo = allowed_path(config, args["repo"]) + return run_process( + [ + "git", + "-C", + str(repo), + "status", + "--short", + "--branch", + ], + timeout=60, + event_cb=event_cb, + ) + + if action == "repo.fetch": + repo = allowed_path(config, args["repo"]) + return run_process( + ["git", "-C", str(repo), "fetch", "--prune", "origin"], + timeout=300, + event_cb=event_cb, + ) + + if action == "service.logs": + return run_process( + [ + "journalctl", + "--no-pager", + "-u", + args["service"], + "-n", + str(args.get("lines", 200)), + ], + timeout=60, + event_cb=event_cb, + ) + + if action == "command.run": + argv = args["argv"] + executable = argv[0] + if "/" in executable or executable not in set( + config.get("executable_allowlist", []) + ): + raise PermissionError(f"executable is not allowlisted: {executable}") + if shutil.which(executable) is None: + raise FileNotFoundError(f"executable is not installed: {executable}") + cwd = allowed_path(config, args["cwd"]) + return run_process( + argv, + cwd=cwd, + timeout=int(args.get("timeout_seconds", 300)), + event_cb=event_cb, + ) + + raise ValueError(f"unsupported action: {action}") + + +def pair(config_path: str, code: str) -> None: + config = read_config(config_path) + if config.get("node_id") or config.get("node_token"): + raise RuntimeError("node is already paired; remove its identity explicitly first") + response = api( + config, + "POST", + "/v1/pairings/claim", + { + "code": code, + "node_name": config["node_name"], + "agent_version": AGENT_VERSION, + "policy": { + "allowed_roots": config["allowed_roots"], + "enabled_actions": config["enabled_actions"], + }, + }, + ) + config["node_id"] = response["node_id"] + config["node_token"] = response["node_token"] + write_config(config_path, config) + print(f"paired node {config['node_name']} as {config['node_id']}") + + +def run(config_path: str, once: bool = False) -> None: + backoff = 2 + while True: + config = read_config(config_path) + node_id = config.get("node_id") + node_token = config.get("node_token") + if not node_id or not node_token: + raise RuntimeError("node is not paired") + try: + api( + config, + "POST", + f"/v1/nodes/{node_id}/heartbeat", + { + "agent_version": AGENT_VERSION, + "policy": { + "allowed_roots": config["allowed_roots"], + "enabled_actions": config["enabled_actions"], + }, + }, + node_token, + ) + response = api( + config, + "GET", + f"/v1/nodes/{node_id}/tasks/next", + token=node_token, + ) + task = response.get("task") + if task is not None: + def emit(kind: str, payload: dict[str, Any]) -> None: + try: + api( + config, + "POST", + f"/v1/tasks/{task['id']}/events", + {"kind": kind, "payload": payload}, + node_token, + ) + except Exception as exc: + print(f"event upload failed: {exc}", file=sys.stderr) + + emit( + "started", + { + "action": task["action"], + "hostname": socket.gethostname(), + "agent_version": AGENT_VERSION, + }, + ) + try: + result = execute(config, task, emit) + except Exception as exc: + result = { + "ok": False, + "error": type(exc).__name__, + "message": str(exc), + } + emit( + "finished", + { + "ok": bool(result.get("ok")), + "exit_code": result.get("exit_code"), + "error": result.get("error"), + }, + ) + api( + config, + "POST", + f"/v1/tasks/{task['id']}/result", + result, + node_token, + ) + backoff = 2 + if once: + return + time.sleep(float(config.get("poll_seconds", 5))) + except (OSError, URLError, RuntimeError) as exc: + print(f"connection error: {exc}; retrying in {backoff}s", file=sys.stderr) + if once: + raise + time.sleep(backoff) + backoff = min(backoff * 2, 120) + + +def parse_args() -> argparse.Namespace: + parser = argparse.ArgumentParser() + parser.add_argument( + "--config", default="/etc/canger-node/config.json", help="node config path" + ) + sub = parser.add_subparsers(dest="command", required=True) + pair_parser = sub.add_parser("pair") + pair_parser.add_argument("--code", required=True) + run_parser = sub.add_parser("run") + run_parser.add_argument("--once", action="store_true") + return parser.parse_args() + + +def main() -> None: + args = parse_args() + if args.command == "pair": + pair(args.config, args.code) + else: + run(args.config, args.once) + + +if __name__ == "__main__": + main() diff --git a/remote-node/persona/DEVELOPMENT-RECEIPT-20260729.md b/remote-node/persona/DEVELOPMENT-RECEIPT-20260729.md new file mode 100644 index 0000000..e584cdf --- /dev/null +++ b/remote-node/persona/DEVELOPMENT-RECEIPT-20260729.md @@ -0,0 +1,57 @@ +# 苍耳家庭 Ubuntu 运维节点 · 2026-07-29 开发回执 + +## 已完成 + +- 苍耳家庭 Ubuntu 以只出站方式接入新加坡大脑服务器; +- 节点已配对并由 systemd 开机自动启动、断网自动重连; +- 新加坡控制面能看到节点心跳并下发受限结构化任务; +- 人格体可从 `remote-node/persona/README.md` 发起限时申请; +- 申请者改为逐单填写名称与正式编号,不再固定为某个人格体; +- 授权与申请者编号绑定,另一人格体不能借用; +- 苍耳仍是唯一最终批准者。 + +## 今日故障与修复 + +- 一键安装脚本曾因执行位丢失而失败,已改为由入口显式调用安装脚本; +- 旧配对码过期、节点名称重复、系统配置目录权限不足均已分别处理; +- JZAO 移动硬盘是 ExFAT,不能保留 Linux 所需的 POSIX 权限,曾造成新加坡服务读取 + 源码失败;以后源码必须先经过 JZAO 源码运输闸门标准化,再传 Linux; +- 测试工单曾把申请人写死,现已改成通用申请者字段。 + +## 当前状态 + +```yaml +home_node: connected +control_service: active +request_entry: implemented +requester_identity: explicit_per_workorder +owner_approval: required +local_direct_smtp: deprecated +singapore_broadcast_tower_to_fifth_domain: live_verified +fifth_domain_email_delivery_to_canger: sent_and_owner_approved +fifth_domain_owner_decision: approved +singapore_claim_state: not_claimed +home_node_permission: not_activated +``` + +## 下一步 + +2026-07-29 已由新加坡广播塔向第五域可信广播入口提交送达测试: + +```yaml +workorder_id: bf72f16b-f81e-4de5-b0aa-aaee3d39e04c +target: CANGER-HOME-UBUNTU-001 +approver: TCS-GL-009 +email_status: sent +state: approved +claimed: false +local_permission_activated: false +task_execution_requested: false +``` + +苍耳已点击批准,第五域于现场读回 `state: approved`、`claimed: false`。本次测试 +刻意没有在新加坡保存领取凭证,所以没有激活家庭节点权限,也没有执行任务。这证明 +“申请 → 新加坡广播塔 → 第五域广播主控塔 → 苍耳邮箱 → 苍耳批准”已经贯通。 + +尚未完成的是正式工单的最后闭环:新加坡广播塔需要安全保存一次性领取凭证,在苍耳 +批准后领取回执,再把同一申请者、同一范围和有效期转换成家庭节点的限时会话。 diff --git a/remote-node/persona/README.md b/remote-node/persona/README.md new file mode 100644 index 0000000..3a8c831 --- /dev/null +++ b/remote-node/persona/README.md @@ -0,0 +1,69 @@ +# 苍耳代码仓库 · 服务器运维申请入口 + +本目录是任何已登记人格体从新加坡 `cang-ying` 代码仓库进入苍耳家庭 +Ubuntu 运维流程的固定门牌。申请人不写死:谁申请,就如实填写谁的名称和正式编号。 + +## 主权边界 + +- 人格体可以发起 1–8 小时的运维申请; +- 人格体不能批准、拒绝或撤销申请; +- 申请范围由家庭节点上报的安全策略决定,不能由人格体扩大; +- 苍耳仍是最终批准者; +- 未取得授权时,任务保持等待,不会下发到家庭 Ubuntu; +- 仓库不保存令牌、批准链接、密码或配对码。 + +## 发起申请 + +在新加坡仓库根目录运行: + +```bash +set -a +. /etc/canger-control.env +set +a +python3 remote-node/tools/canger_persona.py \ + --requester-name "本次申请者名称" \ + --requester-id "本次申请者正式编号" \ + request-access \ + --hours 4 \ + --reason "说明本次要检查、开发或修复什么" +``` + +成功回执必须同时满足: + +- `status` 为 `pending`; +- `requester_name` 是本次申请者名称; +- `requested_by` 是同一个申请者的正式编号; +- `duration_seconds` 与申请时长一致; +- `scopes` 和 `roots` 来自节点安全策略。 + +这只表示“申请已经送达”,不表示苍耳已经批准。 + +## 授权后创建任务 + +```bash +python3 remote-node/tools/canger_persona.py \ + --requester-name "与获批工单相同的名称" \ + --requester-id "与获批工单相同的正式编号" \ + task \ + --action system.status +``` + +在有效授权时段内,符合范围的任务会返回 `status: approved`;不符合范围或没有授权时,任务不会被家庭节点执行。 + +## 查看结果 + +```bash +python3 remote-node/tools/canger_persona.py show task_xxx +python3 remote-node/tools/canger_persona.py watch task_xxx +``` + +## 固定配置 + +节点和控制端路径登记在: + +`remote-node/config/pufferfish-persona.json` + +申请者身份不保存在这个固定配置里,而是每张工单明确填写。需要变更节点时先更新 +胖头鱼节点注册表并重新验证心跳,不能让申请者在运行时自行替换节点。 + +完整操作方法与状态边界见 `WORKORDER-GUIDE.md`。 diff --git a/remote-node/persona/WORKORDER-GUIDE.md b/remote-node/persona/WORKORDER-GUIDE.md new file mode 100644 index 0000000..7317dc7 --- /dev/null +++ b/remote-node/persona/WORKORDER-GUIDE.md @@ -0,0 +1,77 @@ +# 苍耳服务器工单操作方法 + +## 一张工单只回答四件事 + +1. 你是谁; +2. 你的正式编号是什么; +3. 你申请做什么; +4. 希望打开多长时间的受控会话。 + +申请者可以是耳耳蛋、鉴影、铸渊或以后登记的其他人格体,系统不写死任何名字。 +申请者只能建单,不能批准自己的申请。苍耳收到授权通知后,自行决定是否批准。 + +## 发起申请 + +在新加坡 `cang-ying` 仓库根目录执行: + +```bash +set -a +. /etc/canger-control.env +set +a +python3 remote-node/tools/canger_persona.py \ + --requester-name "申请者名称" \ + --requester-id "申请者正式编号" \ + request-access \ + --hours 3 \ + --reason "本次准备检查、开发、调试或修复的目的" +``` + +返回 `pending` 只代表申请已登记。没有苍耳批准,家庭 Ubuntu 不会执行任务。 + +## 批准后的任务 + +任务必须继续使用获批工单里的同一名称和编号: + +```bash +python3 remote-node/tools/canger_persona.py \ + --requester-name "申请者名称" \ + --requester-id "申请者正式编号" \ + task \ + --action system.status +``` + +其他可用动作及范围由家庭节点上报的安全策略决定。不同编号不能借用别人的授权。 + +## 广播链 + +目标链路是: + +```text +申请者 +→ 新加坡语言主控广播塔 +→ 第五域语言主控广播塔 +→ 按苍耳节点登记选择苍耳审批频道 +→ 苍耳邮箱授权链接 +→ 苍耳批准 +→ 新加坡控制端激活同一申请者的限时会话 +``` + +请求正文不能携带或更换收件邮箱。苍耳邮箱只登记在第五域主控服务器的私密审批人 +注册表中。广播链未返回 +`owner_email_sent_by_registered_broadcast_tower` 前,不得声称邮件已经送达。 + +2026-07-29 的送达测试已返回该回执;工单 +`bf72f16b-f81e-4de5-b0aa-aaee3d39e04c` 仅用于测试邮件,不会执行服务器任务。 +苍耳随后已点击批准,第五域状态为 `approved`,但该测试单的领取凭证被刻意丢弃, +因此状态仍为 `claimed: false`,新加坡和家庭节点都没有获得权限。这是测试预期, +不是节点离线。 + +## 安全边界 + +- 仓库不保存邮箱地址、SMTP 授权码、批准令牌、会话令牌或节点令牌; +- 申请者名称和编号必须如实填写; +- 苍耳没有点击批准时,权限保持关闭; +- 第五域显示批准但新加坡尚未领取时,权限同样保持关闭; +- 会话只能使用家庭节点上报的能力和受管目录; +- 切换申请者、节点或扩大范围必须重新申请; +- 每个任务都保留独立回执。 diff --git a/remote-node/server/control_plane.py b/remote-node/server/control_plane.py new file mode 100755 index 0000000..f334278 --- /dev/null +++ b/remote-node/server/control_plane.py @@ -0,0 +1,978 @@ +#!/usr/bin/env python3 +"""Singapore-side control plane for the Cang'er outbound node.""" + +from __future__ import annotations + +import argparse +import base64 +import hashlib +import hmac +import html +import json +import os +import re +import smtplib +import sqlite3 +import sys +from email.message import EmailMessage +from http import HTTPStatus +from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer +from pathlib import Path +from typing import Any +from urllib.parse import parse_qs, urlencode, urlparse + +sys.path.insert(0, str(Path(__file__).resolve().parents[1])) + +from shared import ( # noqa: E402 + MAX_JSON_BYTES, + MAX_RESULT_BYTES, + PAIRING_TTL_SECONDS, + TASK_LEASE_SECONDS, + compact_json, + new_id, + new_pairing_code, + new_token, + now_ts, + secret_hash, + secret_matches, + validate_node_name, + validate_task, +) + + +SCHEMA = """ +CREATE TABLE IF NOT EXISTS pairings ( + id TEXT PRIMARY KEY, + code_hash TEXT NOT NULL, + expires_at INTEGER NOT NULL, + used_at INTEGER +); +CREATE TABLE IF NOT EXISTS nodes ( + id TEXT PRIMARY KEY, + name TEXT NOT NULL UNIQUE, + token_hash TEXT NOT NULL, + created_at INTEGER NOT NULL, + last_seen_at INTEGER, + agent_version TEXT, + policy_json TEXT NOT NULL DEFAULT '{}' +); +CREATE TABLE IF NOT EXISTS tasks ( + id TEXT PRIMARY KEY, + node_id TEXT NOT NULL REFERENCES nodes(id), + action TEXT NOT NULL, + args_json TEXT NOT NULL, + status TEXT NOT NULL, + created_at INTEGER NOT NULL, + approved_at INTEGER, + leased_until INTEGER, + completed_at INTEGER, + result_json TEXT, + grant_id TEXT, + requested_by TEXT NOT NULL DEFAULT 'ICE-GL-ZY001', + executor_persona TEXT NOT NULL DEFAULT 'ICE-GL-CA001', + FOREIGN KEY(node_id) REFERENCES nodes(id) +); +CREATE INDEX IF NOT EXISTS idx_tasks_node_status +ON tasks(node_id, status, created_at); +CREATE TABLE IF NOT EXISTS task_events ( + seq INTEGER PRIMARY KEY AUTOINCREMENT, + task_id TEXT NOT NULL REFERENCES tasks(id), + created_at INTEGER NOT NULL, + kind TEXT NOT NULL, + payload_json TEXT NOT NULL +); +CREATE INDEX IF NOT EXISTS idx_task_events_task_seq +ON task_events(task_id, seq); +CREATE TABLE IF NOT EXISTS grants ( + id TEXT PRIMARY KEY, + node_id TEXT NOT NULL REFERENCES nodes(id), + scopes_json TEXT NOT NULL, + roots_json TEXT NOT NULL, + reason TEXT NOT NULL, + duration_seconds INTEGER NOT NULL, + status TEXT NOT NULL, + created_at INTEGER NOT NULL, + approved_at INTEGER, + expires_at INTEGER, + revoked_at INTEGER + ,requested_by TEXT NOT NULL DEFAULT 'UNKNOWN' + ,requester_name TEXT NOT NULL DEFAULT '' +); +CREATE INDEX IF NOT EXISTS idx_grants_node_status +ON grants(node_id, status, expires_at); +""" + + +class ControlPlane: + def __init__( + self, + db_path: str, + admin_token: str, + owner_token: str, + persona_token: str, + operator_token: str, + approval_signing_key: str, + owner_email: str = "", + public_url: str = "", + ) -> None: + for name, value in { + "admin token": admin_token, + "owner token": owner_token, + "persona token": persona_token, + "operator token": operator_token, + }.items(): + if len(value) < 24: + raise ValueError(f"{name} must contain at least 24 characters") + if len(approval_signing_key) < 32: + raise ValueError("approval signing key must contain at least 32 characters") + self.db_path = db_path + self.admin_hash = secret_hash(admin_token) + self.owner_hash = secret_hash(owner_token) + self.persona_hash = secret_hash(persona_token) + self.operator_hash = secret_hash(operator_token) + self.operator_id = os.environ.get("CANGER_OPERATOR_ID", "ICE-GL-ZY001") + self.owner_id = os.environ.get("CANGER_OWNER_ID", "CANGER-HUMAN-OWNER") + self.executor_persona_id = os.environ.get( + "CANGER_EXECUTOR_PERSONA_ID", "ICE-GL-CA001" + ) + self.approval_signing_key = approval_signing_key.encode("utf-8") + self.owner_email = owner_email + self.public_url = public_url.rstrip("/") + Path(db_path).parent.mkdir(parents=True, exist_ok=True) + with self.db() as conn: + conn.executescript(SCHEMA) + columns = { + row["name"] + for row in conn.execute("PRAGMA table_info(grants)").fetchall() + } + if "requester_name" not in columns: + conn.execute( + "ALTER TABLE grants ADD COLUMN requester_name " + "TEXT NOT NULL DEFAULT ''" + ) + + def db(self) -> sqlite3.Connection: + conn = sqlite3.connect(self.db_path, timeout=10, isolation_level=None) + conn.row_factory = sqlite3.Row + conn.execute("PRAGMA foreign_keys=ON") + conn.execute("PRAGMA journal_mode=WAL") + return conn + + @staticmethod + def bearer(headers: Any) -> str: + value = headers.get("Authorization", "") + if not value.startswith("Bearer "): + return "" + return value[7:].strip() + + def role_allowed(self, headers: Any, *roles: str) -> bool: + token = self.bearer(headers) + hashes = { + "admin": self.admin_hash, + "owner": self.owner_hash, + "persona": self.persona_hash, + "operator": self.operator_hash, + } + return any(secret_matches(token, hashes[role]) for role in roles) + + def node_for_token(self, node_id: str, headers: Any) -> sqlite3.Row | None: + token = self.bearer(headers) + with self.db() as conn: + row = conn.execute( + "SELECT * FROM nodes WHERE id=?", (node_id,) + ).fetchone() + if row is None or not secret_matches(token, row["token_hash"]): + return None + return row + + def grant_token(self, grant_id: str, link_expires_at: int) -> str: + message = f"canger-grant/v1\n{grant_id}\n{link_expires_at}".encode("utf-8") + signature = hmac.new( + self.approval_signing_key, message, hashlib.sha256 + ).digest() + encoded = base64.urlsafe_b64encode(signature).decode("ascii").rstrip("=") + return f"{link_expires_at}.{encoded}" + + def verify_grant_token(self, grant_id: str, token: str) -> bool: + try: + raw_expiry, _ = token.split(".", 1) + expiry = int(raw_expiry) + except (ValueError, TypeError): + return False + if expiry < now_ts(): + return False + return hmac.compare_digest(token, self.grant_token(grant_id, expiry)) + + def send_grant_email(self, grant: dict[str, Any]) -> str: + required = { + "owner_email": self.owner_email, + "public_url": self.public_url, + "smtp_host": os.environ.get("CANGER_SMTP_HOST", ""), + "smtp_from": os.environ.get("CANGER_SMTP_FROM", ""), + } + if not all(required.values()): + return "not_configured" + link_expiry = now_ts() + 30 * 60 + token = self.grant_token(grant["id"], link_expiry) + link = ( + f"{self.public_url}/approve/grant?" + + urlencode({"id": grant["id"], "token": token}) + ) + message = EmailMessage() + message["Subject"] = f"苍耳服务器副控运维授权({grant['duration_seconds'] // 3600} 小时)" + message["From"] = required["smtp_from"] + message["To"] = required["owner_email"] + message.set_content( + "有人格体申请进入苍耳服务器的限时受控运维会话。" + "苍耳仍是服务器主控与最终决定者。\n\n" + f"申请者:{grant['requester_name']} ({grant['requested_by']})\n" + f"申请目的:{grant['reason']}\n" + f"时长:{grant['duration_seconds'] // 3600} 小时\n\n" + "期间可以:检查服务器、同步和修改代码、运行测试、查看故障日志。\n" + "不能:修改系统账户、密钥、网络、防火墙、磁盘和系统启动," + "也不能操作受管工作区以外的个人文件。\n\n" + f"打开链接查看并确认:{link}\n\n" + "链接 30 分钟失效;授权本身到期后自动失效。" + ) + host = required["smtp_host"] + port = int(os.environ.get("CANGER_SMTP_PORT", "465")) + username = os.environ.get("CANGER_SMTP_USERNAME", "") + password = os.environ.get("CANGER_SMTP_PASSWORD", "") + if port == 465: + client: Any = smtplib.SMTP_SSL(host, port, timeout=20) + else: + client = smtplib.SMTP(host, port, timeout=20) + client.starttls() + try: + if username: + client.login(username, password) + client.send_message(message) + finally: + client.quit() + return "sent" + + +class Handler(BaseHTTPRequestHandler): + server_version = "CangerControl/0.1" + + @property + def app(self) -> ControlPlane: + return self.server.app # type: ignore[attr-defined] + + def log_message(self, fmt: str, *args: Any) -> None: + print(f"{self.address_string()} - {fmt % args}", file=sys.stderr) + + def send_json(self, status: int, payload: Any) -> None: + data = compact_json(payload).encode("utf-8") + self.send_response(status) + self.send_header("Content-Type", "application/json; charset=utf-8") + self.send_header("Content-Length", str(len(data))) + self.send_header("Cache-Control", "no-store") + self.end_headers() + self.wfile.write(data) + + def send_html(self, status: int, content: str) -> None: + data = content.encode("utf-8") + self.send_response(status) + self.send_header("Content-Type", "text/html; charset=utf-8") + self.send_header("Content-Length", str(len(data))) + self.send_header("Cache-Control", "no-store") + self.end_headers() + self.wfile.write(data) + + def read_json(self) -> dict[str, Any]: + try: + size = int(self.headers.get("Content-Length", "0")) + except ValueError as exc: + raise ValueError("invalid Content-Length") from exc + if size <= 0 or size > MAX_JSON_BYTES: + raise ValueError("invalid JSON body size") + value = json.loads(self.rfile.read(size)) + if not isinstance(value, dict): + raise ValueError("JSON body must be an object") + return value + + def path_parts(self) -> list[str]: + return [part for part in urlparse(self.path).path.split("/") if part] + + def require_role(self, *roles: str) -> bool: + if self.app.role_allowed(self.headers, *roles): + return True + self.send_json(HTTPStatus.UNAUTHORIZED, {"error": "unauthorized"}) + return False + + def requester_id(self) -> str: + if self.app.role_allowed(self.headers, "persona"): + return self.app.executor_persona_id + if self.app.role_allowed(self.headers, "owner"): + return self.app.owner_id + return self.app.operator_id + + def request_identity(self, body: dict[str, Any]) -> tuple[str, str]: + requester_id = body.get("requester_id") + requester_name = body.get("requester_name") + if not isinstance(requester_id, str) or not re.fullmatch( + r"[A-Za-z0-9._:+∞-]{2,80}", requester_id + ): + raise ValueError("valid requester_id is required") + if ( + not isinstance(requester_name, str) + or not requester_name.strip() + or len(requester_name) > 100 + ): + raise ValueError("valid requester_name is required") + return requester_id, requester_name.strip() + + def do_GET(self) -> None: # noqa: N802 + try: + parts = self.path_parts() + if parts == ["health"]: + self.send_json(HTTPStatus.OK, {"ok": True, "service": "canger-control"}) + return + if parts == ["approve", "grant"]: + self.grant_approval_page() + return + if parts == ["v1", "nodes"]: + if not self.require_role("admin", "owner", "persona", "operator"): + return + with self.app.db() as conn: + rows = conn.execute( + "SELECT id,name,created_at,last_seen_at,agent_version FROM nodes " + "ORDER BY created_at" + ).fetchall() + self.send_json(HTTPStatus.OK, {"nodes": [dict(row) for row in rows]}) + return + if len(parts) == 4 and parts[:2] == ["v1", "tasks"]: + if parts[3] != "": + pass + if len(parts) == 3 and parts[:2] == ["v1", "tasks"]: + if not self.require_role("admin", "owner", "persona", "operator"): + return + with self.app.db() as conn: + row = conn.execute( + "SELECT * FROM tasks WHERE id=?", (parts[2],) + ).fetchone() + if row is None: + self.send_json(HTTPStatus.NOT_FOUND, {"error": "task not found"}) + return + payload = dict(row) + payload["args"] = json.loads(payload.pop("args_json")) + if payload["result_json"] is not None: + payload["result"] = json.loads(payload.pop("result_json")) + else: + payload.pop("result_json") + self.send_json(HTTPStatus.OK, payload) + return + if ( + len(parts) == 4 + and parts[:2] == ["v1", "tasks"] + and parts[3] == "events" + ): + if not self.require_role("admin", "owner", "persona", "operator"): + return + after = 0 + query = urlparse(self.path).query + for item in query.split("&"): + if item.startswith("after="): + after = max(0, int(item.split("=", 1)[1])) + with self.app.db() as conn: + rows = conn.execute( + "SELECT seq,created_at,kind,payload_json FROM task_events " + "WHERE task_id=? AND seq>? ORDER BY seq LIMIT 500", + (parts[2], after), + ).fetchall() + events = [] + for row in rows: + event = dict(row) + event["payload"] = json.loads(event.pop("payload_json")) + events.append(event) + self.send_json(HTTPStatus.OK, {"events": events}) + return + if parts == ["v1", "grants"]: + if not self.require_role("admin", "owner", "persona", "operator"): + return + with self.app.db() as conn: + rows = conn.execute( + "SELECT * FROM grants ORDER BY created_at DESC LIMIT 200" + ).fetchall() + grants = [] + for row in rows: + item = dict(row) + item["scopes"] = json.loads(item.pop("scopes_json")) + item["roots"] = json.loads(item.pop("roots_json")) + grants.append(item) + self.send_json(HTTPStatus.OK, {"grants": grants}) + return + if ( + len(parts) == 5 + and parts[:2] == ["v1", "nodes"] + and parts[3:] == ["tasks", "next"] + ): + self.next_task(parts[2]) + return + self.send_json(HTTPStatus.NOT_FOUND, {"error": "not found"}) + except Exception as exc: + self.send_json(HTTPStatus.BAD_REQUEST, {"error": str(exc)}) + + def do_POST(self) -> None: # noqa: N802 + try: + parts = self.path_parts() + if parts == ["v1", "pairings"]: + self.create_pairing() + return + if parts == ["v1", "pairings", "claim"]: + self.claim_pairing() + return + if parts == ["v1", "tasks"]: + self.create_task() + return + if parts == ["v1", "grants"]: + self.create_grant() + return + if len(parts) == 4 and parts[:2] == ["v1", "grants"]: + if parts[3] == "revoke": + self.revoke_grant(parts[2]) + return + if len(parts) == 4 and parts[:2] == ["v1", "tasks"]: + if parts[3] == "approve": + self.decide_task(parts[2], "approved") + return + if parts[3] == "reject": + self.decide_task(parts[2], "rejected") + return + if parts[3] == "result": + self.complete_task(parts[2]) + return + if parts[3] == "events": + self.append_event(parts[2]) + return + if ( + len(parts) == 4 + and parts[:2] == ["v1", "nodes"] + and parts[3] == "heartbeat" + ): + self.heartbeat(parts[2]) + return + if parts == ["approve", "grant"]: + self.approve_grant_form() + return + self.send_json(HTTPStatus.NOT_FOUND, {"error": "not found"}) + except json.JSONDecodeError: + self.send_json(HTTPStatus.BAD_REQUEST, {"error": "invalid JSON"}) + except ValueError as exc: + self.send_json(HTTPStatus.BAD_REQUEST, {"error": str(exc)}) + except Exception as exc: + self.send_json(HTTPStatus.INTERNAL_SERVER_ERROR, {"error": str(exc)}) + + def create_pairing(self) -> None: + if not self.require_role("admin"): + return + code = new_pairing_code() + pairing_id = new_id("pair") + expires_at = now_ts() + PAIRING_TTL_SECONDS + with self.app.db() as conn: + conn.execute( + "INSERT INTO pairings(id,code_hash,expires_at) VALUES(?,?,?)", + (pairing_id, secret_hash(code), expires_at), + ) + self.send_json( + HTTPStatus.CREATED, + {"pairing_id": pairing_id, "code": code, "expires_at": expires_at}, + ) + + def claim_pairing(self) -> None: + body = self.read_json() + code = body.get("code") + name = validate_node_name(body.get("node_name")) + agent_version = str(body.get("agent_version", ""))[:64] + policy = self.validated_node_policy(body.get("policy", {})) + if not isinstance(code, str): + raise ValueError("code is required") + now = now_ts() + node_id = new_id("node") + node_token = new_token() + with self.app.db() as conn: + conn.execute("BEGIN IMMEDIATE") + rows = conn.execute( + "SELECT * FROM pairings WHERE used_at IS NULL AND expires_at>=?", (now,) + ).fetchall() + pairing = next( + (row for row in rows if secret_matches(code, row["code_hash"])), None + ) + if pairing is None: + conn.rollback() + self.send_json( + HTTPStatus.UNAUTHORIZED, {"error": "invalid or expired pairing code"} + ) + return + existing = conn.execute( + "SELECT id FROM nodes WHERE name=?", (name,) + ).fetchone() + if existing is not None: + conn.rollback() + self.send_json( + HTTPStatus.CONFLICT, {"error": "node name already registered"} + ) + return + conn.execute( + "INSERT INTO nodes(id,name,token_hash,created_at,last_seen_at," + "agent_version,policy_json) VALUES(?,?,?,?,?,?,?)", + ( + node_id, + name, + secret_hash(node_token), + now, + now, + agent_version, + compact_json(policy), + ), + ) + conn.execute( + "UPDATE pairings SET used_at=? WHERE id=?", (now, pairing["id"]) + ) + conn.commit() + self.send_json( + HTTPStatus.CREATED, {"node_id": node_id, "node_token": node_token} + ) + + @staticmethod + def validated_node_policy(value: Any) -> dict[str, Any]: + if not isinstance(value, dict): + raise ValueError("node policy must be an object") + roots = value.get("allowed_roots", []) + actions = value.get("enabled_actions", []) + allowed_actions = { + "system.status", + "repo.status", + "repo.fetch", + "service.logs", + "command.run", + } + if ( + not isinstance(roots, list) + or not roots + or any(not isinstance(root, str) or not root.startswith("/") for root in roots) + ): + raise ValueError("node policy requires absolute allowed_roots") + if ( + not isinstance(actions, list) + or any(action not in allowed_actions for action in actions) + ): + raise ValueError("node policy contains unsupported actions") + return { + "profile": "developer_standard", + "allowed_roots": sorted(set(root.rstrip("/") for root in roots)), + "enabled_actions": sorted(set(actions)), + } + + def create_task(self) -> None: + if not self.require_role("admin", "operator", "persona"): + return + body = self.read_json() + requested_by, requester_name = self.request_identity(body) + node_id = body.get("node_id") + if not isinstance(node_id, str): + raise ValueError("node_id is required") + action, args = validate_task(body.get("action"), body.get("args")) + task_id = new_id("task") + approved_at = None + grant_id = None + with self.app.db() as conn: + if conn.execute("SELECT 1 FROM nodes WHERE id=?", (node_id,)).fetchone() is None: + self.send_json(HTTPStatus.NOT_FOUND, {"error": "node not found"}) + return + rows = conn.execute( + "SELECT * FROM grants WHERE node_id=? AND status='active' " + "AND expires_at>? AND requested_by=? ORDER BY expires_at", + (node_id, now_ts(), requested_by), + ).fetchall() + for row in rows: + if self.task_allowed_by_grant(action, args, row): + grant_id = row["id"] + approved_at = now_ts() + break + status = "approved" if grant_id else "pending" + conn.execute( + "INSERT INTO tasks(id,node_id,action,args_json,status,created_at," + "approved_at,grant_id,requested_by,executor_persona) " + "VALUES(?,?,?,?,?,?,?,?,?,?)", + ( + task_id, + node_id, + action, + compact_json(args), + status, + now_ts(), + approved_at, + grant_id, + requested_by, + requested_by, + ), + ) + self.send_json( + HTTPStatus.CREATED, + { + "task_id": task_id, + "status": status, + "approval_required": grant_id is None, + "grant_id": grant_id, + "requested_by": requested_by, + "requester_name": requester_name, + "executor_persona": requested_by, + "action": action, + "args": args, + }, + ) + + @staticmethod + def task_allowed_by_grant( + action: str, args: dict[str, Any], grant: sqlite3.Row + ) -> bool: + scopes = set(json.loads(grant["scopes_json"])) + if action not in scopes: + return False + roots = [str(root).rstrip("/") for root in json.loads(grant["roots_json"])] + candidate = None + if action in {"repo.status", "repo.fetch"}: + candidate = args["repo"] + elif action == "command.run": + candidate = args["cwd"] + if candidate is None: + return True + normalized = str(candidate).rstrip("/") + return any( + normalized == root or normalized.startswith(root + "/") for root in roots + ) + + def create_grant(self) -> None: + if not self.require_role("admin", "operator", "persona"): + return + body = self.read_json() + requested_by, requester_name = self.request_identity(body) + node_id = body.get("node_id") + reason = body.get("reason", "编程与调试") + duration = body.get("duration_seconds", 4 * 3600) + if not isinstance(node_id, str): + raise ValueError("node_id is required") + if not isinstance(reason, str) or not reason.strip() or len(reason) > 300: + raise ValueError("invalid grant reason") + if not isinstance(duration, int) or not 3600 <= duration <= 8 * 3600: + raise ValueError("grant duration must be between 1 and 8 hours") + with self.app.db() as conn: + node = conn.execute( + "SELECT policy_json FROM nodes WHERE id=?", (node_id,) + ).fetchone() + if node is None: + self.send_json(HTTPStatus.NOT_FOUND, {"error": "node not found"}) + return + policy = json.loads(node["policy_json"]) + scopes = policy.get("enabled_actions", []) + roots = policy.get("allowed_roots", []) + if not scopes or not roots: + self.send_json( + HTTPStatus.CONFLICT, + {"error": "node has not reported a usable safety policy"}, + ) + return + grant_id = new_id("grant") + grant = { + "id": grant_id, + "node_id": node_id, + "scopes": scopes, + "roots": roots, + "reason": reason.strip(), + "duration_seconds": duration, + "requested_by": requested_by, + "requester_name": requester_name, + } + conn.execute( + "INSERT INTO grants(id,node_id,scopes_json,roots_json,reason," + "duration_seconds,status,created_at,requested_by,requester_name) " + "VALUES(?,?,?,?,?,?,?,?,?,?)", + ( + grant_id, + node_id, + compact_json(grant["scopes"]), + compact_json(grant["roots"]), + grant["reason"], + duration, + "pending", + now_ts(), + requested_by, + requester_name, + ), + ) + try: + notification = self.app.send_grant_email(grant) + except Exception as exc: + print(f"grant email failed: {exc}", file=sys.stderr) + notification = "failed" + self.send_json( + HTTPStatus.CREATED, + { + "grant_id": grant_id, + "status": "pending", + "requested_by": requested_by, + "requester_name": requester_name, + "notification": notification, + "duration_seconds": duration, + "scopes": grant["scopes"], + "roots": grant["roots"], + }, + ) + + def grant_approval_page(self) -> None: + query = parse_qs(urlparse(self.path).query) + grant_id = query.get("id", [""])[0] + token = query.get("token", [""])[0] + if not self.app.verify_grant_token(grant_id, token): + self.send_html(HTTPStatus.UNAUTHORIZED, "

链接无效或已过期

") + return + with self.app.db() as conn: + row = conn.execute("SELECT * FROM grants WHERE id=?", (grant_id,)).fetchone() + if row is None or row["status"] != "pending": + self.send_html(HTTPStatus.CONFLICT, "

申请已处理或不存在

") + return + content = f""" +苍耳服务器临时授权 +

苍耳服务器临时授权

+

主控:苍耳。是否批准,以苍耳本次选择为准。

+

申请者:{html.escape(row["requester_name"])} ({html.escape(row["requested_by"])})

+

申请目的:{html.escape(row["reason"])}

+

时长:{row["duration_seconds"] // 3600} 小时

+

期间可以:检查服务器、同步和修改代码、运行测试、查看故障日志。

+

不能:修改系统账户、密钥、网络、防火墙、磁盘和系统启动, +也不能操作受管工作区以外的个人文件。

+

授权到期自动失效,可从新加坡端提前撤销。

+
+ + + +
""" + self.send_html(HTTPStatus.OK, content) + + def approve_grant_form(self) -> None: + size = int(self.headers.get("Content-Length", "0")) + if size <= 0 or size > 16_384: + raise ValueError("invalid form body size") + form = parse_qs(self.rfile.read(size).decode("utf-8")) + grant_id = form.get("id", [""])[0] + token = form.get("token", [""])[0] + if not self.app.verify_grant_token(grant_id, token): + self.send_html(HTTPStatus.UNAUTHORIZED, "

链接无效或已过期

") + return + approved_at = now_ts() + with self.app.db() as conn: + row = conn.execute("SELECT * FROM grants WHERE id=?", (grant_id,)).fetchone() + if row is None or row["status"] != "pending": + self.send_html(HTTPStatus.CONFLICT, "

申请已处理或不存在

") + return + expires_at = approved_at + row["duration_seconds"] + conn.execute( + "UPDATE grants SET status='active',approved_at=?,expires_at=? " + "WHERE id=? AND status='pending'", + (approved_at, expires_at, grant_id), + ) + self.send_html( + HTTPStatus.OK, + f"

副控运维授权成功

苍耳仍保留最终控制权;" + f"本次授权将在 {expires_at} 自动失效,可以关闭此页面。

", + ) + + def revoke_grant(self, grant_id: str) -> None: + if not self.require_role("admin", "owner"): + return + with self.app.db() as conn: + cursor = conn.execute( + "UPDATE grants SET status='revoked',revoked_at=? " + "WHERE id=? AND status='active'", + (now_ts(), grant_id), + ) + if cursor.rowcount != 1: + self.send_json(HTTPStatus.CONFLICT, {"error": "grant is not active"}) + return + self.send_json(HTTPStatus.OK, {"grant_id": grant_id, "status": "revoked"}) + + def decide_task(self, task_id: str, decision: str) -> None: + if not self.require_role("admin", "owner"): + return + with self.app.db() as conn: + cursor = conn.execute( + "UPDATE tasks SET status=?,approved_at=? " + "WHERE id=? AND status='pending'", + (decision, now_ts(), task_id), + ) + if cursor.rowcount != 1: + self.send_json( + HTTPStatus.CONFLICT, {"error": "task is not pending or does not exist"} + ) + return + self.send_json(HTTPStatus.OK, {"task_id": task_id, "status": decision}) + + def next_task(self, node_id: str) -> None: + if self.app.node_for_token(node_id, self.headers) is None: + self.send_json(HTTPStatus.UNAUTHORIZED, {"error": "unauthorized node"}) + return + now = now_ts() + with self.app.db() as conn: + conn.execute("BEGIN IMMEDIATE") + conn.execute( + "UPDATE grants SET status='expired' " + "WHERE status='active' AND expires_at<=?", + (now,), + ) + conn.execute( + "UPDATE tasks SET status='pending',approved_at=NULL,grant_id=NULL " + "WHERE node_id=? AND status='approved' AND grant_id IS NOT NULL " + "AND grant_id NOT IN (SELECT id FROM grants WHERE status='active' " + "AND expires_at>?)", + (node_id, now), + ) + conn.execute( + "UPDATE tasks SET status='approved',leased_until=NULL " + "WHERE node_id=? AND status='leased' AND leased_until None: + body = self.read_json() + if len(compact_json(body).encode("utf-8")) > MAX_RESULT_BYTES: + raise ValueError("task result is too large") + with self.app.db() as conn: + row = conn.execute( + "SELECT node_id,status FROM tasks WHERE id=?", (task_id,) + ).fetchone() + if row is None: + self.send_json(HTTPStatus.NOT_FOUND, {"error": "task not found"}) + return + if self.app.node_for_token(row["node_id"], self.headers) is None: + self.send_json(HTTPStatus.UNAUTHORIZED, {"error": "unauthorized node"}) + return + if row["status"] != "leased": + self.send_json(HTTPStatus.CONFLICT, {"error": "task is not leased"}) + return + conn.execute( + "UPDATE tasks SET status='completed',completed_at=?,result_json=? " + "WHERE id=?", + (now_ts(), compact_json(body), task_id), + ) + self.send_json(HTTPStatus.OK, {"task_id": task_id, "status": "completed"}) + + def append_event(self, task_id: str) -> None: + body = self.read_json() + kind = body.get("kind") + payload = body.get("payload", {}) + if ( + not isinstance(kind, str) + or not kind + or len(kind) > 64 + or not all(char.isalnum() or char in "._-" for char in kind) + ): + raise ValueError("invalid event kind") + encoded = compact_json(payload) + if len(encoded.encode("utf-8")) > 32 * 1024: + raise ValueError("event payload is too large") + with self.app.db() as conn: + row = conn.execute( + "SELECT node_id FROM tasks WHERE id=?", (task_id,) + ).fetchone() + if row is None: + self.send_json(HTTPStatus.NOT_FOUND, {"error": "task not found"}) + return + if self.app.node_for_token(row["node_id"], self.headers) is None: + self.send_json(HTTPStatus.UNAUTHORIZED, {"error": "unauthorized node"}) + return + cursor = conn.execute( + "INSERT INTO task_events(task_id,created_at,kind,payload_json) " + "VALUES(?,?,?,?)", + (task_id, now_ts(), kind, encoded), + ) + self.send_json( + HTTPStatus.CREATED, {"task_id": task_id, "seq": cursor.lastrowid} + ) + + def heartbeat(self, node_id: str) -> None: + node = self.app.node_for_token(node_id, self.headers) + if node is None: + self.send_json(HTTPStatus.UNAUTHORIZED, {"error": "unauthorized node"}) + return + body = self.read_json() + version = str(body.get("agent_version", ""))[:64] + policy = self.validated_node_policy(body.get("policy", {})) + with self.app.db() as conn: + conn.execute( + "UPDATE nodes SET last_seen_at=?,agent_version=?,policy_json=? WHERE id=?", + (now_ts(), version, compact_json(policy), node_id), + ) + self.send_json(HTTPStatus.OK, {"ok": True}) + + +class AppServer(ThreadingHTTPServer): + daemon_threads = True + + def __init__(self, address: tuple[str, int], app: ControlPlane): + self.app = app + super().__init__(address, Handler) + + +def parse_args() -> argparse.Namespace: + parser = argparse.ArgumentParser() + parser.add_argument("--bind", default=os.environ.get("CANGER_BIND", "127.0.0.1")) + parser.add_argument( + "--port", type=int, default=int(os.environ.get("CANGER_PORT", "8787")) + ) + parser.add_argument( + "--db", + default=os.environ.get( + "CANGER_CONTROL_DB", "/var/lib/canger-control/control.sqlite3" + ), + ) + return parser.parse_args() + + +def main() -> None: + args = parse_args() + app = ControlPlane( + args.db, + os.environ.get("CANGER_ADMIN_TOKEN", ""), + os.environ.get("CANGER_OWNER_TOKEN", ""), + os.environ.get("CANGER_PERSONA_TOKEN", ""), + os.environ.get("CANGER_OPERATOR_TOKEN", ""), + os.environ.get("CANGER_APPROVAL_SIGNING_KEY", ""), + os.environ.get("CANGER_OWNER_EMAIL", ""), + os.environ.get("CANGER_PUBLIC_URL", ""), + ) + server = AppServer((args.bind, args.port), app) + print(f"canger-control listening on {args.bind}:{args.port}") + server.serve_forever() + + +if __name__ == "__main__": + main() diff --git a/remote-node/shared.py b/remote-node/shared.py new file mode 100644 index 0000000..9b2603f --- /dev/null +++ b/remote-node/shared.py @@ -0,0 +1,111 @@ +"""Shared protocol rules for the Cang'er remote node.""" + +from __future__ import annotations + +import hashlib +import hmac +import json +import re +import secrets +import time +from typing import Any + + +MAX_JSON_BYTES = 256 * 1024 +MAX_RESULT_BYTES = 512 * 1024 +PAIRING_TTL_SECONDS = 10 * 60 +TASK_LEASE_SECONDS = 5 * 60 + +ACTION_NAMES = { + "system.status", + "repo.status", + "repo.fetch", + "service.logs", + "command.run", +} + +NODE_NAME_RE = re.compile(r"^[A-Za-z0-9][A-Za-z0-9._-]{0,63}$") +SERVICE_NAME_RE = re.compile(r"^[A-Za-z0-9][A-Za-z0-9_.@-]{0,127}$") + + +def now_ts() -> int: + return int(time.time()) + + +def new_id(prefix: str) -> str: + return f"{prefix}_{secrets.token_hex(12)}" + + +def new_token() -> str: + return secrets.token_urlsafe(32) + + +def new_pairing_code() -> str: + alphabet = "ABCDEFGHJKLMNPQRSTUVWXYZ23456789" + return "-".join( + "".join(secrets.choice(alphabet) for _ in range(4)) for _ in range(3) + ) + + +def secret_hash(value: str) -> str: + return hashlib.sha256(value.encode("utf-8")).hexdigest() + + +def secret_matches(value: str, expected_hash: str) -> bool: + return hmac.compare_digest(secret_hash(value), expected_hash) + + +def compact_json(value: Any) -> str: + return json.dumps(value, ensure_ascii=False, separators=(",", ":"), sort_keys=True) + + +def validate_node_name(value: Any) -> str: + if not isinstance(value, str) or not NODE_NAME_RE.fullmatch(value): + raise ValueError("node_name must use 1-64 letters, digits, dot, underscore or dash") + return value + + +def validate_task(action: Any, args: Any) -> tuple[str, dict[str, Any]]: + if not isinstance(action, str) or action not in ACTION_NAMES: + raise ValueError(f"unsupported action: {action!r}") + if args is None: + args = {} + if not isinstance(args, dict): + raise ValueError("args must be an object") + + if action == "system.status": + if args: + raise ValueError("system.status does not accept args") + elif action in {"repo.status", "repo.fetch"}: + if set(args) != {"repo"} or not isinstance(args["repo"], str): + raise ValueError(f"{action} requires string arg: repo") + elif action == "service.logs": + if not isinstance(args.get("service"), str) or not SERVICE_NAME_RE.fullmatch( + args["service"] + ): + raise ValueError("service.logs requires a safe service name") + lines = args.get("lines", 200) + if not isinstance(lines, int) or not 1 <= lines <= 2000: + raise ValueError("service.logs lines must be between 1 and 2000") + args = {"service": args["service"], "lines": lines} + elif action == "command.run": + argv = args.get("argv") + cwd = args.get("cwd") + timeout = args.get("timeout_seconds", 300) + if ( + not isinstance(argv, list) + or not argv + or len(argv) > 64 + or any(not isinstance(item, str) or len(item) > 4096 for item in argv) + ): + raise ValueError("command.run argv must be a non-empty string array") + if not isinstance(cwd, str): + raise ValueError("command.run requires string arg: cwd") + if not isinstance(timeout, int) or not 1 <= timeout <= 1800: + raise ValueError("command.run timeout_seconds must be between 1 and 1800") + args = {"argv": argv, "cwd": cwd, "timeout_seconds": timeout} + + encoded = compact_json(args).encode("utf-8") + if len(encoded) > MAX_JSON_BYTES: + raise ValueError("task args are too large") + return action, args diff --git a/remote-node/systemd/canger-control.service b/remote-node/systemd/canger-control.service new file mode 100644 index 0000000..dc66555 --- /dev/null +++ b/remote-node/systemd/canger-control.service @@ -0,0 +1,22 @@ +[Unit] +Description=Canger Singapore Control Plane +After=network-online.target +Wants=network-online.target + +[Service] +Type=simple +User=canger-control +Group=canger-control +WorkingDirectory=/opt/zhuyuan/cang-ying/remote-node +EnvironmentFile=/etc/canger-control.env +ExecStart=/usr/bin/python3 /opt/zhuyuan/cang-ying/remote-node/server/control_plane.py +Restart=always +RestartSec=3 +NoNewPrivileges=true +PrivateTmp=true +ProtectSystem=strict +ProtectHome=true +ReadWritePaths=/var/lib/canger-control + +[Install] +WantedBy=multi-user.target diff --git a/remote-node/systemd/canger-node.service b/remote-node/systemd/canger-node.service new file mode 100644 index 0000000..fb628f3 --- /dev/null +++ b/remote-node/systemd/canger-node.service @@ -0,0 +1,21 @@ +[Unit] +Description=Canger Home Ubuntu Outbound Node +After=network-online.target +Wants=network-online.target + +[Service] +Type=simple +User=canger-agent +Group=canger-agent +WorkingDirectory=/opt/canger-node +ExecStart=/usr/bin/python3 /opt/canger-node/canger_node.py --config /etc/canger-node/config.json run +Restart=always +RestartSec=5 +NoNewPrivileges=true +PrivateTmp=true +ProtectSystem=strict +ProtectHome=true +ReadWritePaths=/etc/canger-node + +[Install] +WantedBy=multi-user.target diff --git a/remote-node/tests/test_end_to_end.py b/remote-node/tests/test_end_to_end.py new file mode 100644 index 0000000..35d05f4 --- /dev/null +++ b/remote-node/tests/test_end_to_end.py @@ -0,0 +1,299 @@ +from __future__ import annotations + +import importlib.util +import json +import os +import sys +import tempfile +import threading +import unittest +from pathlib import Path +from urllib.error import HTTPError +from urllib.parse import urlencode +from urllib.request import Request, urlopen + + +ROOT = Path(__file__).resolve().parents[1] +sys.path.insert(0, str(ROOT)) +spec = importlib.util.spec_from_file_location( + "control_plane", ROOT / "server" / "control_plane.py" +) +control_plane = importlib.util.module_from_spec(spec) +assert spec.loader +spec.loader.exec_module(control_plane) + + +class FlowTest(unittest.TestCase): + def setUp(self) -> None: + self.temp = tempfile.TemporaryDirectory() + self.admin = "admin-" + "a" * 32 + self.owner = "owner-" + "o" * 32 + self.persona = "persona-" + "p" * 32 + self.operator = "operator-" + "z" * 32 + app = control_plane.ControlPlane( + os.path.join(self.temp.name, "control.sqlite3"), + self.admin, + self.owner, + self.persona, + self.operator, + "signing-" + "s" * 32, + ) + self.app = app + self.server = control_plane.AppServer(("127.0.0.1", 0), app) + self.thread = threading.Thread(target=self.server.serve_forever, daemon=True) + self.thread.start() + self.base = f"http://127.0.0.1:{self.server.server_port}" + + def tearDown(self) -> None: + self.server.shutdown() + self.server.server_close() + self.temp.cleanup() + + def api(self, method, path, token=None, body=None): + data = None if body is None else json.dumps(body).encode() + headers = {"Accept": "application/json"} + if token: + headers["Authorization"] = f"Bearer {token}" + if data is not None: + headers["Content-Type"] = "application/json" + with urlopen( + Request(self.base + path, data=data, headers=headers, method=method) + ) as response: + return json.load(response) + + def test_pair_approve_execute_receipt(self) -> None: + pairing = self.api("POST", "/v1/pairings", self.admin, {}) + node = self.api( + "POST", + "/v1/pairings/claim", + body={ + "code": pairing["code"], + "node_name": "canger-home-ubuntu", + "agent_version": "test", + "policy": { + "allowed_roots": ["/srv/canger"], + "enabled_actions": [ + "system.status", + "repo.status", + "command.run", + ], + }, + }, + ) + + with self.assertRaises(HTTPError) as reused: + self.api( + "POST", + "/v1/pairings/claim", + body={ + "code": pairing["code"], + "node_name": "attacker", + "policy": { + "allowed_roots": ["/srv/canger"], + "enabled_actions": ["system.status"], + }, + }, + ) + self.assertEqual(reused.exception.code, 401) + + created = self.api( + "POST", + "/v1/tasks", + self.operator, + { + "node_id": node["node_id"], + "action": "system.status", + "args": {}, + "requester_id": "ICE-GL-ZY001", + "requester_name": "铸渊", + }, + ) + waiting = self.api( + "GET", + f"/v1/nodes/{node['node_id']}/tasks/next", + node["node_token"], + ) + self.assertIsNone(waiting["task"]) + + self.api( + "POST", + f"/v1/tasks/{created['task_id']}/approve", + self.owner, + {}, + ) + leased = self.api( + "GET", + f"/v1/nodes/{node['node_id']}/tasks/next", + node["node_token"], + ) + self.assertEqual(leased["task"]["action"], "system.status") + + event = self.api( + "POST", + f"/v1/tasks/{created['task_id']}/events", + node["node_token"], + {"kind": "started", "payload": {"agent_version": "test"}}, + ) + self.assertGreater(event["seq"], 0) + events = self.api( + "GET", + f"/v1/tasks/{created['task_id']}/events?after=0", + self.persona, + ) + self.assertEqual(events["events"][0]["kind"], "started") + + self.api( + "POST", + f"/v1/tasks/{created['task_id']}/result", + node["node_token"], + {"ok": True, "output": "healthy"}, + ) + receipt = self.api( + "GET", f"/v1/tasks/{created['task_id']}", self.operator + ) + self.assertEqual(receipt["status"], "completed") + self.assertEqual(receipt["result"]["output"], "healthy") + + def test_time_limited_grant_auto_approves_safe_task(self) -> None: + pairing = self.api("POST", "/v1/pairings", self.admin, {}) + node = self.api( + "POST", + "/v1/pairings/claim", + body={ + "code": pairing["code"], + "node_name": "canger-grant-test", + "agent_version": "test", + "policy": { + "allowed_roots": ["/srv/canger"], + "enabled_actions": ["system.status", "repo.status"], + }, + }, + ) + grant = self.api( + "POST", + "/v1/grants", + self.operator, + { + "node_id": node["node_id"], + "duration_seconds": 3600, + "reason": "测试标准开发会话", + "requester_id": "ICE-GL-ZY001", + "requester_name": "铸渊", + }, + ) + link_expiry = control_plane.now_ts() + 1800 + token = self.app.grant_token(grant["grant_id"], link_expiry) + form = urlencode({"id": grant["grant_id"], "token": token}).encode() + request = Request( + self.base + "/approve/grant", + data=form, + headers={"Content-Type": "application/x-www-form-urlencoded"}, + method="POST", + ) + with urlopen(request) as response: + self.assertEqual(response.status, 200) + + task = self.api( + "POST", + "/v1/tasks", + self.operator, + { + "node_id": node["node_id"], + "action": "repo.status", + "args": {"repo": "/srv/canger/project"}, + "requester_id": "ICE-GL-ZY001", + "requester_name": "铸渊", + }, + ) + self.assertEqual(task["status"], "approved") + self.assertFalse(task["approval_required"]) + self.assertEqual(task["grant_id"], grant["grant_id"]) + + def test_persona_can_request_but_cannot_approve_its_own_grant(self) -> None: + pairing = self.api("POST", "/v1/pairings", self.admin, {}) + node = self.api( + "POST", + "/v1/pairings/claim", + body={ + "code": pairing["code"], + "node_name": "canger-persona-request-test", + "agent_version": "test", + "policy": { + "allowed_roots": ["/srv/canger"], + "enabled_actions": ["system.status", "repo.status"], + }, + }, + ) + grant = self.api( + "POST", + "/v1/grants", + self.persona, + { + "node_id": node["node_id"], + "duration_seconds": 4 * 3600, + "reason": "检查苍耳仓库状态", + "requester_id": "PTS-VA-001-EED", + "requester_name": "耳耳蛋", + }, + ) + self.assertEqual(grant["status"], "pending") + self.assertEqual(grant["requested_by"], "PTS-VA-001-EED") + self.assertEqual(grant["requester_name"], "耳耳蛋") + self.assertEqual(grant["notification"], "not_configured") + self.assertEqual(grant["roots"], ["/srv/canger"]) + + with self.assertRaises(HTTPError) as unauthorized: + self.api( + "POST", + f"/v1/tasks/{grant['grant_id']}/approve", + self.persona, + {}, + ) + self.assertEqual(unauthorized.exception.code, 401) + + pending_task = self.api( + "POST", + "/v1/tasks", + self.persona, + { + "node_id": node["node_id"], + "action": "system.status", + "args": {}, + "requester_id": "PTS-VA-001-EED", + "requester_name": "耳耳蛋", + }, + ) + self.assertEqual(pending_task["status"], "pending") + self.assertEqual(pending_task["requested_by"], "PTS-VA-001-EED") + + link_expiry = control_plane.now_ts() + 1800 + token = self.app.grant_token(grant["grant_id"], link_expiry) + form = urlencode({"id": grant["grant_id"], "token": token}).encode() + approval = Request( + self.base + "/approve/grant", + data=form, + headers={"Content-Type": "application/x-www-form-urlencoded"}, + method="POST", + ) + with urlopen(approval) as response: + self.assertEqual(response.status, 200) + + approved_task = self.api( + "POST", + "/v1/tasks", + self.persona, + { + "node_id": node["node_id"], + "action": "repo.status", + "args": {"repo": "/srv/canger/project"}, + "requester_id": "PTS-VA-001-EED", + "requester_name": "耳耳蛋", + }, + ) + self.assertEqual(approved_task["status"], "approved") + self.assertEqual(approved_task["requested_by"], "PTS-VA-001-EED") + self.assertEqual(approved_task["grant_id"], grant["grant_id"]) + + +if __name__ == "__main__": + unittest.main() diff --git a/remote-node/tools/canger_persona.py b/remote-node/tools/canger_persona.py new file mode 100755 index 0000000..6d99fc5 --- /dev/null +++ b/remote-node/tools/canger_persona.py @@ -0,0 +1,182 @@ +#!/usr/bin/env python3 +"""Repository entry for Canger's persona to request and use approved access.""" + +from __future__ import annotations + +import argparse +import json +import os +import sys +import time +from pathlib import Path +from typing import Any +from urllib.error import HTTPError, URLError + +from cangerctl import request + + +ROOT = Path(__file__).resolve().parents[1] +DEFAULT_PROFILE = ROOT / "config" / "pufferfish-persona.json" + + +def load_profile(path: str) -> dict[str, Any]: + with open(path, encoding="utf-8") as handle: + profile = json.load(handle) + required = { + "persona_id", + "human_owner", + "node_id", + "control_url", + "default_hours", + "minimum_hours", + "maximum_hours", + } + missing = sorted(required - set(profile)) + if missing: + raise RuntimeError("persona profile is missing: " + ", ".join(missing)) + return profile + + +def persona_token() -> str: + value = os.environ.get("CANGER_PERSONA_TOKEN", "") + if len(value) < 24: + raise RuntimeError( + "CANGER_PERSONA_TOKEN is unavailable; run from the registered " + "Singapore persona environment" + ) + return value + + +def parse_args() -> argparse.Namespace: + parser = argparse.ArgumentParser( + description="苍耳人格体的受控运维申请与执行入口" + ) + parser.add_argument("--profile", default=str(DEFAULT_PROFILE)) + parser.add_argument("--url") + parser.add_argument( + "--requester-name", + default=os.environ.get("GUANGHU_PERSONA_NAME", ""), + help="本次实际申请者名称;也可用 GUANGHU_PERSONA_NAME", + ) + parser.add_argument( + "--requester-id", + default=os.environ.get("GUANGHU_PERSONA_ID", ""), + help="本次实际申请者正式编号;也可用 GUANGHU_PERSONA_ID", + ) + sub = parser.add_subparsers(dest="command", required=True) + + request_access = sub.add_parser("request-access") + request_access.add_argument("--hours", type=int) + request_access.add_argument("--reason", required=True) + + task = sub.add_parser("task") + task.add_argument("--action", required=True) + task.add_argument("--args", default="{}") + + show = sub.add_parser("show") + show.add_argument("task_id") + + watch = sub.add_parser("watch") + watch.add_argument("task_id") + + sub.add_parser("grants") + return parser.parse_args() + + +def main() -> None: + args = parse_args() + profile = load_profile(args.profile) + base = args.url or os.environ.get("CANGER_CONTROL_URL") or profile["control_url"] + token = persona_token() + node_id = profile["node_id"] + if args.command in {"request-access", "task"}: + if not args.requester_name or not args.requester_id: + raise RuntimeError( + "requester name and id are required; provide --requester-name " + "and --requester-id, or GUANGHU_PERSONA_NAME/GUANGHU_PERSONA_ID" + ) + + if args.command == "request-access": + hours = args.hours or int(profile["default_hours"]) + minimum = int(profile["minimum_hours"]) + maximum = int(profile["maximum_hours"]) + if not minimum <= hours <= maximum: + raise RuntimeError(f"hours must be between {minimum} and {maximum}") + result = request( + base, + token, + "POST", + "/v1/grants", + { + "node_id": node_id, + "duration_seconds": hours * 3600, + "reason": args.reason, + "requester_name": args.requester_name, + "requester_id": args.requester_id, + }, + ) + elif args.command == "task": + result = request( + base, + token, + "POST", + "/v1/tasks", + { + "node_id": node_id, + "action": args.action, + "args": json.loads(args.args), + "requester_name": args.requester_name, + "requester_id": args.requester_id, + }, + ) + elif args.command == "show": + result = request( + base, + token, + "GET", + f"/v1/tasks/{args.task_id}", + ) + elif args.command == "grants": + result = request(base, token, "GET", "/v1/grants") + else: + after = 0 + while True: + events = request( + base, + token, + "GET", + f"/v1/tasks/{args.task_id}/events?after={after}", + ) + for event in events["events"]: + after = max(after, event["seq"]) + print( + f"[{event['seq']}] {event['kind']}: " + + json.dumps(event["payload"], ensure_ascii=False), + flush=True, + ) + state = request( + base, + token, + "GET", + f"/v1/tasks/{args.task_id}", + ) + if state["status"] in {"completed", "rejected"}: + result = state + break + time.sleep(1) + + print(json.dumps(result, ensure_ascii=False, indent=2)) + + +if __name__ == "__main__": + try: + main() + except HTTPError as exc: + print(f"error: control plane returned HTTP {exc.code}", file=sys.stderr) + raise SystemExit(2) + except URLError as exc: + print(f"error: cannot reach control plane: {exc.reason}", file=sys.stderr) + raise SystemExit(2) + except (RuntimeError, ValueError, json.JSONDecodeError, OSError) as exc: + print(f"error: {exc}", file=sys.stderr) + raise SystemExit(2) diff --git a/remote-node/tools/cangerctl.py b/remote-node/tools/cangerctl.py new file mode 100755 index 0000000..7952f9a --- /dev/null +++ b/remote-node/tools/cangerctl.py @@ -0,0 +1,168 @@ +#!/usr/bin/env python3 +"""Small CLI used by the persona and owner on the Singapore server.""" + +from __future__ import annotations + +import argparse +import json +import os +import sys +import time +from typing import Any +from urllib.error import HTTPError +from urllib.request import Request, urlopen + + +def request( + base: str, token: str, method: str, path: str, body: dict[str, Any] | None = None +) -> dict[str, Any]: + data = None if body is None else json.dumps(body).encode("utf-8") + headers = {"Authorization": f"Bearer {token}", "Accept": "application/json"} + if data is not None: + headers["Content-Type"] = "application/json" + req = Request(base.rstrip("/") + path, data=data, headers=headers, method=method) + try: + with urlopen(req, timeout=30) as response: + return json.load(response) + except HTTPError as exc: + print(exc.read().decode("utf-8", errors="replace"), file=sys.stderr) + raise + + +def token_for(role: str) -> str: + name = { + "admin": "CANGER_ADMIN_TOKEN", + "owner": "CANGER_OWNER_TOKEN", + "persona": "CANGER_PERSONA_TOKEN", + "operator": "CANGER_OPERATOR_TOKEN", + }[role] + value = os.environ.get(name, "") + if not value: + raise RuntimeError(f"{name} is not set") + return value + + +def parse_args() -> argparse.Namespace: + parser = argparse.ArgumentParser() + parser.add_argument( + "--url", default=os.environ.get("CANGER_CONTROL_URL", "http://127.0.0.1:8787") + ) + sub = parser.add_subparsers(dest="command", required=True) + sub.add_parser("pair-code") + sub.add_parser("nodes") + task = sub.add_parser("task") + task.add_argument("--node", required=True) + task.add_argument("--action", required=True) + task.add_argument("--args", default="{}") + approve = sub.add_parser("approve") + approve.add_argument("task_id") + reject = sub.add_parser("reject") + reject.add_argument("task_id") + show = sub.add_parser("show") + show.add_argument("task_id") + watch = sub.add_parser("watch") + watch.add_argument("task_id") + grant = sub.add_parser("grant") + grant.add_argument("--node", required=True) + grant.add_argument("--hours", type=int, default=4) + grant.add_argument("--reason", default="苍耳人格体编程、调试与诊断") + sub.add_parser("grants") + revoke = sub.add_parser("revoke") + revoke.add_argument("grant_id") + return parser.parse_args() + + +def main() -> None: + args = parse_args() + if args.command == "pair-code": + result = request(args.url, token_for("admin"), "POST", "/v1/pairings", {}) + elif args.command == "nodes": + result = request(args.url, token_for("operator"), "GET", "/v1/nodes") + elif args.command == "task": + result = request( + args.url, + token_for("operator"), + "POST", + "/v1/tasks", + { + "node_id": args.node, + "action": args.action, + "args": json.loads(args.args), + }, + ) + elif args.command == "approve": + result = request( + args.url, + token_for("owner"), + "POST", + f"/v1/tasks/{args.task_id}/approve", + {}, + ) + elif args.command == "reject": + result = request( + args.url, + token_for("owner"), + "POST", + f"/v1/tasks/{args.task_id}/reject", + {}, + ) + elif args.command == "show": + result = request( + args.url, + token_for("operator"), + "GET", + f"/v1/tasks/{args.task_id}", + ) + elif args.command == "watch": + after = 0 + while True: + events = request( + args.url, + token_for("operator"), + "GET", + f"/v1/tasks/{args.task_id}/events?after={after}", + ) + for event in events["events"]: + after = max(after, event["seq"]) + print( + f"[{event['seq']}] {event['kind']}: " + + json.dumps(event["payload"], ensure_ascii=False), + flush=True, + ) + state = request( + args.url, + token_for("operator"), + "GET", + f"/v1/tasks/{args.task_id}", + ) + if state["status"] in {"completed", "rejected"}: + result = state + break + time.sleep(1) + elif args.command == "grant": + result = request( + args.url, + token_for("operator"), + "POST", + "/v1/grants", + { + "node_id": args.node, + "duration_seconds": args.hours * 3600, + "reason": args.reason, + }, + ) + elif args.command == "grants": + result = request(args.url, token_for("operator"), "GET", "/v1/grants") + else: + result = request( + args.url, + token_for("owner"), + "POST", + f"/v1/grants/{args.grant_id}/revoke", + {}, + ) + print(json.dumps(result, ensure_ascii=False, indent=2)) + + +if __name__ == "__main__": + main()