From 1b253510e1abf66390e09ad376975de08c14996a Mon Sep 17 00:00:00 2001 From: jxxghp Date: Fri, 21 Aug 2026 22:29:19 +0800 Subject: [PATCH] ci: block new async blocking calls --- .github/workflows/test.yml | 3 + app/api/endpoints/agent.py | 7 +- .../backend-architecture-next-stage.md | 11 + pytest.ini | 1 + scripts/architecture/async_blocking.py | 224 ++++++++++++++++++ .../architecture/async-blocking-baseline.json | 3 + tests/test_async_blocking_gate.py | 31 +++ 7 files changed, 277 insertions(+), 3 deletions(-) create mode 100644 scripts/architecture/async_blocking.py create mode 100644 tests/fixtures/architecture/async-blocking-baseline.json create mode 100644 tests/test_async_blocking_gate.py diff --git a/.github/workflows/test.yml b/.github/workflows/test.yml index 37c9a3822..03a86f735 100644 --- a/.github/workflows/test.yml +++ b/.github/workflows/test.yml @@ -56,6 +56,9 @@ jobs: - name: Check complexity ratchet run: uv run --locked --no-sync python scripts/architecture/complexity.py + - name: Check async blocking ratchet + run: uv run --locked --no-sync python scripts/architecture/async_blocking.py + pytest: runs-on: ubuntu-latest name: Unit Tests (${{ matrix.shard }}) diff --git a/app/api/endpoints/agent.py b/app/api/endpoints/agent.py index bf6374ff9..dfee17022 100644 --- a/app/api/endpoints/agent.py +++ b/app/api/endpoints/agent.py @@ -13,6 +13,7 @@ from pathlib import Path from threading import Lock from typing import Any, AsyncIterator, Callable, Optional, Union +import aiofiles from fastapi import Depends, File, Form, HTTPException, Request, UploadFile, status from fastapi.concurrency import run_in_threadpool from fastapi.responses import FileResponse, StreamingResponse @@ -758,7 +759,7 @@ async def _save_web_agent_upload(upload_file: UploadFile, target_path: Path) -> """ size = 0 try: - with target_path.open("wb") as output: + async with aiofiles.open(target_path, "wb") as output: while True: chunk = await upload_file.read(WEB_AGENT_UPLOAD_CHUNK_SIZE) if not chunk: @@ -769,9 +770,9 @@ async def _save_web_agent_upload(upload_file: UploadFile, target_path: Path) -> status_code=status.HTTP_413_REQUEST_ENTITY_TOO_LARGE, detail="附件超过 32MB,无法发送给智能助手", ) - output.write(chunk) + await output.write(chunk) except Exception: - target_path.unlink(missing_ok=True) + await run_in_threadpool(target_path.unlink, missing_ok=True) raise finally: await upload_file.close() diff --git a/docs/refactor/backend-architecture-next-stage.md b/docs/refactor/backend-architecture-next-stage.md index 77a1bdbb5..31db556fc 100644 --- a/docs/refactor/backend-architecture-next-stage.md +++ b/docs/refactor/backend-architecture-next-stage.md @@ -833,6 +833,17 @@ OTel 初始化只能位于 Startup/Adapter;Domain/Application 只依赖 no-op- 4. 同步 Module 由 dispatcher 线程池兼容,不要求第三方插件立刻 async 化; 5. 只在测量证明有收益时改用 async 第三方 client。 +**实施记录(2026-08-21)**: + +- 新增 `scripts/architecture/async_blocking.py`,AST 扫描 `app/api`、`app/agent`、`app/application` 的 + async 函数,覆盖直接 `open`/Path 读写遍历、`requests`、`time.sleep`、同步 `subprocess` 和目录遍历。 +- scanner 通过局部类型流识别 `aiofiles` 与 `anyio.AsyncPath`,不会把正确异步 I/O 记成债务;baseline + 只允许调用减少/删除,新增或次数增长均在 CI architecture job 失败。 +- Web Agent 上传已从 `Path.open/write/unlink` 改为 `aiofiles` 写入和统一 `run_in_threadpool` 清理;当前仅保留 + ActivityLog 为保证 `O_EXCL` 原子创建使用的一处 `os.open` 精确债务,不泛化豁免整个文件或目录。 +- pytest 全局启用 `asyncio_debug`,专项测试验证实际 loop debug 状态;AST ratchet 与 46 个 Agent 流式回归 + 通过。同步第三方 Module 仍由 dispatcher 的 `app.runtime.execution.run_in_threadpool` 兼容。 + ## 6. 推荐执行队列 下表是默认的提交顺序,不表示所有任务必须由同一个 AI 连续完成。一个 AI 一次只领取一行;如果发现前置条件未满足,应停止实施并回报证据,不得顺手扩大范围。 diff --git a/pytest.ini b/pytest.ini index 3c60316d6..71ec18b26 100644 --- a/pytest.ini +++ b/pytest.ini @@ -5,6 +5,7 @@ timeout = 120 timeout_method = thread asyncio_mode = strict asyncio_default_fixture_loop_scope = function +asyncio_debug = true # 仅对「无法在本仓修复根因」的已知上游/三方弃用告警做精确忽略,保持测试输出干净、 # 让本仓自身的新告警更醒目。本仓代码引发的告警一律不在此忽略,应在源码/用例处修复。 filterwarnings = diff --git a/scripts/architecture/async_blocking.py b/scripts/architecture/async_blocking.py new file mode 100644 index 000000000..58a201f2a --- /dev/null +++ b/scripts/architecture/async_blocking.py @@ -0,0 +1,224 @@ +"""检测关键 async 路径中新引入的直接阻塞调用。""" + +from __future__ import annotations + +import argparse +import ast +import json +from collections import Counter +from pathlib import Path + +PROJECT_ROOT = Path(__file__).resolve().parents[2] +DEFAULT_BASELINE = PROJECT_ROOT / "tests/fixtures/architecture/async-blocking-baseline.json" +SCAN_ROOTS = ("app/api", "app/agent", "app/application") +BLOCKING_EXACT = { + "open", + "time.sleep", + "subprocess.call", + "subprocess.check_call", + "subprocess.check_output", + "subprocess.Popen", + "subprocess.run", + "os.listdir", + "os.scandir", + "os.walk", + "requests.delete", + "requests.get", + "requests.head", + "requests.patch", + "requests.post", + "requests.put", + "requests.request", +} +BLOCKING_ATTRIBUTES = { + "glob", + "iterdir", + "open", + "read_bytes", + "read_text", + "rglob", + "write_bytes", + "write_text", +} + + +def _call_name(node: ast.Call) -> str: + """把简单名称和属性调用还原为点分文本。""" + parts = [] + target: ast.expr = node.func + while isinstance(target, ast.Attribute): + parts.append(target.attr) + target = target.value + if isinstance(target, ast.Name): + parts.append(target.id) + return ".".join(reversed(parts)) + + +def _root_name(expression: ast.expr) -> str | None: + """返回属性调用最左侧的变量名。""" + while isinstance(expression, ast.Attribute): + expression = expression.value + return expression.id if isinstance(expression, ast.Name) else None + + +class _AsyncPathCollector(ast.NodeVisitor): + """用局部数据流识别 anyio AsyncPath 变量及其派生值。""" + + def __init__(self, function: ast.AsyncFunctionDef) -> None: + """从参数注解初始化 AsyncPath 变量集合。""" + self.paths = { + argument.arg + for argument in (*function.args.posonlyargs, *function.args.args) + if argument.annotation and "AsyncPath" in ast.unparse(argument.annotation) + } + self.path_collections: set[str] = set() + + def visit_Assign(self, node: ast.Assign) -> None: + """识别 AsyncPath 构造和已知路径的 `/` 派生赋值。""" + is_async_path = ( + isinstance(node.value, ast.Call) + and _call_name(node.value).endswith("AsyncPath") + ) or ( + isinstance(node.value, ast.BinOp) + and isinstance(node.value.left, ast.Name) + and node.value.left.id in self.paths + ) + if is_async_path: + for target in node.targets: + if isinstance(target, ast.Name): + self.paths.add(target.id) + self.generic_visit(node) + + def visit_AnnAssign(self, node: ast.AnnAssign) -> None: + """识别 `list[AsyncPath]` 等路径集合。""" + if isinstance(node.target, ast.Name): + annotation = ast.unparse(node.annotation) + if "AsyncPath" in annotation: + if "list" in annotation or "List" in annotation: + self.path_collections.add(node.target.id) + else: + self.paths.add(node.target.id) + self.generic_visit(node) + + def visit_AsyncFor(self, node: ast.AsyncFor) -> None: + """AsyncPath.iterdir 产出的元素仍是 AsyncPath。""" + if isinstance(node.target, ast.Name) and isinstance(node.iter, ast.Call): + receiver = _root_name(node.iter.func) + if receiver in self.paths: + self.paths.add(node.target.id) + self.generic_visit(node) + + def visit_For(self, node: ast.For) -> None: + """从 `list[AsyncPath]` 迭代得到的元素仍是 AsyncPath。""" + if ( + isinstance(node.target, ast.Name) + and isinstance(node.iter, ast.Name) + and node.iter.id in self.path_collections + ): + self.paths.add(node.target.id) + self.generic_visit(node) + + def visit_FunctionDef(self, node: ast.FunctionDef) -> None: + """不分析嵌套同步函数。""" + + def visit_AsyncFunctionDef(self, node: ast.AsyncFunctionDef) -> None: + """不分析嵌套异步函数。""" + + +class _AsyncCallVisitor(ast.NodeVisitor): + """只收集一个 async 函数本体中的阻塞调用,不进入嵌套函数定义。""" + + def __init__(self, async_paths: set[str]) -> None: + """初始化违规计数器和异步文件对象白名单。""" + self.calls: Counter[str] = Counter() + self._async_paths = async_paths + + def visit_Call(self, node: ast.Call) -> None: + """记录命中阻塞名单的调用并继续遍历参数表达式。""" + name = _call_name(node) + attribute = name.rsplit(".", 1)[-1] + receiver = _root_name(node.func) + async_safe = name.startswith("aiofiles.") or receiver in self._async_paths + if not async_safe and ( + name in BLOCKING_EXACT or attribute in BLOCKING_ATTRIBUTES + ): + self.calls[name or attribute] += 1 + self.generic_visit(node) + + def visit_FunctionDef(self, node: ast.FunctionDef) -> None: + """嵌套同步函数不属于外层 async 的直接执行体。""" + + def visit_AsyncFunctionDef(self, node: ast.AsyncFunctionDef) -> None: + """嵌套异步函数由模块级收集器单独治理。""" + + +def _async_functions(tree: ast.Module): + """产出模块顶层及类直接拥有的 async 函数限定名。""" + for node in tree.body: + if isinstance(node, ast.AsyncFunctionDef): + yield node.name, node + elif isinstance(node, ast.ClassDef): + for method in node.body: + if isinstance(method, ast.AsyncFunctionDef): + yield f"{node.name}.{method.name}", method + + +def collect_async_blocking(root: Path = PROJECT_ROOT) -> dict[str, int]: + """扫描关键目录并以文件、函数、调用名聚合存量次数。""" + debt: Counter[str] = Counter() + for scan_root in SCAN_ROOTS: + for path in sorted((root / scan_root).rglob("*.py")): + tree = ast.parse(path.read_text(encoding="utf-8"), filename=str(path)) + relative = path.relative_to(root).as_posix() + for qualname, function in _async_functions(tree): + collector = _AsyncPathCollector(function) + for statement in function.body: + collector.visit(statement) + visitor = _AsyncCallVisitor(collector.paths) + for statement in function.body: + visitor.visit(statement) + for call_name, count in visitor.calls.items(): + debt[f"{relative}:{qualname}:{call_name}"] += count + return dict(sorted(debt.items())) + + +def compare_async_blocking( + baseline: dict[str, int], current: dict[str, int] +) -> list[str]: + """允许存量减少或删除,拒绝新增阻塞调用和调用次数增长。""" + problems = [] + for entry, count in current.items(): + previous = baseline.get(entry) + if previous is None: + problems.append(f"新增 async 阻塞调用:{entry} x{count}") + elif count > previous: + problems.append(f"async 阻塞调用增长:{entry} x{count}>{previous}") + return problems + + +def main() -> int: + """执行 async 阻塞 baseline check 或显式 write。""" + parser = argparse.ArgumentParser(description=__doc__) + parser.add_argument("--write", action="store_true", help="写入当前阻塞债务") + parser.add_argument("--baseline", type=Path, default=DEFAULT_BASELINE) + args = parser.parse_args() + current = collect_async_blocking() + if args.write: + args.baseline.parent.mkdir(parents=True, exist_ok=True) + args.baseline.write_text( + json.dumps(current, ensure_ascii=False, indent=2, sort_keys=True) + "\n", + encoding="utf-8", + ) + print(f"已写入 {args.baseline.relative_to(PROJECT_ROOT)}") + return 0 + baseline = json.loads(args.baseline.read_text(encoding="utf-8")) + problems = compare_async_blocking(baseline, current) + if problems: + print("\n".join(problems)) + return 1 + print("async 阻塞调用 ratchet 通过") + return 0 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/tests/fixtures/architecture/async-blocking-baseline.json b/tests/fixtures/architecture/async-blocking-baseline.json new file mode 100644 index 000000000..89dbb3b6d --- /dev/null +++ b/tests/fixtures/architecture/async-blocking-baseline.json @@ -0,0 +1,3 @@ +{ + "app/agent/middleware/activity_log.py:ActivityLogMiddleware._append_activity:os.open": 1 +} diff --git a/tests/test_async_blocking_gate.py b/tests/test_async_blocking_gate.py new file mode 100644 index 000000000..3b00f5c23 --- /dev/null +++ b/tests/test_async_blocking_gate.py @@ -0,0 +1,31 @@ +"""async 阻塞调用 ratchet 与 debug 模式测试。""" + +import asyncio + +import pytest + +from scripts.architecture.async_blocking import compare_async_blocking + + +def test_async_blocking_ratchet_allows_removal_and_rejects_growth() -> None: + """存量减少合法,新增或增加直接阻塞调用必须失败。""" + baseline = { + "app/api/a.py:stream:Path.read_text": 2, + "app/agent/a.py:run:time.sleep": 1, + } + reduced = {"app/api/a.py:stream:Path.read_text": 1} + increased = { + "app/api/a.py:stream:Path.read_text": 3, + "app/application/a.py:load:requests.get": 1, + } + + assert compare_async_blocking(baseline, reduced) == [] + problems = compare_async_blocking(baseline, increased) + assert any("增长" in problem for problem in problems) + assert any("新增" in problem for problem in problems) + + +@pytest.mark.asyncio +async def test_asyncio_debug_is_enabled_for_async_tests() -> None: + """专项异步测试必须启用慢 callback 和阻塞诊断所需的 debug 模式。""" + assert asyncio.get_running_loop().get_debug() is True