From 605e4eecfeeb34df5549a45b87655ead01203d97 Mon Sep 17 00:00:00 2001 From: caoqianming Date: Wed, 2 Sep 2026 08:15:31 +0800 Subject: [PATCH] =?UTF-8?q?feat(sandbox):=20=E5=A2=9E=E5=8A=A0=E5=85=B1?= =?UTF-8?q?=E4=BA=AB=E6=89=A7=E8=A1=8C=E5=AE=B9=E9=87=8F=E4=B8=8E=E4=BE=9D?= =?UTF-8?q?=E8=B5=96=E8=A7=82=E6=B5=8B?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- CHANGELOG.md | 2 + DESIGN.md | 12 +- PROGRESS.md | 8 +- RUN.md | 25 +- config/agent.yaml | 12 +- core/executor_docker.py | 51 ++-- core/procs.py | 111 ++++++++- core/sandbox/capacity.py | 222 ++++++++++++++++++ core/sandbox/package_scans.py | 57 +++++ core/sandbox/pool.py | 104 +++++++- core/storage/models.py | 24 ++ ...0260901_1000_0038_sandbox_package_scans.py | 37 +++ deploy/sandbox/Dockerfile | 4 + deploy/sandbox/package_scan.py | 84 +++++++ tests/test_executor_docker.py | 13 + tests/test_sandbox_capacity_packages.py | 183 +++++++++++++++ tests/test_web_routes_nodb.py | 2 + tools/check_process.py | 16 +- web/admin.py | 49 +++- web/background.py | 9 +- web/static/js/admin.js | 47 ++++ 21 files changed, 1000 insertions(+), 72 deletions(-) create mode 100644 core/sandbox/capacity.py create mode 100644 core/sandbox/package_scans.py create mode 100644 db/migrations/versions/20260901_1000_0038_sandbox_package_scans.py create mode 100644 deploy/sandbox/package_scan.py create mode 100644 tests/test_sandbox_capacity_packages.py diff --git a/CHANGELOG.md b/CHANGELOG.md index 884a68e..b91d93b 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -8,6 +8,8 @@ ## Unreleased +- Sandbox 长任务会在整机繁忙时安全排队,后台任务可查看或取消排队状态;管理员可实时查看执行容量,并分析用户临时安装的 Python 依赖及占用。 + - 软件作业与轮次导航统一收在对话右侧,并增加醒目的数量徽章;既可悬停速览,也可互斥固定展开,进行中或未读作业会直接通过数字提醒。 - DeepSeek Flash 的思考过程恢复实时显示,不再等到正文开始后才集中出现。 diff --git a/DESIGN.md b/DESIGN.md index 12a5436..405e8e4 100644 --- a/DESIGN.md +++ b/DESIGN.md @@ -253,11 +253,11 @@ scheduled_jobs(§8.5) channel_bindings(§8.7,判别列+JSONB) ### 7.5 沙盒:Per-user 容器 + Per-tool exec -选型:**每 user 长驻容器**(文件模型本就以 user root 为安全边界,per-task 会切碎共享工作区)+ **每 tool 一次 docker exec**(exec 级 timeout/cwd/统计)+ 空闲 5 分钟回收 + bind mount user root→`/workspace`。 +选型:**每 user 长驻容器**(文件模型本就以 user root 为安全边界,per-task 会切碎共享工作区)+ **每 tool 一次 docker exec**(exec 级 timeout/cwd/统计)+ 无执行活动 10 分钟回收 + bind mount user root→`/workspace`。 **边界划分**:Control plane 留宿主(auth/DB/files 校验/SSE/LLM/受控 web 工具/配额审计),Execution plane 进容器(shell/run_python/任意生成代码)。目标不是"所有操作进容器",是"所有不可信执行不能在宿主"——否则凭据反被带进执行面。 -**硬限制**:cgroup CPU/mem、pids-limit、exec timeout、并发数、read-only rootfs、tmpfs /tmp、no-new-privileges、drop ALL caps、非 root、`--shm-size`。**软配额**:按 user 计 DB(磁盘/LLM cost/wall time/流量/并发),超额 429。**网络**:默认 deny outbound,搜索抓取走宿主受控工具。 +**硬限制**:单容器 4 GiB/2 CPU、pids 1024、`/dev/shm` 512 MiB、`/tmp` 1 GiB、exec timeout、read-only rootfs、no-new-privileges、drop ALL caps、非 root。整机重型执行由宿主共享文件账本 + advisory lock 统一准入:前后台合计 6、后台最多 4、单用户合计 2;蓝绿/多 Web 实例共享同一账本,前台排队不计 timeout 且可取消,宿主 MemAvailable 低于阈值时暂停新放行、不杀存量。轻量 fs 工具不占重型槽。普通容器显式跟踪 active exec,reaper 只回收无 active exec 且超 TTL 的容器。**软配额**:按 user 计 DB(磁盘/LLM cost/wall time/流量),超额 429。**网络**:默认 deny outbound,搜索抓取走宿主受控工具。 **落地清单(Stage C 硬协议,实施按此对账)**: 1. **网络 blocklist 硬编码段**(任一缺失=未完成):`169.254/16`(metadata)、内网三段、CGNAT `100.64/10`;**PG 实际 IP 单独再 block**(belt-and-suspenders)。**容器自身 loopback(`-o lo`)显式放行**——netns 隔离下容器内 127.0.0.1 到不了宿主,DROP 它无安全收益且误伤容器内 IPC(2026-07 实锤:puppeteer↔chromium DevTools 走 127.0.0.1,被 DROP 导致 mermaid 渲染 90 天 0 成功)。 @@ -276,7 +276,7 @@ scheduled_jobs(§8.5) channel_bindings(§8.7,判别列+JSONB) | gVisor → Firecracker/e2b | 合规客户 / 单机 100+ user / 兼容墙 | 每 VM 100MB+ 不划算;e2b 与 storage_root 自持冲突 | | docker exec → 容器内 tool-runner RPC | exec 开销 >30% 持续两周 / 长驻服务 / 单轮 >20 次调用 | 自管清理+观测损失 >> 200ms×N;美学统一 ≠ 理由 | -**Image 体积 / 多 user 资源 / 加包**(2026-05-28):① image 大 ≠ 运行时吃资源(layer 共享、不 exec 只是磁盘字节);② 瓶颈在并发 exec 不在 idle 容器,杠杆全在运行时限制;③ 新增依赖 = base 收敛 + **per-user venv**(`/.venv/`,bind mount 回收不丢;不放共享 volume——install 脚本是任意代码,破坏隔离)+ 使用频次沉淀进 base。 +**Image 体积 / 多 user 资源 / 加包**(2026-09-01 修订):① image 大 ≠ 运行时吃资源(layer 共享、不 exec 只是磁盘字节);② 瓶颈在并发 exec 不在 idle 容器,杠杆在共享准入与 cgroup;③ requirements 保持全局 site-packages,非 root + readonly rootfs + `HOME=/tmp` 使普通 pip 临时落到 `/tmp/.local`、cache 落 `/tmp/.cache`。不建 venv、不持久化用户包、不加额外包目录;容器删除前由镜像内可信扫描器只解析 `.dist-info`,与构建期基础清单比较 added/override/reinstall,空扫描不入库、失败不阻塞删除。Admin 只分析依赖会话,不自动改基础镜像。 ### 7.6 / 7.7 改造项与阶段 @@ -407,15 +407,15 @@ scheduled_jobs(§8.5) channel_bindings(§8.7,判别列+JSONB) **决策**:把长进程从 zcbot 进程树上摘下来,OS 就是"任务组件",不引 Celery/RQ/队列。`shell`/`run_python` 加 `background=true`(**模型判断**:预计 >~1min 走后台;前台超时报错里提示改后台——判断错了纠错路径只有一步;用户显式指令永远优先)。状态协议**纯文件**:`/.zcbot_procs///{proc.json, output.log, exit_code}`——dotfile 用户不可见(同 `.zcbot_tmp` 惯例),`exit_code` 文件出现是唯一终态信号,状态判定全靠文件 + 现场探测(pid / 容器 running),**无常驻登记,天然扛重启**。查询/终止走配套 `check_process` 工具(host in-process,两种 backend 通吃)。 - **host backend**:detach 独立 wrapper(`core/proc_wrapper.py`,stdlib-only,sys.executable 直跑不依赖 PYTHONPATH):限时(默认 7200s / cap 86400s)、超时杀进程树记 124、日志截尾 10MB、最后写 exit_code。 -- **docker backend**:**专用容器** `zcbot-proc-`(pool.run_proc_container,同款硬化 + iptables init),`product=proc` + 无 instance label —— 与 sandbox 容器的 idle reaper / shutdown_all 生命周期**解耦**,dockerd 托管,蓝绿切换/实例重启不中断。不用 `docker exec -d` 进 sandbox 容器:idle 5min reaper + 启动 shutdown_all 会把长进程随容器带走。 -- **回收**:check_process 见终态顺手 rm 容器;web lifespan 每小时 `procs.sweep`(终态目录 7d TTL / exited 孤儿容器),幂等,蓝绿双实例同时跑无害。 +- **docker backend**:**专用容器** `zcbot-proc-`(pool.run_proc_container,同款硬化 + iptables init),`product=proc` + 无 instance label —— 与普通 sandbox 的 idle reaper / shutdown_all 生命周期**解耦**,dockerd 托管,蓝绿切换/实例重启不中断。无容量时先写 `state=queued`,调度器获得共享前后台槽后才启动,故排队时间不计运行 timeout;`check_process` 兼容 queued/running/finished/lost,并可取消 queued。 +- **回收**:后台仍运行即有活动,不适用普通容器 idle TTL;check_process 见终态顺手扫描临时包并 rm 容器,web lifespan 约 30 秒 sweep/dispatch(终态目录仍保留 7d),幂等,蓝绿双实例以 proc 目录锁避免重复启动。运行中容器由 Docker 托管,服务重启恢复 queued 并保留 running。 - **通知/可视**:不做服务端推送 —— 前端轮询 `GET /v1/procs`(用户级,纯文件读取,仅有 running proc 时 5s 一拉):`[Background]` 工具结果卡本身活化(spinner+跳秒+停止按钮,与前台工具卡同体验,历史重渲同样恢复;`POST .../procs//kill`)、running→终态弹 toast(跨 task 也提醒,点击跳转)。proc 完成时刻往往没有活跃 run,SSE 通道根本不在,轮询是诚实的选型。 ### 8.13 扫描件 PDF 直读:方舟文档理解,不接外部 OCR(✅ 2026-07-21) 缺口:markitdown 只抽 PDF 文本层,扫描件(老标准/检测报告/红头指南,建材院高频)转出为空=死路。**选复用 seed-2.0-lite 的方舟文档理解**(chat file 内容块,PDF 整本 base64 内联)新增 `read_document`:零新供应商(敏感文档不出已有豆包面)、零新基础设施、记账复用 vision 通道;实测 ~1300 输入 token/页(约 1 厘/页)、100 页全覆盖、17MB 内联可用。**不选专用解析 API**(MinerU/Textin:版面还原最好,但申报书/专利底稿要上传新第三方 + 免费额度政策不稳);**不选本地 OCR**(PaddleOCR 类:镜像塞推理依赖,需求未量化前过度投资);**不选 file_url/file_id 传址**(前者要给用户文件开免认证公网直链=新安全面、开发机 NAT 后还跑不通;后者要接 TOS 多落一份存储;base64 是零新增面的唯一形态,行业惯例 chat 端点也不收 multipart)。防上下文爆:多页 OCR 强制 `save_md` 落盘只返预览。**升级信号**:>100 页/>30MB 巨件成高频 → 接 TOS 走 file_id;要高保真版面/公式还原 → 再评 MinerU。probe/smoke 留仓(`scripts/probe_ark_doc.py` / `smoke_read_document.py`)。 - **对话锁(前端)**:bg proc 运行期间该 task 的 composer 锁定(发送→停止,Enter 拦截),观感与前台执行完全一致 —— 后台化的收益定位为「进程扛超时/服务重启」,**不改变"一个任务同时只做一件事"的对话心智**;完成的那次轮询解锁 + toast「可继续对话」。锁只在前端,服务端不 409:「停止」入口必须可达,且多设备/渠道绕过前端锁属可接受边缘(等的是同一个进程,发了消息也不冲突)。 -- **防失控**:每用户并发 running 上限(`ZCBOT_MAX_BG_PROCS` 默 3);前台默认超时不放大(它是逼模型做前台/后台选择的杠杆)。 +- **防失控**:前后台统一使用整机 6 / 后台 4 / 单用户 2 的共享准入;前台默认超时不放大(它是逼模型做前台/后台选择的杠杆)。 **边界(防滑坡)**:只覆盖「单个本地长进程」。①**外部异步作业**(seedance 等 submit/poll 形态)不进这里——工具内轮询 + `resume_task_id` 续查已够;②**job 链/依赖/自动重试**不做——那是 workflow 引擎,编排的唯一归属是 agent loop(模型 check 后自己决定下一步),同 §6 拒绝编排的理由;③**完成后自动续跑 run**不做——zcbot 的长任务产物多为终点交付物(与 Claude Code"build 是中间步骤"不同),自动续跑=无人在场烧 token,通知给人、下一步由人/下次对话决定。 diff --git a/PROGRESS.md b/PROGRESS.md index 25193e7..4b70038 100644 --- a/PROGRESS.md +++ b/PROGRESS.md @@ -2,7 +2,7 @@ > 配合 `DESIGN.md`。本文件只记 phase 状态、决策偏差、文件量、下一步。每条 1-2 句:做了啥 + 关键判断;细节查 `git log` / `git diff` / `DESIGN §7.9`。 -最后更新:2026-08-31(Windows Node Host/WPF 重构阶段 0–7;未发版) +最后更新:2026-09-01(Sandbox 共享容量、后台队列与临时依赖观测;未发版) --- @@ -20,6 +20,8 @@ --- ## 已完成关键能力 +- **09-01 / Unreleased / Sandbox 执行容量与临时依赖观测**:默认单容器资源提升到 4 GiB/2 CPU、pids 1024、shm 512 MiB、tmp 1 GiB,普通容器无活动 10 分钟回收;新增跨蓝绿/多 Web 实例共享的整机 6、后台 4、单用户 2 重型槽位,前台可取消排队且排队不计 timeout,活动计数防长命令被 reaper 误删,内存压力只暂停新放行。Docker 后台 proc 无槽时持久化 queued、约 30 秒恢复调度,运行中仍由 Docker 托管,终态目录继续保留 7 天。镜像构建基础 Python 包清单,容器删除前可信扫描 `/tmp/.local` 的 `.dist-info` 并以单会话表幂等记录非空差异;Admin 新增实时容量与 7/30 天/全部依赖会话统计。新增 0038 migration,未连接或迁移生产数据库、未启动真实专业软件。 + - **08-31 / Unreleased / Windows Node 完整 WPF 管理界面**:本机任务、ANSYS 固定验收、诊断复制和数据目录迁移全部进入 WPF;任务页每秒刷新本地持久化记录并保持选择,验收支持进度、停止、通过报告校验与机器级执行门,数据迁移继续经过活动任务检查、停连、复制、SHA-256 校验、失败回滚和重启。新增 `JobMonitorService`、`AnsysAcceptanceService`、`DataRootManagementService` 与 `DiagnosticService`,删除经典 `ConfigurationForm` 和旧 `TrayApplicationContext`,WinForms 只保留托盘及受控系统对话框互操作;117 项 .NET 行为测试、28 项 Windows Node 源码专项、格式与 diff 检查通过,未启动专业软件、连接生产服务或读写真实节点数据。 - **08-31 / Unreleased / Windows Node Host/WPF 重构阶段 0–6 第二切片**:WPF 在节点概览、首次注册和运行设置基础上新增专业软件页,Origin、ANSYS、Blender 卡片通过 `SoftwareManagementService` 统一获取位置、runtime 和 probe 快照,可选择或恢复自动检测位置、安装/更新独立 Python 3.12 环境并取消安装。选择对话框仍限制在 Tray WinForms 互操作层,ViewModel 不访问注册表或启动进程;全机执行门会禁用安装,`ManagedRuntimeInstaller` 入口继续二次拒绝活动任务并保持临时验证、原子替换和失败 rollback。同步完成 WPF 视觉打磨:主题 token、圆角控件模板、导航选中态、状态色带/徽标、软件可用性标识、忙碌进度和窄窗口操作换行统一落地。ANSYS 验收、本机任务和数据目录迁移仍从经典配置进入。73 项 .NET 行为测试、191 项 Origin/ANSYS/Blender/合同/Node/UI 专项、格式与 diff 检查、Release 发布包及隔离 GUI/headless 烟测通过,未安装 runtime、启动专业软件、连接生产服务或读写真实节点数据。 @@ -456,10 +458,10 @@ core/ark_client.py 105 ← 火山方舟 HTTP 客户端 core/asr_xfyun.py 170 ← 讯飞语音听写 IAT wss 客户端(整段 PCM→文本;web 语音输入用,diag: scripts/diag_asr.py) core/asr_lfasr.py 250 ← 讯飞录音文件转写 LFASR 客户端(异步订单 + 说话人分离;transcribe_audio 工具底座,diag: scripts/diag_lfasr.py) core/agent_builder.py 649 ← 装配 lib:build_agent/system prompt(工具注册块已迁 tool_registry) -core/executor.py / sandbox/{network,pool}.py / executor_docker.py ← Executor ABC + Docker per-user 容器池 +core/executor.py / sandbox/{network,pool,capacity,package_scans}.py / executor_docker.py / procs.py ← Executor ABC + Docker per-user 容器池 + 宿主共享执行槽 + bg proc 文件队列 + 临时依赖扫描 tools/{base,output,fs,shell,run_python,skill_tool,skill_authoring,media_common,seedream,seedance,gpt_image,look_at_image,read_document,image_ref,web_search,web_fetch,documents,materials_project,transcribe_audio,office_to_pdf,external_systems}.py ← media_common=媒体五工具共享原语;external_systems=host-side 外部系统元工具 main.py ~210 ← 入口:web / db / probe / user / sandbox check -db/migrations/versions/ 0001-0030 +db/migrations/versions/ 0001-0038 web/app.py ~210 ← 工厂 + lifespan 编排(07-23 拆分;路由在 routers/,协程在 background 等) web/routers/*.py ← 含 external_systems 用户连接与 software_nodes 节点路由 web/{background,scheduler_runner,wechat_runner}.py ← lifespan 后台协程按域析出 diff --git a/RUN.md b/RUN.md index ee441b3..7fa791c 100644 --- a/RUN.md +++ b/RUN.md @@ -93,9 +93,8 @@ # 单次 LLM 请求超时(秒),默 600(与 litellm 默认一致,显式化 + 可调)。长思考模型 # 被掐("600s 无字节 → run 标 error")再调大;流式正常出 chunk 不触发。 # ZCBOT_LLM_TIMEOUT_S=600 - # 后台进程(bg proc,DESIGN §8.12;shell/run_python background=true):每用户并发上限,默 3。 + # 后台进程与前台重型执行统一受 sandbox.max_* / ZCBOT_MAX_* 共享容量控制。 # 状态锚 /.zcbot_procs/,后台默认限时 7200s(工具 timeout 参数可调,cap 86400)。 - # ZCBOT_MAX_BG_PROCS=3 # 静态多文件 Web 项目预览会话 TTL(秒),默认 7 天,平台限制 5 分钟~30 天。 # 改后重启 web 生效;只影响新建/再次发布的预览。 # ZCBOT_WEB_PREVIEW_TTL_SECONDS=604800 @@ -678,12 +677,18 @@ sudo -u zcbot docker network create zcbot-sandbox-net # 容器 runtime(切 gVisor 用 runsc,Firecracker 用 kata;默 runc) # ZCBOT_SANDBOX_RUNTIME= # 空闲多少秒回收(默 300) -# ZCBOT_SANDBOX_IDLE_TTL=300 +# ZCBOT_SANDBOX_IDLE_TTL=600 # 资源限制(优先级 env > yaml `sandbox.*` > 默);改后重启 web 新起容器生效 -# ZCBOT_SANDBOX_MEMORY=2g -# ZCBOT_SANDBOX_CPUS=1.0 -# ZCBOT_SANDBOX_PIDS_LIMIT=256 +# ZCBOT_SANDBOX_MEMORY=4g +# ZCBOT_SANDBOX_CPUS=2.0 +# ZCBOT_SANDBOX_PIDS_LIMIT=1024 # ZCBOT_SANDBOX_SHM_SIZE=512m # chromium/mmdc 渲 mermaid 的 /dev/shm(docker 默 64MB 不够会挂超时) +# ZCBOT_SANDBOX_TMP_SIZE=1g +# ZCBOT_MAX_ACTIVE_EXECS=6 +# ZCBOT_MAX_BACKGROUND_EXECS=4 +# ZCBOT_MAX_ACTIVE_EXECS_PER_USER=2 +# 三个并发 env 只用于向下收紧,代码硬上限固定为 6 / 4 / 2。 +# ZCBOT_MIN_MEM_AVAILABLE=1g # 低于阈值暂停新放行,不杀正在执行的任务 # PG 实际 IP,逗号分隔。defense-in-depth ── 即便落内网三段(§7.5 #1), # init.sh 再加一遍 DROP 规则。生产部署必填。 ZCBOT_PG_IPS=10.1.2.3,10.1.2.4 @@ -705,7 +710,7 @@ ZCBOT_SANDBOX_BACKEND=docker .venv/bin/python main.py web # 触发任一 shell / run_python 消息后,容器应已起 sudo -u zcbot docker ps --filter label=zcbot.product=sandbox # 应看到 zcbot-sandbox-,STATUS = Up ... -# 5 分钟无新消息后 reaper 自动 rm +# 10 分钟无执行活动后 reaper 自动 rm ``` 也可直接起一个测试容器单验 hardening(不依赖 web 进程): @@ -717,10 +722,10 @@ sudo -u zcbot docker run -d \ --label zcbot.product=sandbox \ --label zcbot.user_id=$USER_ID \ --network zcbot-sandbox-net \ - --read-only --tmpfs /tmp:exec,size=512m,mode=1777 \ + --read-only --tmpfs /tmp:exec,size=1g,mode=1777 --shm-size=512m \ --cap-drop=ALL --cap-add=NET_ADMIN \ --security-opt=no-new-privileges \ - --pids-limit=256 --memory=2g --cpus=1.0 \ + --pids-limit=1024 --memory=4g --cpus=2.0 \ -v /opt/zcbot/workspace/users/$USER_ID:/workspace \ -e ZCBOT_PG_IPS=10.1.2.3 \ zcbot-sandbox:latest @@ -985,7 +990,7 @@ sudo xfs_quota -x -c "limit -p bhard=10g zcbot_" /opt | prod 想把 workspace 落独立数据盘 | **别用 env / 别指 ROOT 外绝对路径**(workspace 锚定 ROOT,ROOT 外会让文件面板 / agent / 新建 task 三家分叉)。用 **bind mount** 把 `/data/...` 接到 `ROOT/workspace`,逻辑路径不变,DB 不用改。详「workspace 落独立数据盘」段 | | 文件面板"目录尚未创建"但文件确实在 / agent 写的文件面板看不到 | workspace 被指到了 ROOT 外(旧 `ZCBOT_WORKSPACE_DIR` 绝对路径残留)→ 文件面板走 `resolve_workspace` 看一处、agent 走 DB `from_db_path`(锚 ROOT)看另一处。删掉 env、改用 bind mount(见上段),三家归一 | | `docker run zcbot-sandbox:latest` 报 `Unable to find image` | 镜像没 build。`sudo -u zcbot docker build -f deploy/sandbox/Dockerfile --build-arg HOST_UID=$(id -u zcbot) --build-arg HOST_GID=$(id -g zcbot) -t zcbot-sandbox:latest .` | -| 后台进程(bg proc)疑似残留 / 想手工排查 | docker 模式:`docker ps --filter label=zcbot.product=proc` 列在跑的 proc 容器,`docker rm -f zcbot-proc-` 手杀;host 模式:状态锚在 `/.zcbot_procs///`(proc.json 有 pid,exit_code 文件在 = 已结束)。web 进程每小时 sweep 自动回收(终态目录 7 天 TTL);正常终止走前端停止按钮 / `check_process(action="kill")` | +| 后台进程(bg proc)疑似残留 / 想手工排查 | docker 模式:`docker ps --filter label=zcbot.product=proc` 列在跑的 proc 容器;状态锚在 `/.zcbot_procs///`,`proc.json` 的 `state=queued` 表示等待共享容量,`exit_code` 存在表示已结束。web 约每 30 秒调度 queued 并回收终态容器,终态目录仍保留 7 天;正常终止/取消排队走前端停止按钮或 `check_process(action="kill")`。若手工 `docker rm -f`,状态会收敛为 lost。 | | 后台进程状态显示 lost | 进程没留退出码就没了 —— 宿主重启 / OOM killer / docker daemon 重启把它带走(bg proc 扛 zcbot 重启和蓝绿,但不扛宿主级重启)。output.log 保留到中断为止,需重新 background=true 发起 | | 镜像 build pip 报 `THESE PACKAGES DO NOT MATCH THE HASHES FROM THE REQUIREMENTS FILE`(本仓 requirements 未钉 hash) | **不是被篡改、也不是 require-hashes**:镜像 index 声明的 wheel hash 与它实际吐出的文件字节不符 = 该镜像存的文件损坏 / 截断(2026-06-03 腾讯源就这么坏过 litellm-1.87.0)。换源重 build:`PIP_INDEX_URL=https://pypi.tuna.tsinghua.edu.cn/simple/ sudo -E bash deploy/update.sh`。验真伪:`https://pypi.org/pypi///json` 看官方 sha256 是哪边对。与下面"版本滞后(Could not find)"是两回事 | | 镜像 build pip 报 `ReadTimeoutError: HTTPSConnectionPool(host='files.pythonhosted.org', ...)` | 境内访问 PyPI 抖动。加 `--build-arg PIP_INDEX_URL=https://pypi.tuna.tsinghua.edu.cn/simple/`(清华,现默认)或腾讯 / 阿里源,详 RUN.md「镜像构建」段。Dockerfile 已把 pip timeout 拉到 60s,主因仍是源不通而非超时 | diff --git a/config/agent.yaml b/config/agent.yaml index 1404485..8295f46 100644 --- a/config/agent.yaml +++ b/config/agent.yaml @@ -66,14 +66,20 @@ shutdown: cancel_grace_seconds: 15 # 超时转 cancel 后再给的退场宽限 # Sandbox 容器资源限制(docker run flag,env 可 override);改后重启 web 生效, -# 新起的容器用新值,已 running 的不变(idle 5min 回收后下次起)。 +# 新起的容器用新值,已 running 的不变(idle 10min 回收后下次起)。 sandbox: - memory: 2g # --memory (env: ZCBOT_SANDBOX_MEMORY) - cpus: 1.0 # --cpus (env: ZCBOT_SANDBOX_CPUS) + memory: 4g # --memory (env: ZCBOT_SANDBOX_MEMORY) + cpus: 2.0 # --cpus (env: ZCBOT_SANDBOX_CPUS) pids_limit: 1024 # --pids-limit (env: ZCBOT_SANDBOX_PIDS_LIMIT);线程也计数, # chromium headless 一次 ~150-200 线程 + 超时残留进程,256 会被 # 打满致 mmdc 渲 mermaid 必崩(pthread_create EAGAIN) shm_size: 512m # --shm-size (env: ZCBOT_SANDBOX_SHM_SIZE);chromium/mmdc 渲 mermaid 的 /dev/shm,默 64MB 不够会挂 + tmp_size: 1g # /tmp tmpfs (env: ZCBOT_SANDBOX_TMP_SIZE) + idle_ttl_seconds: 600 # 普通容器无 exec 活动 10 分钟回收 + max_active_execs: 6 # 宿主硬上限;前后台统一计数 + max_background_execs: 4 + max_active_execs_per_user: 2 + min_mem_available: 1g # 低于阈值暂停新放行,不杀已运行任务 # 容器 DNS server 显式配置(docker run --dns,容器 /etc/resolv.conf 直接写, # 绕过 docker daemon 上游 DNS 探测路径;腾讯云轻量 / 部分云上 daemon 探测 # systemd-resolved 上游会失败,导致 embedded DNS 127.0.0.11 forward 出去也跪)。 diff --git a/core/executor_docker.py b/core/executor_docker.py index 73ba5df..04efcb0 100644 --- a/core/executor_docker.py +++ b/core/executor_docker.py @@ -25,7 +25,7 @@ run_python tmp .py 落 host 侧 `/.zcbot_tmp//.py`(bin Cancel limitation(第一版接受): - docker exec 客户端断开后,容器内 server 端进程**不会**因此终止 —— 这是 docker 设计 -- 第一版只杀 docker CLI(Popen.kill());容器内残留进程靠 idle 5min reaper / 下次 +- 第一版只杀 docker CLI(Popen.kill());容器内残留进程靠 idle 10min reaper / 下次 ensure 时 rm -f 兜底 - 升级触发(§7.5 #3 PGID 协议):用户反馈"取消了但还在烧 CPU" / 多次 cancel 后 容器内进程堆积 → 启用「ZCBOT_EXEC_ID env + PGID 写文件 + 二次 exec kill」协议 @@ -145,8 +145,12 @@ class DockerExecutor(Executor): return ToolResult(content=f"[Error] unknown tool: {name}", exit_code=2) try: if name == "shell": + if not args.get("background"): + return self._call_heavy_foreground(name, args, ctx) return self._exec_shell(args, ctx) if name == "run_python": + if not args.get("background"): + return self._call_heavy_foreground(name, args, ctx) return self._exec_python(args, ctx) if name in FS_TOOLS: return self._exec_fs_tool(name, args, ctx) @@ -157,6 +161,17 @@ class DockerExecutor(Executor): ) return ToolResult(content=f"[Error] unhandled container tool: {name}", exit_code=2) + def _call_heavy_foreground(self, name: str, args: Dict[str, Any], ctx: ExecCtx) -> ToolResult: + """前台重型工具先排共享槽;排队不进入命令 timeout。""" + with self.pool.capacity.foreground(str(self.user_id), ctx.cancel_check) as admitted: + if not admitted: + return ToolResult(content="[Error] command cancelled by user while queued", exit_code=130) + self.pool.exec_started(self.user_id) + try: + return self._exec_shell(args, ctx) if name == "shell" else self._exec_python(args, ctx) + finally: + self.pool.exec_finished(self.user_id) + # ── shell ──────────────────────────────────────────────── def _exec_shell(self, args: Dict[str, Any], ctx: ExecCtx) -> ToolResult: @@ -248,7 +263,7 @@ class DockerExecutor(Executor): def _exec_background(self, name: str, args: Dict[str, Any], ctx: ExecCtx) -> ToolResult: """background=true:专用容器跑长进程,立即返回 proc_id。 - 为什么不是 `docker exec -d` 进 sandbox 容器:sandbox 有 idle 5min reaper + + 为什么不是 `docker exec -d` 进 sandbox 容器:sandbox 有 idle 10min reaper + 启动时 shutdown_all,长进程会随容器陪葬。专用容器 product=proc 与这两条 生命周期解耦(pool.run_proc_container),状态协议同 host 模式(core/procs.py): runner.sh 结束时写 exit_code,check_process/sweep 负责回收容器。 @@ -256,15 +271,6 @@ class DockerExecutor(Executor): from core import procs anchor = self.user_root - if procs.count_running(anchor) >= procs.MAX_RUNNING_PER_USER: - return ToolResult( - content=( - f"[Error] 已有 {procs.MAX_RUNNING_PER_USER} 个后台进程在跑(上限)。" - f"用 check_process 查看,等待完成或 kill 掉不需要的再启动。" - ), - exit_code=2, - ) - raw_timeout = args.get("timeout") fg_default = 60 if name == "shell" else 120 try: @@ -334,29 +340,16 @@ class DockerExecutor(Executor): "timeout_s": timeout_s, "created_at": time.strftime("%Y-%m-%dT%H:%M:%S"), "created_ts": time.time(), + "user_id": str(self.user_id), + "exec_user": self.exec_user, + "state": "queued", } procs.write_meta(d, meta) - try: - container = self.pool.run_proc_container(self.user_id, proc_id, str(d)) - except Exception as e: - return ToolResult(content=f"[Error] 后台容器启动失败: {e}", exit_code=1) - meta["container"] = container - procs.write_meta(d, meta) - - argv = self._docker_exec_argv( - container, extra_env=_sandbox_env({"PYTHONIOENCODING": "utf-8"}), detach=True - ) + ["bash", f"{cdir}/runner.sh"] - r = subprocess.run(argv, capture_output=True, text=True, timeout=60) - if r.returncode != 0: - subprocess.run(["docker", "rm", "-f", container], capture_output=True) - return ToolResult( - content=f"[Error] 后台进程启动失败: {(r.stderr or '').strip()[:300]}", - exit_code=1, - ) + started = procs.start_queued_docker(meta, d, self.pool) return ToolResult( content=( - f"[Background] 已启动后台进程 proc_id={proc_id}({display[:150]})," + f"[Background] 已{'启动' if started else '排队'}后台进程 proc_id={proc_id}({display[:150]})," f"最长运行 {timeout_s}s。\n" f"用 check_process(proc_id=\"{proc_id}\") 查进度和日志。进程独立于本轮对话运行," f"服务重启也不中断。现在可以继续其他工作;若无事可做,结束回合并告知用户稍后询问进度。" diff --git a/core/procs.py b/core/procs.py index 5423667..b87d871 100644 --- a/core/procs.py +++ b/core/procs.py @@ -45,9 +45,6 @@ PROCS_SUBDIR = ".zcbot_procs" DEFAULT_TIMEOUT_S = 7200 MAX_TIMEOUT_S = 86400 -# 每用户并发 bg 进程上限(跨 task 合计);防模型失控起一堆 -MAX_RUNNING_PER_USER = int(os.getenv("ZCBOT_MAX_BG_PROCS", "3")) - # 终态 proc 目录保留时长,sweep 超期删除 FINISHED_TTL_S = 7 * 86400 @@ -151,12 +148,14 @@ def _container_running(name: str) -> bool: def status_of(meta: Dict[str, Any], d: Path) -> Tuple[str, Optional[int]]: - """→ ("running"|"finished"|"lost", exit_code|None)。exit_code 文件是唯一终态锚。""" + """→ queued/running/finished/lost。exit_code 文件是唯一终态锚。""" try: raw = (d / "exit_code").read_text(encoding="utf-8").strip() return "finished", int(raw) except (OSError, ValueError): pass + if meta.get("state") == "queued": + return "queued", None if meta.get("backend") == "docker": if _container_running(str(meta.get("container") or "")): return "running", None @@ -304,10 +303,24 @@ def kill_proc(meta: Dict[str, Any], d: Path) -> str: st, _ = status_of(meta, d) if st == "finished": return "进程已结束,无需终止" + if st == "queued": + meta["state"] = "finished" + meta["killed"] = True + write_meta(d, meta) + (d / "exit_code").write_text("137", encoding="utf-8") + return "已取消排队" if meta.get("backend") == "docker": name = str(meta.get("container") or "") if name: + _scan_proc_container(meta, d, name) subprocess.run(["docker", "rm", "-f", name], capture_output=True, timeout=30) + try: + from core.sandbox import get_pool + pool = get_pool() + if pool is not None: + pool.capacity.release(meta.get("capacity_lease_id")) + except Exception: + pass else: wrapper_pid = int(meta.get("wrapper_pid") or 0) child_pid = int(meta.get("child_pid") or 0) @@ -338,6 +351,69 @@ def kill_proc(meta: Dict[str, Any], d: Path) -> str: return "已终止" +def start_queued_docker(meta: Dict[str, Any], d: Path, pool) -> bool: + """尝试为 queued proc 获取共享槽并启动;无槽返回 False,保持排队。""" + from core.file_store import interprocess_file_lock + with interprocess_file_lock(d / "dispatch.lock", timeout_seconds=None): + meta = read_meta(d) or meta + if status_of(meta, d)[0] != "queued": + return False + uid = str(meta.get("user_id") or "") + lease_id = f"bg-{meta.get('proc_id')}" + lease = pool.capacity.try_acquire(uid, "background", lease_id=lease_id, + proc_id=str(meta.get("proc_id"))) + if not lease: + return False + try: + from uuid import UUID + container = pool.run_proc_container(UUID(uid), str(meta["proc_id"]), str(d)) + cdir = f"/workspace/{PROCS_SUBDIR}/{meta['task_id']}/{meta['proc_id']}" + argv = ["docker", "exec", "--user", str(meta.get("exec_user") or "zcbot"), + "--workdir", str(meta.get("cwd") or "/workspace"), "-d", + "-e", "PYTHONPATH=/sandbox:/workspace", "-e", "HOME=/tmp", + "-e", "PYTHONIOENCODING=utf-8", container, "bash", f"{cdir}/runner.sh"] + r = subprocess.run(argv, capture_output=True, text=True, timeout=60) + if r.returncode != 0: + subprocess.run(["docker", "rm", "-f", container], capture_output=True) + raise RuntimeError((r.stderr or "").strip()[:300]) + meta.update({"container": container, "capacity_lease_id": lease, + "state": "running", "started_ts": time.time()}) + write_meta(d, meta) + return True + except Exception as exc: + pool.capacity.release(lease) + meta["last_start_error"] = f"{type(exc).__name__}: {exc}" + write_meta(d, meta) + return False + + +def dispatch_queued(user_root_base: Path, pool) -> int: + started = 0 + base = Path(user_root_base) + if not base.is_dir(): + return 0 + candidates = [] + running_ids: set[str] = set() + for uroot in base.iterdir(): + root = uroot / PROCS_SUBDIR + if not root.is_dir(): + continue + for meta_file in root.glob("*/*/proc.json"): + d = meta_file.parent + meta = read_meta(d) + if meta: + status = status_of(meta, d)[0] + if status == "queued": + candidates.append((float(meta.get("created_ts") or 0), meta, d)) + elif status == "running": + running_ids.add(str(meta.get("proc_id") or d.name)) + pool.capacity.reconcile_background(running_ids) + for _, meta, d in sorted(candidates, key=lambda x: x[0]): + if start_queued_docker(meta, d, pool): + started += 1 + return started + + # ───────────── 清扫(web lifespan 周期调用) ───────────── def sweep(user_root_base: Path, ttl_s: int = FINISHED_TTL_S) -> Dict[str, int]: @@ -365,16 +441,27 @@ def sweep(user_root_base: Path, ttl_s: int = FINISHED_TTL_S) -> Dict[str, int]: if not meta: continue st, _ = status_of(meta, d) - if st != "running" and meta.get("backend") == "docker": + if st in {"finished", "lost"} and meta.get("capacity_lease_id"): + try: + from core.sandbox import get_pool + pool = get_pool() + if pool is not None: + pool.capacity.release(meta.get("capacity_lease_id")) + meta.pop("capacity_lease_id", None) + write_meta(d, meta) + except Exception: + pass + if st not in {"running", "queued"} and meta.get("backend") == "docker": name = str(meta.get("container") or "") if name and _container_exists(name): + _scan_proc_container(meta, d, name) subprocess.run( ["docker", "rm", "-f", name], capture_output=True, timeout=30, ) reaped_containers += 1 age = now - float(meta.get("created_ts") or now) - if st != "running" and age > ttl_s: + if st not in {"running", "queued"} and age > ttl_s: shutil.rmtree(d, ignore_errors=True) removed_dirs += 1 try: @@ -410,6 +497,16 @@ def _container_exists(name: str) -> bool: return False +def _scan_proc_container(meta: Dict[str, Any], d: Path, name: str) -> None: + try: + from uuid import UUID + from core.sandbox.package_scans import scan_and_persist + uid = UUID(str(meta.get("user_id"))) + scan_and_persist(name, uid, "background", str(meta.get("proc_id") or d.name)) + except Exception as exc: + print(f"[sandbox-package-scan] proc finalizer failed container={name}: {type(exc).__name__}: {exc}") + + # ───────────── LLM 面向的描述格式化(check_process / 启动返回共用) ───────────── def format_proc_line(meta: Dict[str, Any]) -> str: @@ -418,7 +515,7 @@ def format_proc_line(meta: Dict[str, Any]) -> str: elapsed = "" ts = meta.get("created_ts") if ts: - elapsed = f" · 已运行 {_fmt_elapsed(time.time() - float(ts))}" if st == "running" else "" + elapsed = f" · 已运行 {_fmt_elapsed(time.time() - float(meta.get('started_ts') or ts))}" if st == "running" else "" tail = f"(exit {ec})" if st == "finished" else "" return ( f"- {meta.get('proc_id')} [{st}{tail}] {meta.get('kind')}: " diff --git a/core/sandbox/capacity.py b/core/sandbox/capacity.py new file mode 100644 index 0000000..d742788 --- /dev/null +++ b/core/sandbox/capacity.py @@ -0,0 +1,222 @@ +"""宿主级 Sandbox 重型执行容量。 + +状态保存在 workspace/.sandbox 下并由 advisory file lock 串行化,因此蓝绿和多 +Web 进程共享同一组槽位。业务 DB 不承载实时状态;后台排队事实源仍是 +``.zcbot_procs``,这里只记录已经获得槽位的租约和短暂的前台候选。 +""" +from __future__ import annotations + +import json +import os +import time +import uuid +from contextlib import contextmanager +from pathlib import Path +from typing import Callable, Dict, Iterator, Optional + +from core.file_store import atomic_write_text, interprocess_file_lock + +DEFAULT_MAX_ACTIVE_EXECS = 6 +DEFAULT_MAX_BACKGROUND_EXECS = 4 +DEFAULT_MAX_ACTIVE_EXECS_PER_USER = 2 +DEFAULT_MIN_MEM_AVAILABLE_BYTES = 1024 ** 3 + + +def _parse_bytes(value: object, default: int) -> int: + text = str(value or "").strip().lower() + if not text: + return default + units = {"k": 1024, "kb": 1024, "m": 1024**2, "mb": 1024**2, + "g": 1024**3, "gb": 1024**3} + for suffix, factor in sorted(units.items(), key=lambda x: -len(x[0])): + if text.endswith(suffix): + return int(float(text[:-len(suffix)]) * factor) + return int(text) + + +def mem_available_bytes() -> Optional[int]: + try: + for line in Path("/proc/meminfo").read_text(encoding="ascii").splitlines(): + if line.startswith("MemAvailable:"): + return int(line.split()[1]) * 1024 + except (OSError, ValueError, IndexError): + return None + return None + + +class ExecCapacity: + def __init__(self, state_dir: Path, cfg: Optional[dict] = None) -> None: + cfg = cfg or {} + self.state_dir = Path(state_dir) + self.state_path = self.state_dir / "exec-capacity.json" + self.lock_path = self.state_dir / "exec-capacity.lock" + self.max_active = max(1, min(DEFAULT_MAX_ACTIVE_EXECS, int(os.getenv("ZCBOT_MAX_ACTIVE_EXECS") or cfg.get("max_active_execs") or DEFAULT_MAX_ACTIVE_EXECS))) + self.max_background = max(1, min(DEFAULT_MAX_BACKGROUND_EXECS, int(os.getenv("ZCBOT_MAX_BACKGROUND_EXECS") or cfg.get("max_background_execs") or DEFAULT_MAX_BACKGROUND_EXECS))) + self.max_per_user = max(1, min(DEFAULT_MAX_ACTIVE_EXECS_PER_USER, int(os.getenv("ZCBOT_MAX_ACTIVE_EXECS_PER_USER") or cfg.get("max_active_execs_per_user") or DEFAULT_MAX_ACTIVE_EXECS_PER_USER))) + self.min_mem_available = _parse_bytes( + os.getenv("ZCBOT_MIN_MEM_AVAILABLE") or cfg.get("min_mem_available"), + DEFAULT_MIN_MEM_AVAILABLE_BYTES, + ) + + def _read(self) -> dict: + try: + data = json.loads(self.state_path.read_text(encoding="utf-8")) + if isinstance(data, dict): + data.setdefault("leases", {}) + data.setdefault("foreground_queue", []) + data.setdefault("containers", {}) + return data + except (OSError, ValueError): + pass + return {"leases": {}, "foreground_queue": [], "containers": {}} + + def _write(self, state: dict) -> None: + atomic_write_text(self.state_path, json.dumps(state, ensure_ascii=False, sort_keys=True)) + + def _prune(self, state: dict) -> None: + leases = state["leases"] + for key, lease in list(leases.items()): + if lease.get("kind") == "foreground": + pid = int(lease.get("owner_pid") or 0) + if pid and not self._pid_alive(pid): + leases.pop(key, None) + state["foreground_queue"] = [ + q for q in state["foreground_queue"] + if self._pid_alive(int(q.get("owner_pid") or 0)) + ] + + @staticmethod + def _pid_alive(pid: int) -> bool: + try: + os.kill(pid, 0) + return True + except PermissionError: + return True + except OSError: + return False + + def _can_admit(self, state: dict, user_id: str, kind: str) -> bool: + leases = list(state["leases"].values()) + if len(leases) >= self.max_active: + return False + if sum(1 for x in leases if x.get("user_id") == user_id) >= self.max_per_user: + return False + if kind == "background" and sum(1 for x in leases if x.get("kind") == "background") >= self.max_background: + return False + available = mem_available_bytes() + return available is None or available >= self.min_mem_available + + def try_acquire(self, user_id: str, kind: str, *, lease_id: Optional[str] = None, proc_id: Optional[str] = None) -> Optional[str]: + lease_id = lease_id or uuid.uuid4().hex + with interprocess_file_lock(self.lock_path, timeout_seconds=None): + state = self._read() + self._prune(state) + if lease_id in state["leases"]: + return None + if not self._can_admit(state, str(user_id), kind): + self._write(state) + return None + state["leases"][lease_id] = { + "lease_id": lease_id, "user_id": str(user_id), "kind": kind, + "proc_id": proc_id, "owner_pid": os.getpid(), "started_ts": time.time(), + } + self._write(state) + return lease_id + + def acquire_foreground(self, user_id: str, cancel_check: Optional[Callable[[], bool]] = None) -> Optional[str]: + ticket = uuid.uuid4().hex + with interprocess_file_lock(self.lock_path, timeout_seconds=None): + state = self._read(); self._prune(state) + state["foreground_queue"].append({"ticket": ticket, "user_id": str(user_id), "created_ts": time.time(), "owner_pid": os.getpid()}) + self._write(state) + while True: + with interprocess_file_lock(self.lock_path, timeout_seconds=None): + state = self._read(); self._prune(state) + queue = state["foreground_queue"] + eligible = next((q for q in queue if q.get("user_id") == str(user_id)), None) + if eligible and eligible.get("ticket") == ticket and self._can_admit(state, str(user_id), "foreground"): + queue[:] = [q for q in queue if q.get("ticket") != ticket] + state["leases"][ticket] = {"lease_id": ticket, "user_id": str(user_id), "kind": "foreground", "owner_pid": os.getpid(), "started_ts": time.time()} + self._write(state) + return ticket + self._write(state) + if cancel_check is not None and cancel_check(): + self.cancel_waiter(ticket) + return None + time.sleep(0.1) + + def cancel_waiter(self, ticket: str) -> None: + with interprocess_file_lock(self.lock_path, timeout_seconds=None): + state = self._read() + state["foreground_queue"] = [q for q in state["foreground_queue"] if q.get("ticket") != ticket] + self._write(state) + + def release(self, lease_id: Optional[str]) -> None: + if not lease_id: + return + with interprocess_file_lock(self.lock_path, timeout_seconds=None): + state = self._read(); state["leases"].pop(lease_id, None); self._write(state) + + def reconcile_background(self, running_proc_ids: set[str]) -> None: + """服务重启后保留 Docker 仍在跑的租约,清启动窗口遗留的孤儿租约。""" + with interprocess_file_lock(self.lock_path, timeout_seconds=None): + state = self._read() + for key, lease in list(state["leases"].items()): + if (lease.get("kind") == "background" + and str(lease.get("proc_id") or "") not in running_proc_ids + and not self._pid_alive(int(lease.get("owner_pid") or 0))): + state["leases"].pop(key, None) + self._write(state) + + def touch_container(self, name: str, *, active_delta: int = 0) -> None: + with interprocess_file_lock(self.lock_path, timeout_seconds=None): + state = self._read() + row = state["containers"].setdefault(name, {"active_execs": 0}) + row["active_execs"] = max(0, int(row.get("active_execs") or 0) + active_delta) + row["last_active_ts"] = time.time() + self._write(state) + + def remove_container(self, name: str) -> None: + with interprocess_file_lock(self.lock_path, timeout_seconds=None): + state = self._read(); state["containers"].pop(name, None); self._write(state) + + def reconcile_containers(self, running_names: set[str]) -> None: + with interprocess_file_lock(self.lock_path, timeout_seconds=None): + state = self._read() + for name in list(state["containers"]): + if name not in running_names: + state["containers"].pop(name, None) + self._write(state) + + @contextmanager + def foreground(self, user_id: str, cancel_check: Optional[Callable[[], bool]] = None) -> Iterator[bool]: + lease = self.acquire_foreground(user_id, cancel_check) + try: + yield lease is not None + finally: + self.release(lease) + + def snapshot(self) -> dict: + with interprocess_file_lock(self.lock_path, timeout_seconds=None): + state = self._read(); self._prune(state); self._write(state) + leases = list(state["leases"].values()) + fg = sum(x.get("kind") == "foreground" for x in leases) + bg = sum(x.get("kind") == "background" for x in leases) + by_user: Dict[str, int] = {} + for x in leases: + by_user[x["user_id"]] = by_user.get(x["user_id"], 0) + 1 + available_mem = mem_available_bytes() + memory_paused = available_mem is not None and available_mem < self.min_mem_available + containers = state.get("containers", {}) + now = time.time() + return { + "limits": {"active": self.max_active, "background": self.max_background, "per_user": self.max_per_user}, + "foreground_running": fg, "foreground_queued": len(state["foreground_queue"]), + "background_running": bg, "per_user": by_user, + "admit_available": 0 if memory_paused else max(0, self.max_active - len(leases)), + "memory_paused": memory_paused, + "mem_available_bytes": available_mem, + "cpu_load": list(os.getloadavg()) if hasattr(os, "getloadavg") else None, + "active_sandbox_containers": len(containers), + "idle_reap_candidates": sum(int(x.get("active_execs") or 0) == 0 and float(x.get("last_active_ts") or now) < now - getattr(self, "idle_ttl_seconds", 600) for x in containers.values()), + } diff --git a/core/sandbox/package_scans.py b/core/sandbox/package_scans.py new file mode 100644 index 0000000..005d86c --- /dev/null +++ b/core/sandbox/package_scans.py @@ -0,0 +1,57 @@ +"""Sandbox 临时 Python 包扫描与 best-effort 持久化。""" +from __future__ import annotations +import json +import subprocess +from datetime import datetime, timezone +from typing import Optional +from uuid import UUID + + +def scan_container(container: str) -> Optional[dict]: + try: + scan = subprocess.run(["docker", "exec", "--user", "zcbot", "--workdir", "/sandbox", container, + "python", "-I", "-S", "/sandbox/package_scan.py"], capture_output=True, + text=True, timeout=30) + if scan.returncode != 0: + raise RuntimeError(scan.stderr.strip()) + packages = json.loads(scan.stdout) + if not isinstance(packages, list): + raise ValueError("scanner output is not a package list") + inspect = subprocess.run(["docker", "inspect", "--format={{.Image}}", container], + capture_output=True, text=True, timeout=15) + py = subprocess.run(["docker", "exec", "--workdir", "/sandbox", container, "python", "-I", "-S", "-c", + "import sys;print('.'.join(map(str,sys.version_info[:3])))"], + capture_output=True, text=True, timeout=15) + if inspect.returncode != 0 or py.returncode != 0: + raise RuntimeError("failed to inspect image or Python version") + return {"packages": packages, "image_digest": inspect.stdout.strip(), + "python_version": py.stdout.strip(), + "total_installed_bytes": sum(int(x.get("installed_bytes") or 0) for x in packages)} + except Exception as exc: + print(f"[sandbox-package-scan] failed container={container}: {type(exc).__name__}: {exc}") + return None + + +def persist_scan(*, container_session_id: str, user_id: UUID, execution_kind: str, + result: Optional[dict]) -> None: + if not result or not result.get("packages"): + return + try: + from core.storage import session_scope + from core.storage.models import SandboxPackageScan + with session_scope() as s: + s.merge(SandboxPackageScan( + container_session_id=container_session_id, user_id=user_id, + execution_kind=execution_kind, image_digest=result["image_digest"], + python_version=result["python_version"], packages=result["packages"], + package_count=len(result["packages"]), + total_installed_bytes=result["total_installed_bytes"], + finished_at=datetime.now(timezone.utc), + )) + except Exception as exc: + print(f"[sandbox-package-scan] persist failed session={container_session_id}: {type(exc).__name__}: {exc}") + + +def scan_and_persist(container: str, user_id: UUID, kind: str, session_id: str) -> None: + persist_scan(container_session_id=session_id, user_id=user_id, + execution_kind=kind, result=scan_container(container)) diff --git a/core/sandbox/pool.py b/core/sandbox/pool.py index fb10e8b..f245536 100644 --- a/core/sandbox/pool.py +++ b/core/sandbox/pool.py @@ -56,11 +56,11 @@ LABEL_INSTANCE_KEY = "zcbot.instance" INSTANCE = os.getenv("ZCBOT_INSTANCE", "").strip() DEFAULT_IMAGE = "zcbot-sandbox:latest" -DEFAULT_IDLE_TTL_SECONDS = 300 +DEFAULT_IDLE_TTL_SECONDS = 600 # 容器资源限制默认值(可被 yaml `sandbox.*` / env override,详 SandboxPool ctor) -DEFAULT_MEMORY = "2g" -DEFAULT_CPUS = "1.0" +DEFAULT_MEMORY = "4g" +DEFAULT_CPUS = "2.0" # pids-limit 把**线程**也计数:chromium headless 一次启动(browser+GPU+network+renderer # 多进程)就 ~150-200 线程,256 本来就贴边;shell 超时又只杀 host 侧 docker CLI,容器内 # mmdc+chromium 树留着(见 executor_docker.py 头注 Cancel limitation),残留几棵就把配额 @@ -70,8 +70,9 @@ DEFAULT_PIDS_LIMIT = 1024 # chromium(mmdc 渲 mermaid / puppeteer)默认走 /dev/shm,docker 不传 --shm-size 时 # 只给 64MB,起不来就一直挂到 timeout。镜像备的 puppeteer-config 有 --disable-dev-shm-usage, # 但模型不一定用那份;这里从根上把 /dev/shm 撑到够用,任何 chromium 路径都不再挂。 -# 从 --memory(默 2g)里切,512m 是上限非占用(tmpfs 按需用)。 +# 从 --memory(默 4g)里切,512m 是上限非占用(tmpfs 按需用)。 DEFAULT_SHM_SIZE = "512m" +DEFAULT_TMP_SIZE = "1g" def container_name(user_id: UUID) -> str: @@ -102,6 +103,11 @@ def _container_running(name: str) -> bool: return r.returncode == 0 and r.stdout.strip() == "true" +def _container_session_id(name: str) -> str: + r = subprocess.run(["docker", "inspect", "--format={{.Id}}", name], capture_output=True, text=True) + return r.stdout.strip() if r.returncode == 0 else name + + class SandboxPool: def __init__( self, @@ -115,7 +121,9 @@ class SandboxPool: cpus: Optional[str] = None, pids_limit: Optional[int] = None, shm_size: Optional[str] = None, + tmp_size: Optional[str] = None, dns: Optional[List[str]] = None, + capacity_cfg: Optional[Dict[str, object]] = None, ) -> None: """ user_root_base: per-user 子树父目录,典型 `/users`。bind mount 源 @@ -134,9 +142,9 @@ class SandboxPool: pg_ips: 逗号分隔的 PG IP 串,塞容器 `ZCBOT_PG_IPS` env,init.sh 加 DROP 规则 (env `ZCBOT_PG_IPS`)。defense-in-depth ── 即便落内网三段。 memory/cpus/pids_limit/shm_size: - 容器资源限制,默 2g/1.0/256/512m;env(`ZCBOT_SANDBOX_MEMORY` 等) + 容器资源限制,默 4g/2.0/1024/512m;env(`ZCBOT_SANDBOX_MEMORY` 等) override caller 参数 override 默认。改后重启 web 生效,新起的 - 容器用新值;已 running 不变(idle 5min 回收后下次起按新值)。 + 容器用新值;已 running 不变(idle 10min 回收后下次起按新值)。 shm_size 撑 chromium 的 /dev/shm(默 64MB 不够,mmdc 渲图会挂)。 """ self.user_root_base = user_root_base @@ -153,6 +161,7 @@ class SandboxPool: self.memory = os.getenv("ZCBOT_SANDBOX_MEMORY") or memory or DEFAULT_MEMORY self.cpus = os.getenv("ZCBOT_SANDBOX_CPUS") or cpus or DEFAULT_CPUS self.shm_size = os.getenv("ZCBOT_SANDBOX_SHM_SIZE") or shm_size or DEFAULT_SHM_SIZE + self.tmp_size = os.getenv("ZCBOT_SANDBOX_TMP_SIZE") or tmp_size or DEFAULT_TMP_SIZE self.pids_limit = int( os.getenv("ZCBOT_SANDBOX_PIDS_LIMIT") or (pids_limit if pids_limit is not None else DEFAULT_PIDS_LIMIT) @@ -166,6 +175,10 @@ class SandboxPool: self._dict_lock = threading.Lock() # 保护 _locks / _last_active 的字典级 race self._locks: Dict[UUID, threading.Lock] = {} self._last_active: Dict[UUID, int] = {} + self._active_execs: Dict[UUID, int] = {} + from .capacity import ExecCapacity + self.capacity = ExecCapacity(self.user_root_base.parent / ".sandbox", capacity_cfg) + self.capacity.idle_ttl_seconds = self.idle_ttl def _lock_for(self, user_id: UUID) -> threading.Lock: with self._dict_lock: @@ -179,15 +192,18 @@ class SandboxPool: name = container_name(user_id) if _container_running(name): self._last_active[user_id] = _now() + self.capacity.touch_container(name) return name if _container_exists(name): # stopped / crashed ── rm 重起。iptables 规则随容器生命周期重新 apply。 + self._scan_before_remove(name, user_id, "foreground") subprocess.run( ["docker", "rm", "-f", name], capture_output=True, check=False, ) self._docker_run(user_id, name) self._last_active[user_id] = _now() + self.capacity.touch_container(name) return name def _ensure_resolv_conf_file(self) -> Optional[Path]: @@ -255,7 +271,7 @@ class SandboxPool: "--network", NETWORK_NAME, # §7.5 硬限制(任一缺失视为 hardening 未完成) "--read-only", # rootfs read-only - "--tmpfs", "/tmp:exec,size=512m,mode=1777", # 可写临时区,exec 允许 (run_python 写脚本) + "--tmpfs", f"/tmp:exec,size={self.tmp_size},mode=1777", f"--shm-size={self.shm_size}", # chromium/mmdc 的 /dev/shm,默 64MB 不够会挂(DEFAULT_SHM_SIZE) "--cap-drop=ALL", # 默全丢 "--cap-add=NET_ADMIN", # init.sh 配 iptables 需要;exec 进来的 uid 1000 拿不到 @@ -305,20 +321,48 @@ class SandboxPool: def mark_active(self, user_id: UUID) -> None: """每次 `docker exec` 完调一次,刷新 idle 计时。""" self._last_active[user_id] = _now() + self.capacity.touch_container(container_name(user_id)) + + def exec_started(self, user_id: UUID) -> None: + """普通容器内 exec 开始;显式计数防止 idle reaper 误删长命令。""" + with self._dict_lock: + self._active_execs[user_id] = self._active_execs.get(user_id, 0) + 1 + self._last_active[user_id] = _now() + self.capacity.touch_container(container_name(user_id), active_delta=1) + + def exec_finished(self, user_id: UUID) -> None: + with self._dict_lock: + n = max(0, self._active_execs.get(user_id, 0) - 1) + if n: + self._active_execs[user_id] = n + else: + self._active_execs.pop(user_id, None) + self._last_active[user_id] = _now() + self.capacity.touch_container(container_name(user_id), active_delta=-1) + + def _scan_before_remove(self, name: str, user_id: UUID, kind: str, session_id: Optional[str] = None) -> None: + """扫描/落库均 best-effort,调用方无论如何继续删除容器。""" + try: + from .package_scans import scan_and_persist + scan_and_persist(name, user_id, kind, session_id or _container_session_id(name)) + except Exception as exc: + print(f"[sandbox-package-scan] finalizer failed container={name}: {type(exc).__name__}: {exc}") def reap_idle(self) -> List[str]: """杀超过 idle_ttl 没活跃的容器。返回已杀容器名列表(供日志 / 审计)。""" removed: List[str] = [] cutoff = _now() - self.idle_ttl for uid, ts in list(self._last_active.items()): - if ts < cutoff: + if ts < cutoff and self._active_execs.get(uid, 0) == 0: name = container_name(uid) + self._scan_before_remove(name, uid, "foreground") r = subprocess.run( ["docker", "rm", "-f", name], capture_output=True, text=True, ) if r.returncode == 0: removed.append(name) + self.capacity.remove_container(name) # 无论 rm 成功与否,从 dict 移除 ── 失败则下次启动靠 shutdown_all 兜底 del self._last_active[uid] return removed @@ -341,6 +385,21 @@ class SandboxPool: if list_r.returncode != 0 or not list_r.stdout.strip(): return [] ids = list_r.stdout.strip().splitlines() + for cid in ids: + info = subprocess.run( + ["docker", "inspect", "--format={{index .Config.Labels \"zcbot.user_id\"}}", cid], + capture_output=True, text=True, + ) + try: + uid = UUID(info.stdout.strip()) + except (ValueError, AttributeError): + continue + self._scan_before_remove(cid, uid, "foreground", cid) + try: + name_r = subprocess.run(["docker", "inspect", "--format={{.Name}}", cid], capture_output=True, text=True) + self.capacity.remove_container(name_r.stdout.strip().lstrip("/") or cid) + except Exception: + pass subprocess.run( ["docker", "rm", "-f", *ids], capture_output=True, text=True, @@ -349,6 +408,31 @@ class SandboxPool: self._last_active.clear() return ids + def runtime_snapshot(self) -> dict: + try: + listed = subprocess.run( + ["docker", "ps", "--format={{.Names}}", "--filter", f"label={LABEL_PRODUCT_KEY}={LABEL_PRODUCT_VALUE}"], + capture_output=True, text=True, timeout=15, + ) + if listed.returncode == 0: + self.capacity.reconcile_containers(set(listed.stdout.split())) + except (OSError, subprocess.TimeoutExpired): + pass + snap = self.capacity.snapshot() + now = _now() + snap["idle_ttl_seconds"] = self.idle_ttl + try: + from core import procs + items = [] + for uroot in self.user_root_base.iterdir(): + items.extend(procs.list_all_procs(uroot)) + snap["background_queued"] = sum(x.get("_status") == "queued" for x in items) + snap["running_proc_containers"] = sum(x.get("_status") == "running" and x.get("backend") == "docker" for x in items) + except OSError: + snap["background_queued"] = 0 + snap["running_proc_containers"] = 0 + return snap + def setup_pool( user_root_base: Path, @@ -376,12 +460,16 @@ def setup_pool( cpus = cfg.get("cpus") pids_limit = cfg.get("pids_limit") shm_size = cfg.get("shm_size") + tmp_size = cfg.get("tmp_size") return SandboxPool( user_root_base=user_root_base, repo_root=repo_root, + idle_ttl=int(str(cfg.get("idle_ttl_seconds"))) if cfg.get("idle_ttl_seconds") is not None else None, memory=memory if isinstance(memory, str) else None, cpus=str(cpus) if cpus is not None else None, pids_limit=int(str(pids_limit)) if pids_limit is not None else None, shm_size=shm_size if isinstance(shm_size, str) else None, + tmp_size=tmp_size if isinstance(tmp_size, str) else None, dns=[str(x) for x in dns_cfg], + capacity_cfg=cfg, ) diff --git a/core/storage/models.py b/core/storage/models.py index 6d1ef46..563c7e2 100644 --- a/core/storage/models.py +++ b/core/storage/models.py @@ -804,3 +804,27 @@ class ExternalSystemAudit(Base): created_at: Mapped[datetime] = mapped_column( DateTime(timezone=True), server_default=func.now(), nullable=False ) + + +class SandboxPackageScan(Base): + """一个发现临时 Python 包的容器会话一行(空扫描不入库)。""" + + __tablename__ = "sandbox_package_scans" + __table_args__ = ( + Index("ix_sandbox_package_scans_finished_at", "finished_at"), + Index("ix_sandbox_package_scans_user_finished", "user_id", "finished_at"), + ) + container_session_id: Mapped[str] = mapped_column(Text, primary_key=True) + user_id: Mapped[UUID] = mapped_column( + PG_UUID(as_uuid=True), ForeignKey("users.user_id", ondelete="CASCADE"), nullable=False + ) + execution_kind: Mapped[str] = mapped_column(Text, nullable=False) + image_digest: Mapped[str] = mapped_column(Text, nullable=False) + python_version: Mapped[str] = mapped_column(Text, nullable=False) + packages: Mapped[list[Any]] = mapped_column(JSONB, nullable=False) + package_count: Mapped[int] = mapped_column(Integer, nullable=False) + total_installed_bytes: Mapped[int] = mapped_column(BigInteger, nullable=False) + finished_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), nullable=False) + created_at: Mapped[datetime] = mapped_column( + DateTime(timezone=True), server_default=func.now(), nullable=False + ) diff --git a/db/migrations/versions/20260901_1000_0038_sandbox_package_scans.py b/db/migrations/versions/20260901_1000_0038_sandbox_package_scans.py new file mode 100644 index 0000000..077686e --- /dev/null +++ b/db/migrations/versions/20260901_1000_0038_sandbox_package_scans.py @@ -0,0 +1,37 @@ +"""Add sandbox package scan sessions. + +Revision ID: 0038 +Revises: 0037 +Create Date: 2026-09-01 +""" +from collections.abc import Sequence +import sqlalchemy as sa +from alembic import op +from sqlalchemy.dialects import postgresql + +revision: str = "0038" +down_revision: str | None = "0037" +branch_labels: str | Sequence[str] | None = None +depends_on: str | Sequence[str] | None = None + + +def upgrade() -> None: + op.create_table( + "sandbox_package_scans", + sa.Column("container_session_id", sa.Text(), primary_key=True), + sa.Column("user_id", postgresql.UUID(as_uuid=True), sa.ForeignKey("users.user_id", ondelete="CASCADE"), nullable=False), + sa.Column("execution_kind", sa.Text(), nullable=False), + sa.Column("image_digest", sa.Text(), nullable=False), + sa.Column("python_version", sa.Text(), nullable=False), + sa.Column("packages", postgresql.JSONB(), nullable=False), + sa.Column("package_count", sa.Integer(), nullable=False), + sa.Column("total_installed_bytes", sa.BigInteger(), nullable=False), + sa.Column("finished_at", sa.DateTime(timezone=True), nullable=False), + sa.Column("created_at", sa.DateTime(timezone=True), server_default=sa.func.now(), nullable=False), + ) + op.create_index("ix_sandbox_package_scans_finished_at", "sandbox_package_scans", ["finished_at"]) + op.create_index("ix_sandbox_package_scans_user_finished", "sandbox_package_scans", ["user_id", "finished_at"]) + + +def downgrade() -> None: + op.drop_table("sandbox_package_scans") diff --git a/deploy/sandbox/Dockerfile b/deploy/sandbox/Dockerfile index f2e6df7..e92cc36 100644 --- a/deploy/sandbox/Dockerfile +++ b/deploy/sandbox/Dockerfile @@ -64,6 +64,9 @@ RUN grep -v '\[host-only\]' /tmp/requirements.txt > /tmp/requirements.sandbox.tx -r /tmp/requirements.sandbox.txt \ && rm /tmp/requirements.txt /tmp/requirements.sandbox.txt +# 构建期基础包清单;运行时扫描只与它比较,不 import/执行用户包。 +RUN mkdir -p /sandbox && python -c "import importlib.metadata as m,json; c=lambda s:s.lower().replace('_','-').replace('.','-'); print(json.dumps({c(d.metadata['Name']):d.version for d in m.distributions() if d.metadata.get('Name')},sort_keys=True))" > /sandbox/base-python-packages.json + # 持久化 pip 源到 /etc/pip.conf ── 让运行时模型用 `pip install foo` 也走 mirror, # 不只 build 时。zcbot user / root 都吃这个 global 配置。 RUN printf '[global]\nindex-url = %s\ntimeout = 60\n%s\n' \ @@ -176,6 +179,7 @@ COPY core/__init__.py core/file_store.py core/kb_lock.py /sandbox/core/ COPY core/sandbox/tool_runner.py /sandbox/tool_runner.py COPY deploy/sandbox/init.sh /init.sh +COPY deploy/sandbox/package_scan.py /sandbox/package_scan.py RUN chmod +x /init.sh # 默认 cwd /workspace ── 但每次 `docker exec --workdir /workspace/` 会覆盖 diff --git a/deploy/sandbox/package_scan.py b/deploy/sandbox/package_scan.py new file mode 100644 index 0000000..0b712f7 --- /dev/null +++ b/deploy/sandbox/package_scan.py @@ -0,0 +1,84 @@ +"""可信临时包扫描器;只解析 dist-info 元数据,不 import 用户代码。""" +from __future__ import annotations +import csv +import email.parser +import json +from pathlib import Path + +BASE = Path("/sandbox/base-python-packages.json") +ROOT = Path("/tmp/.local/lib") + + +def safe_file(root: Path, path: Path) -> bool: + try: + path.relative_to(root) + except ValueError: + return False + current = root + for part in path.relative_to(root).parts: + current = current / part + if current.is_symlink(): + return False + return path.is_file() + + +def safe_dir(root: Path, path: Path) -> bool: + try: + rel = path.relative_to(root) + except ValueError: + return False + current = root + for part in rel.parts: + current = current / part + if current.is_symlink(): + return False + return path.is_dir() + + +def canonical(name: str) -> str: + return name.lower().replace("_", "-").replace(".", "-") + + +def scan_packages(root: Path, base: dict) -> list[dict]: + out = [] + if root.is_dir(): + for site in sorted(root.glob("python*/site-packages")): + if not safe_dir(root, site): + continue + for dist in sorted(site.glob("*.dist-info")): + if len(out) >= 5000: + return out + meta_path = dist / "METADATA" + if not safe_file(site, meta_path) or meta_path.stat().st_size > 1024 * 1024: + continue + msg = email.parser.Parser().parsestr(meta_path.read_text(encoding="utf-8", errors="replace"), headersonly=True) + name, version = msg.get("Name", "").strip(), msg.get("Version", "").strip() + if not name: + continue + total = 0 + record = dist / "RECORD" + if safe_file(site, record) and record.stat().st_size <= 32 * 1024 * 1024: + for index, row in enumerate(csv.reader(record.read_text(encoding="utf-8", errors="replace").splitlines())): + if index >= 200_000: + break + if not row: + continue + candidate = site / row[0] + if safe_file(site, candidate): + try: total += candidate.stat().st_size + except OSError: pass + key = canonical(name) + old = base.get(key) + change = "added" if old is None else ("reinstall" if old == version else "override") + out.append({"name": name, "version": version, "base_version": old, + "change": change, "direct_requested": safe_file(site, dist / "REQUESTED") or safe_file(site, dist / "direct_url.json"), + "installed_bytes": total}) + return out + + +def main() -> None: + base = json.loads(BASE.read_text(encoding="utf-8")) + print(json.dumps(scan_packages(ROOT, base), ensure_ascii=False)) + +if __name__ == "__main__": + main() diff --git a/tests/test_executor_docker.py b/tests/test_executor_docker.py index 0612192..f0f61cb 100644 --- a/tests/test_executor_docker.py +++ b/tests/test_executor_docker.py @@ -18,6 +18,7 @@ import tempfile import threading import time import unittest +from contextlib import contextmanager from pathlib import Path from unittest.mock import MagicMock, patch from uuid import uuid4 @@ -35,6 +36,18 @@ class FakePool: def __init__(self): self.ensure_calls = [] self.mark_active_calls = [] + self.active = 0 + self.capacity = self + + @contextmanager + def foreground(self, user_id, cancel_check=None): + yield True + + def exec_started(self, user_id): + self.active += 1 + + def exec_finished(self, user_id): + self.active -= 1 def ensure(self, user_id): name = f"zcbot-sandbox-{user_id}" diff --git a/tests/test_sandbox_capacity_packages.py b/tests/test_sandbox_capacity_packages.py new file mode 100644 index 0000000..cc65fec --- /dev/null +++ b/tests/test_sandbox_capacity_packages.py @@ -0,0 +1,183 @@ +from __future__ import annotations +import json +import os +import tempfile +import threading +import time +import unittest +from datetime import datetime, timezone +from types import SimpleNamespace +from pathlib import Path +from unittest.mock import MagicMock, patch +from uuid import uuid4 + +from core.sandbox.capacity import ExecCapacity +from core.sandbox.pool import SandboxPool +from core.sandbox.package_scans import persist_scan +from deploy.sandbox.package_scan import scan_packages +from core import procs +from web.admin import _sandbox_package_stats + + +class CapacityTests(unittest.TestCase): + def test_env_overrides_capacity_config(self): + with tempfile.TemporaryDirectory() as td, patch.dict(os.environ, { + "ZCBOT_MAX_ACTIVE_EXECS": "5", "ZCBOT_MAX_BACKGROUND_EXECS": "3", + "ZCBOT_MAX_ACTIVE_EXECS_PER_USER": "1", "ZCBOT_MIN_MEM_AVAILABLE": "2g", + }, clear=False): + cap = ExecCapacity(Path(td), {"max_active_execs": 9, "max_background_execs": 8, + "max_active_execs_per_user": 7, "min_mem_available": "4g"}) + self.assertEqual((cap.max_active, cap.max_background, cap.max_per_user, cap.min_mem_available), + (5, 3, 1, 2 * 1024**3)) + + def test_config_cannot_raise_hard_limits(self): + with tempfile.TemporaryDirectory() as td: + cap = ExecCapacity(Path(td), {"max_active_execs": 99, "max_background_execs": 99, + "max_active_execs_per_user": 99}) + self.assertEqual((cap.max_active, cap.max_background, cap.max_per_user), (6, 4, 2)) + + def test_global_per_user_and_background_limits(self): + with tempfile.TemporaryDirectory() as td, patch("core.sandbox.capacity.mem_available_bytes", return_value=10**12): + cap = ExecCapacity(Path(td), {"max_active_execs": 3, "max_background_execs": 1, "max_active_execs_per_user": 2}) + a = cap.try_acquire("u1", "background", lease_id="a") + b = cap.try_acquire("u1", "foreground", lease_id="b") + self.assertIsNotNone(a); self.assertIsNotNone(b) + self.assertIsNone(cap.try_acquire("u1", "foreground", lease_id="c")) + self.assertIsNone(cap.try_acquire("u2", "background", lease_id="d")) + self.assertIsNotNone(cap.try_acquire("u2", "foreground", lease_id="e")) + self.assertIsNone(cap.try_acquire("u3", "foreground", lease_id="f")) + + def test_cross_instance_state_and_cancel_queue(self): + with tempfile.TemporaryDirectory() as td, patch("core.sandbox.capacity.mem_available_bytes", return_value=10**12): + one = ExecCapacity(Path(td), {"max_active_execs": 1}) + two = ExecCapacity(Path(td), {"max_active_execs": 1}) + lease = one.try_acquire("u1", "foreground", lease_id="held") + cancelled = threading.Event() + result = [] + t = threading.Thread(target=lambda: result.append(two.acquire_foreground("u2", cancelled.is_set))) + t.start(); time.sleep(.15); cancelled.set(); t.join(2) + self.assertEqual(result, [None]) + self.assertEqual(two.snapshot()["foreground_queued"], 0) + one.release(lease) + + def test_memory_pressure_pauses_new_admission(self): + with tempfile.TemporaryDirectory() as td, patch("core.sandbox.capacity.mem_available_bytes", return_value=100): + cap = ExecCapacity(Path(td), {"min_mem_available": "1g"}) + self.assertIsNone(cap.try_acquire("u", "foreground")) + self.assertTrue(cap.snapshot()["memory_paused"]) + + +class PoolTests(unittest.TestCase): + def test_defaults_and_docker_tmpfs(self): + with tempfile.TemporaryDirectory() as td, patch("core.sandbox.pool.subprocess.run") as run: + run.return_value = MagicMock(returncode=0, stdout="") + pool = SandboxPool(Path(td) / "users") + self.assertEqual((pool.memory, pool.cpus, pool.idle_ttl, pool.tmp_size), ("4g", "2.0", 600, "1g")) + pool._docker_run(uuid4(), "box") + argv = run.call_args_list[-1].args[0] + self.assertIn("/tmp:exec,size=1g,mode=1777", argv) + + def test_active_exec_not_reaped(self): + with tempfile.TemporaryDirectory() as td, patch("core.sandbox.pool._container_running", return_value=True), patch("core.sandbox.pool.subprocess.run") as run: + pool = SandboxPool(Path(td) / "users", idle_ttl=1) + uid = uuid4(); pool._last_active[uid] = int(time.time()) - 10 + pool.exec_started(uid) + self.assertEqual(pool.reap_idle(), []) + run.assert_not_called() + + def test_scan_failure_does_not_block_idle_removal(self): + with tempfile.TemporaryDirectory() as td, \ + patch("core.sandbox.package_scans.scan_and_persist", side_effect=RuntimeError("db down")), \ + patch("core.sandbox.pool.subprocess.run", return_value=MagicMock(returncode=0)): + pool = SandboxPool(Path(td) / "users", idle_ttl=1) + uid = uuid4(); pool._last_active[uid] = int(time.time()) - 10 + self.assertEqual(pool.reap_idle(), [f"zcbot-sandbox-{uid}"]) + + def test_finalizer_helper_swallows_scanner_error(self): + with tempfile.TemporaryDirectory() as td, patch("core.sandbox.package_scans.scan_and_persist", side_effect=RuntimeError("db down")): + pool = SandboxPool(Path(td) / "users") + pool._scan_before_remove("box", uuid4(), "foreground") + + +class PackageScannerTests(unittest.TestCase): + def test_empty_scan_does_not_touch_database(self): + persist_scan(container_session_id="empty", user_id=uuid4(), execution_kind="foreground", + result={"packages": [], "image_digest": "sha256:x", "python_version": "3.12", "total_installed_bytes": 0}) + + def test_added_override_bytes_and_symlink_escape(self): + with tempfile.TemporaryDirectory() as td: + root = Path(td) / "lib"; site = root / "python3.12" / "site-packages"; site.mkdir(parents=True) + pkg = site / "demo.py"; pkg.write_text("abc", encoding="utf-8") + dist = site / "Demo-2.dist-info"; dist.mkdir() + (dist / "METADATA").write_text("Name: Demo\nVersion: 2\n", encoding="utf-8") + (dist / "REQUESTED").write_text("", encoding="utf-8") + (dist / "RECORD").write_text("demo.py,,3\n", encoding="utf-8") + rows = scan_packages(root, {"demo": "1"}) + self.assertEqual(rows[0]["change"], "override") + self.assertEqual(rows[0]["installed_bytes"], 3) + self.assertTrue(rows[0]["direct_requested"]) + self.assertEqual(scan_packages(root, {"demo": "2"})[0]["change"], "reinstall") + self.assertEqual(scan_packages(root, {})[0]["change"], "added") + + def test_metadata_symlink_is_ignored(self): + with tempfile.TemporaryDirectory() as td: + root = Path(td) / "lib"; site = root / "python3.12" / "site-packages"; site.mkdir(parents=True) + outside = Path(td) / "outside"; outside.write_text("Name: Bad\nVersion: 1\n") + dist = site / "Bad.dist-info"; dist.mkdir() + try: + (dist / "METADATA").symlink_to(outside) + except OSError: + self.skipTest("symlink unavailable") + self.assertEqual(scan_packages(root, {}), []) + + +class ProcQueueTests(unittest.TestCase): + def test_queued_status_and_cancel(self): + with tempfile.TemporaryDirectory() as td: + d = Path(td) / "p"; d.mkdir() + meta = {"proc_id": "p1", "backend": "docker", "state": "queued"} + procs.write_meta(d, meta) + self.assertEqual(procs.status_of(meta, d)[0], "queued") + self.assertEqual(procs.kill_proc(meta, d), "已取消排队") + self.assertEqual(procs.status_of(procs.read_meta(d), d), ("finished", 137)) + + def test_start_queued_respects_capacity_then_runs(self): + with tempfile.TemporaryDirectory() as td: + d = Path(td) / "p"; d.mkdir(); (d / "runner.sh").write_text("true") + uid = uuid4() + meta = {"proc_id": "p1", "task_id": "t1", "backend": "docker", "state": "queued", + "user_id": str(uid), "cwd": "/workspace/demo"} + procs.write_meta(d, meta) + pool = MagicMock(); pool.capacity.try_acquire.return_value = None + self.assertFalse(procs.start_queued_docker(meta, d, pool)) + pool.capacity.try_acquire.return_value = "bg-p1" + pool.run_proc_container.return_value = "zcbot-proc-p1" + with patch("core.procs.subprocess.run", return_value=MagicMock(returncode=0, stderr="")): + self.assertTrue(procs.start_queued_docker(meta, d, pool)) + saved = procs.read_meta(d) + self.assertEqual(saved["state"], "running") + self.assertEqual(saved["capacity_lease_id"], "bg-p1") + + +class AdminAggregationTests(unittest.TestCase): + def test_package_sessions_are_aggregated_by_name_version(self): + scans = [ + SimpleNamespace(container_session_id="s1", user_id=uuid4(), execution_kind="foreground", + packages=[{"name":"Demo","version":"1","direct_requested":True,"installed_bytes":100,"change":"added","base_version":None}], + finished_at=datetime.now(timezone.utc)), + SimpleNamespace(container_session_id="s2", user_id=uuid4(), execution_kind="background", + packages=[{"name":"Demo","version":"1","direct_requested":False,"installed_bytes":300,"change":"added","base_version":None}], + finished_at=datetime.now(timezone.utc)), + ] + session = MagicMock(); session.execute.return_value.scalars.return_value.all.return_value = scans + row = _sandbox_package_stats(session, None)["rows"][0] + self.assertEqual(row["session_count"], 2) + self.assertEqual(row["user_count"], 2) + self.assertEqual(row["foreground_sessions"], 1) + self.assertEqual(row["background_sessions"], 1) + self.assertEqual(row["direct_sessions"], 1) + self.assertEqual(row["average_installed_bytes"], 200) + + +if __name__ == "__main__": + unittest.main() diff --git a/tests/test_web_routes_nodb.py b/tests/test_web_routes_nodb.py index da34aa9..2a4ad85 100644 --- a/tests/test_web_routes_nodb.py +++ b/tests/test_web_routes_nodb.py @@ -111,6 +111,8 @@ class AuthGateTests(unittest.TestCase): ("POST", "/v1/tasks"), ("POST", "/v1/asr/transcribe"), ("GET", "/v1/admin/overview"), + ("GET", "/v1/admin/sandbox/capacity"), + ("GET", "/v1/admin/sandbox/packages"), ("GET", "/v1/admin/software-nodes"), ("GET", "/v1/software-jobs"), ("POST", "/v1/software-jobs/00000000-0000-0000-0000-000000000000/cancel"), diff --git a/tools/check_process.py b/tools/check_process.py index 86188ef..2f29f9b 100644 --- a/tools/check_process.py +++ b/tools/check_process.py @@ -90,18 +90,28 @@ class CheckProcessTool(Tool): tail = procs.tail_log(d, max_bytes=tail_bytes) # docker 模式:已结束但专用容器还在空转 → 顺手回收(幂等,sweep 也会兜底) - if st != "running" and meta.get("backend") == "docker": + if st not in {"running", "queued"} and meta.get("backend") == "docker": name = str(meta.get("container") or "") if name: try: + procs._scan_proc_container(meta, d, name) subprocess.run(["docker", "rm", "-f", name], capture_output=True, timeout=30) except (OSError, subprocess.TimeoutExpired): pass + try: + from core.sandbox import get_pool + pool = get_pool() + if pool is not None: + pool.capacity.release(meta.get("capacity_lease_id")) + except Exception: + pass header = f"[check_process] {proc_id} · {meta.get('kind')} · {meta.get('command', '')[:150]}" - if st == "running": - elapsed = procs._fmt_elapsed(time.time() - float(meta.get("created_ts") or time.time())) + if st == "queued": + body = "状态: queued(等待 Sandbox 执行容量;排队时间不计运行超时)\n提示:可用 action='kill' 取消排队。" + elif st == "running": + elapsed = procs._fmt_elapsed(time.time() - float(meta.get("started_ts") or meta.get("created_ts") or time.time())) body = ( f"状态: running(已运行 {elapsed},上限 {meta.get('timeout_s')}s)\n" f"--- 日志尾部 ---\n{tail}\n" diff --git a/web/admin.py b/web/admin.py index 3dca8d3..87ae03b 100644 --- a/web/admin.py +++ b/web/admin.py @@ -24,7 +24,7 @@ from sqlalchemy import func, select, update from core.storage import session_scope from core.storage import usage_report -from core.storage.models import Task, UsageEvent, User, UserDiskUsage +from core.storage.models import SandboxPackageScan, Task, UsageEvent, User, UserDiskUsage from .broker import broker @@ -51,6 +51,37 @@ def _range_cutoff(now: datetime, range_key: str): return None # all / 未知 → 不筛 +def _sandbox_package_stats(s: Any, cutoff: datetime | None) -> dict: + stmt = select(SandboxPackageScan) + if cutoff is not None: + stmt = stmt.where(SandboxPackageScan.finished_at >= cutoff) + scans = s.execute(stmt).scalars().all() + grouped: dict[tuple[str, str], dict] = {} + for scan in scans: + for package in scan.packages: + key = (str(package.get("name") or ""), str(package.get("version") or "")) + row = grouped.setdefault(key, {"name": key[0], "version": key[1], "sessions": set(), + "users": set(), "foreground_sessions": 0, "background_sessions": 0, + "direct_sessions": 0, "total_installed_bytes": 0, "latest_at": None, + "base_versions": set(), "changes": set()}) + row["sessions"].add(scan.container_session_id); row["users"].add(str(scan.user_id)) + row[f"{scan.execution_kind}_sessions"] += 1 + row["direct_sessions"] += int(bool(package.get("direct_requested"))) + row["total_installed_bytes"] += int(package.get("installed_bytes") or 0) + row["base_versions"].add(package.get("base_version")); row["changes"].add(package.get("change")) + if row["latest_at"] is None or scan.finished_at > row["latest_at"]: + row["latest_at"] = scan.finished_at + rows = [] + for row in grouped.values(): + sessions = len(row.pop("sessions")); row["session_count"] = sessions + row["user_count"] = len(row.pop("users")); row["average_installed_bytes"] = row["total_installed_bytes"] // max(1, sessions) + row["base_versions"] = sorted(str(x) for x in row["base_versions"] if x is not None) + row["changes"] = sorted(row["changes"]); row["latest_at"] = row["latest_at"].isoformat() + rows.append(row) + rows.sort(key=lambda x: (-x["session_count"], x["name"].lower(), x["version"])) + return {"rows": rows, "scan_sessions": len(scans)} + + def _runtime_section(app: FastAPI) -> dict: """实时运行态:从 app.state 读内存,无 DB。 @@ -281,6 +312,22 @@ def register_admin_routes(app: FastAPI, require_admin) -> None: "usage": usage_report.usage_overview(s, cutoff_7d), } + @app.get("/v1/admin/sandbox/capacity", tags=["admin"]) + def admin_sandbox_capacity(user_id: UUID = Depends(require_admin)): + """宿主共享实时容量;只读文件/Docker 状态,不写 DB。""" + pool = getattr(app.state, "sandbox_pool", None) + if pool is None: + return {"enabled": False} + return {"enabled": True, **pool.runtime_snapshot()} + + @app.get("/v1/admin/sandbox/packages", tags=["admin"]) + def admin_sandbox_packages(range: str = "30d", user_id: UUID = Depends(require_admin)): + now = datetime.now(timezone.utc) + with session_scope() as s: + result = _sandbox_package_stats(s, _range_cutoff(now, range)) + result["range"] = range + return result + @app.get("/v1/admin/external-system-definitions", tags=["admin"]) def admin_external_system_definitions(user_id: UUID = Depends(require_admin)): from core.external_systems.service import list_external_system_definitions diff --git a/web/background.py b/web/background.py index 8db08ed..539a6db 100644 --- a/web/background.py +++ b/web/background.py @@ -248,7 +248,7 @@ def start_proc_sweeper(cfg: dict) -> asyncio.Task: 已结束容器/过期目录收掉。 """ from core.agent_builder import resolve_workspace - from core.procs import sweep as _procs_sweep + from core.procs import dispatch_queued, sweep as _procs_sweep users_base = resolve_workspace(None, cfg) / "users" async def _proc_sweeper() -> None: @@ -258,14 +258,19 @@ def start_proc_sweeper(cfg: dict) -> asyncio.Task: stats = await loop.run_in_executor( None, _procs_sweep, users_base ) + from core.sandbox import get_pool + pool = get_pool() + started = await loop.run_in_executor(None, dispatch_queued, users_base, pool) if pool is not None else 0 if stats["removed_dirs"] or stats["reaped_containers"]: print(f"[proc-sweep] dirs={stats['removed_dirs']} " f"containers={stats['reaped_containers']}") + if started: + print(f"[proc-sweep] started queued={started}") except asyncio.CancelledError: raise except Exception as e: print(f"[proc-sweep] error: {type(e).__name__}: {e}") - await asyncio.sleep(3600) + await asyncio.sleep(30) return asyncio.create_task(_proc_sweeper(), name="proc-sweeper") diff --git a/web/static/js/admin.js b/web/static/js/admin.js index c8acdd2..5188dcd 100644 --- a/web/static/js/admin.js +++ b/web/static/js/admin.js @@ -14,6 +14,7 @@ const SORT_OPTS = [["cost", "按成本"], ["tokens", "按用量"]]; const SECTIONS = [ ["s-overview", "总览"], ["s-usage", "用量趋势"], ["s-models", "按模型"], ["s-users", "各用户"], ["s-storage", "存储"], + ["s-sandbox-capacity", "Sandbox 容量"], ["s-sandbox-packages", "Sandbox 依赖"], ["s-windows-node", "Windows Node"], ["s-external", "外部系统"], ["s-toolfail", "工具失败"], @@ -50,6 +51,7 @@ let timer = null; let modelRange = "7d", modelSort = "cost"; let userRange = "7d", userSort = "cost", userPage = 0; let storagePage = 0; +let packageRange = "30d"; let tiersData = null; // {tiers, default_tier, catalog};加载一次(改档位 / 看图例用) let externalDefinitions = []; let externalUsers = []; @@ -1052,6 +1054,42 @@ function renderMetrics(d) { renderOpsSummary(); } +function renderSandboxCapacity(d) { + if (!d.enabled) { + $("s-sandbox-capacity").innerHTML = `

