From af6f0b3e0652cf791d3e470a623dc9d9e82de5bb Mon Sep 17 00:00:00 2001 From: Wenjie Zhang Date: Mon, 28 Sep 2026 15:33:16 +0800 Subject: [PATCH] =?UTF-8?q?fix:=20=E9=99=8D=E4=BD=8E=20worker=20=E5=81=A5?= =?UTF-8?q?=E5=BA=B7=E6=A3=80=E6=9F=A5=E4=B8=8E=E5=89=8D=E7=AB=AF=E8=BD=AE?= =?UTF-8?q?=E8=AF=A2=E7=9A=84=E7=A9=BA=E9=97=B2=E5=BC=80=E9=94=80?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../yuxi/services/readiness_service.py | 3 +- .../yuxi/services/run_queue_service.py | 7 +- backend/package/yuxi/services/run_worker.py | 3 +- .../package/yuxi/services/worker_health.py | 35 ++++++++++ .../services/test_worker_health_redis.py | 69 +++++++++++++++++++ .../test_docker_compose_checkpointer.py | 6 +- .../test/unit/services/test_worker_health.py | 62 +++++++++++++++++ docker-compose.prod.yml | 2 +- docker-compose.yml | 2 +- docs/advanced/deployment.md | 2 + .../implemented/2026-09-28-worker-idle-cpu.md | 37 ++++++++++ web/vite.config.js | 1 + 12 files changed, 215 insertions(+), 14 deletions(-) create mode 100644 backend/package/yuxi/services/worker_health.py create mode 100644 backend/test/integration/services/test_worker_health_redis.py create mode 100644 backend/test/unit/services/test_worker_health.py create mode 100644 docs/develop-guides/decisions/implemented/2026-09-28-worker-idle-cpu.md diff --git a/backend/package/yuxi/services/readiness_service.py b/backend/package/yuxi/services/readiness_service.py index 8d9a01365b..dbb8225b93 100644 --- a/backend/package/yuxi/services/readiness_service.py +++ b/backend/package/yuxi/services/readiness_service.py @@ -11,8 +11,6 @@ from sqlalchemy import text from yuxi.services.run_queue_service import ( - WORKER_HEALTH_KEY, - WORKER_HEALTH_MAX_TTL_MS, WORKER_RECONCILIATION_HEALTH_KEY, WORKER_RECONCILIATION_HEALTH_TTL_SECONDS, get_redis_client, @@ -21,6 +19,7 @@ TASK_RECONCILIATION_HEALTH_KEY, TASK_RECONCILIATION_HEALTH_TTL_SECONDS, ) +from yuxi.services.worker_health import WORKER_HEALTH_KEY, WORKER_HEALTH_MAX_TTL_MS from yuxi.storage.postgres.manager import pg_manager READINESS_PROBE_TIMEOUT_SECONDS = float(os.getenv("READINESS_PROBE_TIMEOUT_SECONDS", "2")) diff --git a/backend/package/yuxi/services/run_queue_service.py b/backend/package/yuxi/services/run_queue_service.py index 1bd7a034a8..8ad09aed9a 100644 --- a/backend/package/yuxi/services/run_queue_service.py +++ b/backend/package/yuxi/services/run_queue_service.py @@ -7,18 +7,13 @@ import os from datetime import UTC, datetime +from yuxi.services.worker_health import WORKER_HEALTH_KEY from yuxi.storage.redis import close_async_redis_client, create_arq_redis_pool, get_async_redis_client from yuxi.utils.logging_config import logger RUN_CANCEL_KEY_TTL_SECONDS = int(os.getenv("RUN_CANCEL_KEY_TTL_SECONDS", "1800")) RUN_EVENTS_STREAM_TTL_SECONDS = int(os.getenv("RUN_EVENTS_STREAM_TTL_SECONDS", "7200")) RUN_EVENTS_STREAM_MAXLEN = int(os.getenv("RUN_EVENTS_STREAM_MAXLEN", "0")) -WORKER_HEALTH_CONTRACT = "agent-run-v1" -WORKER_HEALTH_KEY = f"yuxi:worker:health:{WORKER_HEALTH_CONTRACT}" -WORKER_HEALTH_INTERVAL_SECONDS = float(os.getenv("WORKER_HEALTH_INTERVAL_SECONDS", "5")) -if not 0 < WORKER_HEALTH_INTERVAL_SECONDS <= 10: - raise ValueError("WORKER_HEALTH_INTERVAL_SECONDS 必须大于 0 且不超过 10") -WORKER_HEALTH_MAX_TTL_MS = int((WORKER_HEALTH_INTERVAL_SECONDS + 1) * 1000) RUN_RECONCILIATION_SECONDS = 30 WORKER_RECONCILIATION_HEALTH_KEY = f"{WORKER_HEALTH_KEY}:lease-reconciliation" WORKER_RECONCILIATION_HEALTH_TTL_SECONDS = RUN_RECONCILIATION_SECONDS * 2 + 5 diff --git a/backend/package/yuxi/services/run_worker.py b/backend/package/yuxi/services/run_worker.py index 7b66c65605..629dd3998c 100644 --- a/backend/package/yuxi/services/run_worker.py +++ b/backend/package/yuxi/services/run_worker.py @@ -33,8 +33,6 @@ from yuxi.services.input_message_service import restore_chat_input_message from yuxi.services.run_queue_service import ( RUN_RECONCILIATION_SECONDS, - WORKER_HEALTH_INTERVAL_SECONDS, - WORKER_HEALTH_KEY, WORKER_RECONCILIATION_HEALTH_KEY, WORKER_RECONCILIATION_HEALTH_TTL_SECONDS, append_run_stream_event, @@ -59,6 +57,7 @@ resolve_authorized_workdir, resolve_conversation_workdir_path, ) +from yuxi.services.worker_health import WORKER_HEALTH_INTERVAL_SECONDS, WORKER_HEALTH_KEY from yuxi.storage.postgres.manager import pg_manager from yuxi.storage.postgres.models_business import AgentRun, Conversation, Message, User from yuxi.storage.redis import get_arq_redis_settings diff --git a/backend/package/yuxi/services/worker_health.py b/backend/package/yuxi/services/worker_health.py new file mode 100644 index 0000000000..db27109045 --- /dev/null +++ b/backend/package/yuxi/services/worker_health.py @@ -0,0 +1,35 @@ +"""ARQ 消费心跳契约与轻量 Compose 健康检查。""" + +import os +import sys + +from yuxi.storage.redis import RedisConfig, sync_redis_client + +WORKER_HEALTH_CONTRACT = "agent-run-v1" +WORKER_HEALTH_KEY = f"yuxi:worker:health:{WORKER_HEALTH_CONTRACT}" +WORKER_HEALTH_INTERVAL_SECONDS = float(os.getenv("WORKER_HEALTH_INTERVAL_SECONDS", "5")) +if not 0 < WORKER_HEALTH_INTERVAL_SECONDS <= 10: + raise ValueError("WORKER_HEALTH_INTERVAL_SECONDS 必须大于 0 且不超过 10") +WORKER_HEALTH_MAX_TTL_MS = int((WORKER_HEALTH_INTERVAL_SECONDS + 1) * 1000) + + +def main() -> int: + """读取有界心跳租约,失败时仅输出错误类型以避免泄露连接凭据。""" + try: + config = RedisConfig.from_env(socket_timeout=2, socket_connect_timeout=2) + with sync_redis_client(config, ping=False) as client: + with client.pipeline() as pipeline: + pipeline.get(WORKER_HEALTH_KEY) + pipeline.pttl(WORKER_HEALTH_KEY) + value, ttl_ms = pipeline.execute() + if not value or not 0 < ttl_ms <= WORKER_HEALTH_MAX_TTL_MS: + print("worker health lease missing or invalid", file=sys.stderr) + return 1 + except Exception as exc: + print(f"worker health check failed: {type(exc).__name__}", file=sys.stderr) + return 1 + return 0 + + +if __name__ == "__main__": + sys.exit(main()) diff --git a/backend/test/integration/services/test_worker_health_redis.py b/backend/test/integration/services/test_worker_health_redis.py new file mode 100644 index 0000000000..211471f59c --- /dev/null +++ b/backend/test/integration/services/test_worker_health_redis.py @@ -0,0 +1,69 @@ +"""使用真实 ARQ 心跳和 Redis 验证轻量探针。""" + +import asyncio +import uuid + +import pytest +import pytest_asyncio +from arq import create_pool +from arq.worker import Worker +from yuxi.services import worker_health +from yuxi.storage.redis import get_arq_redis_settings + +pytestmark = [pytest.mark.asyncio, pytest.mark.integration] + + +async def unused_job(ctx): + """满足 ARQ 注册约束,健康检查测试不投递任务。""" + + +@pytest_asyncio.fixture +async def health_redis(monkeypatch): + """健康键和队列仅属于本测试,清理不影响共享 worker。""" + redis = await create_pool(get_arq_redis_settings()) + key = f"pytest-worker-health:{uuid.uuid4().hex}" + monkeypatch.setattr(worker_health, "WORKER_HEALTH_KEY", key) + try: + yield redis, key + finally: + await redis.delete(key) + await redis.aclose() + + +async def test_arq_heartbeat_expires_without_renewal(health_redis): + """ARQ 原生心跳可通过探针,停止续租后由 Redis 过期事实拒绝。""" + redis, key = health_redis + worker = Worker( + functions=[unused_job], + redis_pool=redis, + queue_name=f"{key}:queue", + health_check_key=key, + health_check_interval=0.1, + handle_signals=False, + ) + await worker.record_health() + assert await redis.get(key) + assert 0 < await redis.pttl(key) <= 1100 + assert worker_health.main() == 0 + await asyncio.sleep(1.2) + assert await redis.get(key) is None + assert worker_health.main() == 1 + + +@pytest.mark.parametrize("state", ["missing", "empty", "persistent", "excessive"]) +async def test_invalid_redis_leases_fail(health_redis, state): + """Redis 中的非法心跳不能维持健康状态。""" + redis, key = health_redis + if state == "empty": + await redis.set(key, b"", px=1000) + elif state == "persistent": + await redis.set(key, b"alive") + elif state == "excessive": + await redis.set(key, b"alive", px=worker_health.WORKER_HEALTH_MAX_TTL_MS + 60000) + assert worker_health.main() == 1 + + +async def test_unreachable_redis_fails(monkeypatch): + """连接拒绝产生非零结果。""" + monkeypatch.setenv("REDIS_URL", "redis://127.0.0.1:1/0") + assert worker_health.main() == 1 diff --git a/backend/test/unit/config/test_docker_compose_checkpointer.py b/backend/test/unit/config/test_docker_compose_checkpointer.py index 939699abeb..937816b00c 100644 --- a/backend/test/unit/config/test_docker_compose_checkpointer.py +++ b/backend/test/unit/config/test_docker_compose_checkpointer.py @@ -130,8 +130,10 @@ def test_worker_healthcheck_uses_arq_health_contract_in_development_and_producti compose = yaml.safe_load((project_root / filename).read_text()) assert compose["services"]["worker"]["healthcheck"]["test"] == [ - "CMD-SHELL", - "uv run --no-sync --no-dev arq --check server.worker_main.WorkerSettings", + "CMD", + "python", + "-m", + "yuxi.services.worker_health", ] diff --git a/backend/test/unit/services/test_worker_health.py b/backend/test/unit/services/test_worker_health.py new file mode 100644 index 0000000000..2680b9d1ed --- /dev/null +++ b/backend/test/unit/services/test_worker_health.py @@ -0,0 +1,62 @@ +"""轻量健康探针的失败边界和导入隔离。""" + +import os +import subprocess +import sys +from unittest.mock import MagicMock + +import pytest + + +@pytest.mark.parametrize( + "value,ttl,expected", + [ + (b"alive", 1000, 0), + (b"alive", 6000, 0), + (None, -2, 1), + (b"", 1000, 1), + (b"alive", -1, 1), + (b"alive", 0, 1), + (b"alive", 6001, 1), + ], +) +def test_health_requires_live_bounded_lease(monkeypatch, value, ttl, expected): + """缺失、空值、永久或超长租约不能被视为健康。""" + from yuxi.services import worker_health + + client = MagicMock() + client.pipeline.return_value.__enter__.return_value.execute.return_value = (value, ttl) + context = MagicMock() + context.__enter__.return_value = client + monkeypatch.setattr(worker_health, "sync_redis_client", lambda *args, **kwargs: context) + monkeypatch.setattr(worker_health, "WORKER_HEALTH_MAX_TTL_MS", 6000) + assert worker_health.main() == expected + + +def test_health_connection_error_does_not_expose_credentials(monkeypatch, capsys): + """连接失败返回非零且不输出异常中的凭据。""" + from yuxi.services import worker_health + + context = MagicMock() + context.__enter__.side_effect = ConnectionError("redis://user:secret@host/0") + monkeypatch.setattr(worker_health, "sync_redis_client", lambda *args, **kwargs: context) + assert worker_health.main() == 1 + assert "secret" not in capsys.readouterr().err + + +def test_health_import_does_not_load_business_runtime(): + """干净解释器在禁止业务运行时导入时仍可加载探针。""" + script = """ +import importlib.abc +import sys +class BlockBusiness(importlib.abc.MetaPathFinder): + def find_spec(self, fullname, path=None, target=None): + blocked = ('yuxi.services.run_worker', 'yuxi.services.run_queue_service', + 'langgraph', 'sqlalchemy', 'tiktoken') + if fullname.startswith(blocked): + raise AssertionError('health probe imported business runtime: ' + fullname) +sys.meta_path.insert(0, BlockBusiness()) +import yuxi.services.worker_health +""" + result = subprocess.run([sys.executable, "-c", script], capture_output=True, text=True, env=os.environ.copy()) + assert result.returncode == 0, result.stderr diff --git a/docker-compose.prod.yml b/docker-compose.prod.yml index e513d1ec99..8d7ca25437 100644 --- a/docker-compose.prod.yml +++ b/docker-compose.prod.yml @@ -126,7 +126,7 @@ services: command: uv run --no-sync --no-dev python -m server.worker_main restart: unless-stopped healthcheck: - test: ["CMD-SHELL", "uv run --no-sync --no-dev arq --check server.worker_main.WorkerSettings"] + test: ["CMD", "python", "-m", "yuxi.services.worker_health"] interval: 10s timeout: 10s retries: 6 diff --git a/docker-compose.yml b/docker-compose.yml index 31658e3e09..f6ee9fbc92 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -152,7 +152,7 @@ services: command: watchfiles --filter python "python -m server.worker_main" /app/server /app/package restart: unless-stopped healthcheck: - test: ["CMD-SHELL", "uv run --no-sync --no-dev arq --check server.worker_main.WorkerSettings"] + test: ["CMD", "python", "-m", "yuxi.services.worker_health"] interval: 10s timeout: 10s retries: 6 diff --git a/docs/advanced/deployment.md b/docs/advanced/deployment.md index 98cf11f211..4baf8ee554 100644 --- a/docs/advanced/deployment.md +++ b/docs/advanced/deployment.md @@ -158,6 +158,8 @@ curl --fail http://localhost/api/system/ready - `/api/system/health` 只表示 API 进程存活; - `/api/system/ready` 表示启动完成、PostgreSQL/Redis 可用,并且兼容 worker 正在提供健康租约。 +worker 的 Compose 健康检查通过 `python -m yuxi.services.worker_health` 轻量读取 `REDIS_URL` 中的 ARQ 心跳,不加载业务执行依赖。心跳缺失、过期、没有 TTL、TTL 超过约定上界或 Redis 连接失败时检查失败。该心跳表达共享队列的消费健康,多副本部署不能用它判断单个 worker 进程是否失活。 + 就绪接口返回 `ready` 后,再用浏览器完成登录和一次真实对话。健康或就绪状态不能证明知识库、模型、沙盒或外部服务的业务链路正确。 公开头像和智能体图片通过同源 `/minio/public/...` 只读代理访问。不要把 MinIO 的 9000 对象 API 或 9001 控制台暴露到公网;知识库等私有 bucket 不经过该代理。需要单独的静态资源域名时,设置 `MINIO_PUBLIC_URL`,并在域名侧保持同样的只读限制。 diff --git a/docs/develop-guides/decisions/implemented/2026-09-28-worker-idle-cpu.md b/docs/develop-guides/decisions/implemented/2026-09-28-worker-idle-cpu.md new file mode 100644 index 0000000000..b8628dabef --- /dev/null +++ b/docs/develop-guides/decisions/implemented/2026-09-28-worker-idle-cpu.md @@ -0,0 +1,37 @@ +# 降低空闲 worker 探针与开发文件监听开销 + +状态:implemented +类型:bug-fix +Owner:backend/package/yuxi/services/worker_health.py + +## 问题 + +Issue #1078 中周期性 ARQ 健康检查导入完整 worker 业务依赖,空闲时仍消耗 CPU 并超时;Vite 高频轮询进一步增加开发环境负载。目标是减少这两条周期路径的开销,保留失活检测与跨平台文件监听。其他进程峰值和业务执行性能不在范围内。 + +## 决策 + +### 实现方案 + +将 ARQ 心跳键、间隔和 TTL 上界放入轻量 worker_health 模块,生产者和 readiness 消费同一契约。开发与生产 Compose 直接运行该模块,复用 Redis 连接配置读取心跳值及 PTTL;缺失、空值、无过期时间、超长 TTL 或连接失败返回非零。保持现有检查频率、超时与重试。Vite 保留 polling 并设置 1000 ms 间隔。 + +## 替代方案 + +延长检查周期仍保留重依赖导入;关闭探针失去失活检测;关闭 polling 会影响部分 Windows/Docker 挂载的热更新,因此均不采用。 + +## 验证 + +| 验收主张 | 失败面 | 语义 Owner | 直接证据 / 命令 | 负向案例 | 当前结果 | +|---|---|---|---|---|---| +| 探针不加载 worker 业务依赖 | 每次检查重复重导入 | worker_health、Compose | 导入隔离 unit、命令检查 | 阻止导入 run_worker、LangGraph、SQLAlchemy | Passed | +| 心跳仍能表达失活 | 永久或过期键误判健康 | worker_health | unit、真实 Redis / ARQ integration | 缺失、过期、永久、超长 TTL、不可达 | Passed | +| 文件变更仍被监听 | polling 配置失效 | vite.config.js | Vite watcher 实测:1000 ms 轮询,813 ms 后发出 HMR update;web lint、380 unit、build | 文件修改后无通知即失败 | Passed | + +## 后果 + +轮询通知延迟增加到约一秒。心跳仍是共享队列级事实,多副本中不能辨认单个进程的失活;维持现有语义。Linux 测试不能替代 Issue 用户机器上的 CPU 和 Windows 热更新复测。 + +验证使用独占临时 Redis 与 API 镜像,工作树挂载为 `/repo`,`PYTHONPATH=/repo/backend/package`;后端完整 unit 2428 项通过。容器中的 `YUXI_SKILL_PROJECTION_DIR` 指向可写临时目录。针对 Redis 的测试不需要 API 或沙盒,执行 `python -m pytest -p no:cacheprovider --confcutdir=test/integration/services test/integration/services/test_worker_health_redis.py -q`,6 项通过;`--confcutdir` 排除全局 API/沙盒清理 fixture,保留测试自身的随机 Redis 键清理。 + +导入隔离和 Compose 装配由 `test/unit/services/test_worker_health.py` 与 `test/unit/config/test_docker_compose_checkpointer.py` 覆盖,readiness 回归仍通过。工程契约检查及其 62 项测试、Ruff、文档构建与 `git diff --check` 通过。前端和文档使用 `npm run` 执行 package.json 中相同脚本,避免 pnpm 重装跨工作树复用的依赖目录。 + +Windows Docker 热更新、浏览器实际消费 HMR、用户机器的整体空闲 CPU 及 402% 峰值来源为 Not run。探针检查的是共享 ARQ 心跳;API readiness 仍单独检查启动、数据库和 reconciliation 租约。 diff --git a/web/vite.config.js b/web/vite.config.js index 920382b372..f0ced6a3ec 100644 --- a/web/vite.config.js +++ b/web/vite.config.js @@ -102,6 +102,7 @@ export default defineConfig(({ mode }) => { }, watch: { usePolling: true, + interval: 1000, ignored: ['**/node_modules/**', '**/dist/**'] }, host: '0.0.0.0'