Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 1 addition & 2 deletions backend/package/yuxi/services/readiness_service.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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"))
Expand Down
7 changes: 1 addition & 6 deletions backend/package/yuxi/services/run_queue_service.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
3 changes: 1 addition & 2 deletions backend/package/yuxi/services/run_worker.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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
Expand Down
35 changes: 35 additions & 0 deletions backend/package/yuxi/services/worker_health.py
Original file line number Diff line number Diff line change
@@ -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())
69 changes: 69 additions & 0 deletions backend/test/integration/services/test_worker_health_redis.py
Original file line number Diff line number Diff line change
@@ -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
6 changes: 4 additions & 2 deletions backend/test/unit/config/test_docker_compose_checkpointer.py
Original file line number Diff line number Diff line change
Expand Up @@ -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",
]


Expand Down
62 changes: 62 additions & 0 deletions backend/test/unit/services/test_worker_health.py
Original file line number Diff line number Diff line change
@@ -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
2 changes: 1 addition & 1 deletion docker-compose.prod.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
2 changes: 1 addition & 1 deletion docker-compose.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
2 changes: 2 additions & 0 deletions docs/advanced/deployment.md
Original file line number Diff line number Diff line change
Expand Up @@ -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`,并在域名侧保持同样的只读限制。
Expand Down
Original file line number Diff line number Diff line change
@@ -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 租约。
1 change: 1 addition & 0 deletions web/vite.config.js
Original file line number Diff line number Diff line change
Expand Up @@ -102,6 +102,7 @@ export default defineConfig(({ mode }) => {
},
watch: {
usePolling: true,
interval: 1000,
ignored: ['**/node_modules/**', '**/dist/**']
},
host: '0.0.0.0'
Expand Down
Loading