Sandbox 实时容量

Docker Sandbox 未启用。

`; + return; + } + const lim = d.limits || {}; + const users = Object.entries(d.per_user || {}).sort((a,b) => b[1]-a[1]) + .map(([u,n]) => `${escapeHtml(u.slice(0,8))}: ${n}`).join(" · ") || "无"; + $("s-sandbox-capacity").innerHTML = `

Sandbox 实时容量

` + + `
${d.foreground_running || 0}前台执行
` + + `
${d.foreground_queued || 0}前台排队
` + + `
${d.background_running || 0}后台执行
` + + `
${d.background_queued || 0}后台排队
` + + `
${d.admit_available || 0}当前可放行
` + + `
${d.active_sandbox_containers || 0}活跃普通容器
` + + `
${d.idle_reap_candidates || 0}空闲待回收
` + + `
${d.running_proc_containers || 0}运行中 proc
` + + `

硬上限 ${lim.active || 0} · 后台 ${lim.background || 0} · 单用户 ${lim.per_user || 0}` + + ` · MemAvailable ${humanSize(d.mem_available_bytes || 0)} · CPU load ${(d.cpu_load || []).join(" / ") || "—"}` + + `${d.memory_paused ? " · 内存压力暂停放行" : ""}

单用户占用:${users}

`; +} + +function renderSandboxPackages(d) { + const rows = d.rows || []; + const opts = RANGE_OPTS.map(([v,l]) => ``).join(""); + const body = rows.map(r => `${escapeHtml(r.name)}${escapeHtml(r.version)}` + + `${r.session_count}${r.user_count}${r.foreground_sessions}/${r.background_sessions}` + + `${r.direct_sessions}${humanSize(r.average_installed_bytes)}${humanSize(r.total_installed_bytes)}` + + `${escapeHtml((r.changes || []).join("/"))}
${fmtTimeAgo(r.latest_at)}`).join(""); + $("s-sandbox-packages").innerHTML = `

Sandbox 依赖统计

` + + `

只分析临时安装,不自动修改基础镜像;共 ${d.scan_sessions || 0} 个安装会话。

` + + `
${body || ``}
包名版本安装会话独立用户前台/后台直接安装会话平均占用累计占用基础镜像状态 / 最近
暂无临时依赖记录
`; + const select = $("sandbox-package-range"); + if (select) select.onchange = () => { packageRange = select.value; loadSandboxPackages(); }; +} + function showMsg(html) { $("main").innerHTML = `
${html}
`; // 清骨架,错误态独占 } @@ -1199,6 +1237,13 @@ async function loadSoftwareNodes() { } catch (e) { /* overview 统一处理鉴权 */ } } +async function loadSandboxCapacity() { + try { renderSandboxCapacity(await apiGet("/v1/admin/sandbox/capacity")); } catch (e) { /* overview 统一处理 */ } +} +async function loadSandboxPackages() { + try { renderSandboxPackages(await apiGet(`/v1/admin/sandbox/packages?range=${packageRange}`)); } catch (e) { /* overview 统一处理 */ } +} + // overview(固定指标)轮询:拿到后建骨架、渲指标,再顺手刷新四个独立表(保持各自状态) async function refresh() { try { @@ -1212,6 +1257,8 @@ async function refresh() { loadSoftwareNodes(); loadExternalDefinitions(); loadToolFailures(); + loadSandboxCapacity(); + loadSandboxPackages(); } catch (e) { if (e.code !== "auth") showMsg(`加载失败:${escapeHtml(e.message || String(e))}`); }