From c2178a985dbb2dd17928f9e5c4e5233cc7cec7e1 Mon Sep 17 00:00:00 2001 From: zhangshy <994452054@qq.com> Date: Fri, 11 Sep 2026 09:49:21 +0800 Subject: [PATCH] =?UTF-8?q?feat:=20=E5=A2=9E=E5=8A=A0=E9=A3=8E=E6=8E=A7?= =?UTF-8?q?=E5=AE=9A=E6=97=B6=E6=89=AB=E6=8F=8F=20Worker=20=E4=B8=8E?= =?UTF-8?q?=E7=8E=AF=E5=A2=83=E9=85=8D=E7=BD=AE?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .env.example | 6 + app/core/config.py | 5 + app/service/risk_scan_schedule_config.py | 25 +++ app/worker/risk_scan_scheduler.py | 195 ++++++++++++++++++ docs/21-风控业务第二版迁移清单.md | 42 +++- .../风控业务演示文档/15-模块验收与演示清单.md | 4 +- docs/风控业务演示文档/18-当前项目完成进度.md | 177 ++++++++++++++++ docs/风控业务演示文档/README.md | 2 + .../service/test_risk_scan_schedule_config.py | 41 ++++ tests/unit/worker/test_risk_scan_scheduler.py | 127 ++++++++++++ 10 files changed, 622 insertions(+), 2 deletions(-) create mode 100644 app/service/risk_scan_schedule_config.py create mode 100644 app/worker/risk_scan_scheduler.py create mode 100644 docs/风控业务演示文档/18-当前项目完成进度.md create mode 100644 tests/unit/service/test_risk_scan_schedule_config.py create mode 100644 tests/unit/worker/test_risk_scan_scheduler.py diff --git a/.env.example b/.env.example index fdcccea..ebc4c09 100644 --- a/.env.example +++ b/.env.example @@ -44,6 +44,12 @@ WORKER_RETRY_LIMIT=3 WORKER_POLL_SECONDS=1 SSE_CHUNK_CHARACTERS=256 +RISK_SCAN_SCHEDULE_ENABLED=false +RISK_SCAN_INTERVAL_MINUTES=5 +RISK_SCAN_RUN_IMMEDIATELY=false +RISK_SCAN_RETRY_LIMIT=2 +RISK_SCAN_POLL_SECONDS=30 + # 限流:按“用户 + 方法 + 路由”在窗口内计数,超限返回 429 RATE_LIMITED + Retry-After。 # 默认 600/60s(10 QPS)是为本机验收与集成测试留足余量的宽松值;生产按真实容量收紧。 # Redis 不可用时自动降级为放行(限流是保护措施,不能因后端故障拒绝正常请求)。 diff --git a/app/core/config.py b/app/core/config.py index 8973c83..67a011e 100644 --- a/app/core/config.py +++ b/app/core/config.py @@ -58,6 +58,11 @@ class Settings(BaseSettings): worker_lease_seconds: int = Field(default=60, gt=0) worker_retry_limit: int = Field(default=3, ge=0) worker_poll_seconds: float = Field(default=1, gt=0) + risk_scan_schedule_enabled: bool = False + risk_scan_interval_minutes: int = Field(default=5, ge=1, le=1440) + risk_scan_run_immediately: bool = False + risk_scan_retry_limit: int = Field(default=2, ge=0, le=5) + risk_scan_poll_seconds: float = Field(default=30, gt=0) model_config = SettingsConfigDict(env_file=".env", env_file_encoding="utf-8", extra="ignore") diff --git a/app/service/risk_scan_schedule_config.py b/app/service/risk_scan_schedule_config.py new file mode 100644 index 0000000..2019be1 --- /dev/null +++ b/app/service/risk_scan_schedule_config.py @@ -0,0 +1,25 @@ +"""风控定时扫描环境配置读取。""" + +from __future__ import annotations + +from dataclasses import dataclass + +from app.core.config import get_settings + + +@dataclass(frozen=True, slots=True) +class RiskScanScheduleConfig: + enabled: bool = False + interval_minutes: int = 5 + run_immediately: bool = False + retry_limit: int = 2 + +async def load_risk_scan_schedule_config() -> RiskScanScheduleConfig: + """读取环境变量配置;缺失时使用默认关闭值。""" + settings = get_settings() + return RiskScanScheduleConfig( + enabled=settings.risk_scan_schedule_enabled, + interval_minutes=settings.risk_scan_interval_minutes, + run_immediately=settings.risk_scan_run_immediately, + retry_limit=settings.risk_scan_retry_limit, + ) diff --git a/app/worker/risk_scan_scheduler.py b/app/worker/risk_scan_scheduler.py new file mode 100644 index 0000000..ee9249a --- /dev/null +++ b/app/worker/risk_scan_scheduler.py @@ -0,0 +1,195 @@ +"""风控规则扫描定时 Worker。 + +该进程不依赖 Web 进程生命周期,通过 MySQL 咨询锁保证同一时刻只有一个 +调度者执行扫描。默认关闭,配置开启后才执行。 +""" + +from __future__ import annotations + +import argparse +import asyncio +import logging +from collections.abc import AsyncIterator, Awaitable, Callable +from contextlib import AbstractAsyncContextManager, asynccontextmanager +from datetime import UTC, datetime, timedelta +from typing import Any +from uuid import uuid4 + +from sqlalchemy import text + +from app.core.config import get_settings +from app.core.contracts import RequestContext +from app.infrastructure.db import SessionFactory, engine +from app.model.audit import InteractionAudit +from app.service.risk_scan_schedule_config import ( + RiskScanScheduleConfig, + load_risk_scan_schedule_config, +) +from app.service.risk_scan_service import RiskScanService + +SCAN_LOCK_NAME = "jr_risk_scan_schedule" +logger = logging.getLogger(__name__) + +ConfigLoader = Callable[[], Awaitable[RiskScanScheduleConfig]] +ScanExecutor = Callable[[], Awaitable[dict[str, int | str]]] +AuditWriter = Callable[[str, dict[str, Any]], Awaitable[None]] +LockFactory = Callable[[], AbstractAsyncContextManager[bool]] + + +@asynccontextmanager +async def mysql_scan_lock() -> AsyncIterator[bool]: + """使用 MySQL 连接级咨询锁约束跨进程并发。""" + async with SessionFactory() as session: + acquired = bool( + await session.scalar( + text("SELECT GET_LOCK(:name, 0)"), + {"name": SCAN_LOCK_NAME}, + ) + ) + try: + yield acquired + finally: + if acquired: + await session.scalar( + text("SELECT RELEASE_LOCK(:name)"), + {"name": SCAN_LOCK_NAME}, + ) + + +async def default_scan_executor() -> dict[str, int | str]: + context = RequestContext( + user_id="0", + trace_id=f"risk-scan-schedule-{uuid4()}", + roles=("system",), + permissions=("risk:alert:scan",), + data_scope="all", + portal="worker", + ) + async with SessionFactory() as session: + return await RiskScanService(session).scan(context) + + +async def default_audit_writer(status: str, detail: dict[str, Any]) -> None: + async with SessionFactory() as session, session.begin(): + session.add( + InteractionAudit( + actor_type="system", + actor_id=None, + target_customer_id=None, + portal="worker", + action_type=f"risk_scan_scheduled_{status}", + detail=detail, + created_at=datetime.now(UTC).replace(tzinfo=None), + ) + ) + + +class RiskScanSchedulerWorker: + def __init__( + self, + *, + config_loader: ConfigLoader = load_risk_scan_schedule_config, + scan_executor: ScanExecutor = default_scan_executor, + audit_writer: AuditWriter = default_audit_writer, + lock_factory: LockFactory = mysql_scan_lock, + now: Callable[[], datetime] | None = None, + ) -> None: + self.config_loader = config_loader + self.scan_executor = scan_executor + self.audit_writer = audit_writer + self.lock_factory = lock_factory + self.now = now or (lambda: datetime.now(UTC)) + self.last_run_at: datetime | None = None + + async def run_once(self, *, force: bool = False) -> bool: + config = await self.config_loader() + if not config.enabled: + return False + current = self.now() + if not force and not self._is_due(config, current): + return False + + async with self.lock_factory() as acquired: + if not acquired: + logger.info("风控定时扫描由其他 Worker 执行,本轮跳过") + return False + return await self._execute_with_retry(config) + + def _is_due( + self, + config: RiskScanScheduleConfig, + current: datetime, + ) -> bool: + if self.last_run_at is None: + return config.run_immediately + return current - self.last_run_at >= timedelta(minutes=config.interval_minutes) + + async def _execute_with_retry( + self, + config: RiskScanScheduleConfig, + ) -> bool: + attempts = config.retry_limit + 1 + last_error: Exception | None = None + for attempt in range(1, attempts + 1): + try: + result = await self.scan_executor() + self.last_run_at = self.now() + await self.audit_writer( + "succeeded", + { + **result, + "attempt": attempt, + "trace_id": f"risk-scan-schedule-{uuid4()}", + }, + ) + return True + except Exception as error: + last_error = error + logger.warning( + "风控定时扫描失败 attempt=%s/%s", + attempt, + attempts, + exc_info=True, + ) + await self.audit_writer( + "failed", + { + "attempts": attempts, + "error_type": type(last_error).__name__ if last_error else "unknown", + "trace_id": f"risk-scan-schedule-{uuid4()}", + }, + ) + return False + + +async def serve(*, once: bool = False, force: bool = False) -> None: + worker = RiskScanSchedulerWorker() + try: + while True: + try: + await worker.run_once(force=force) + except Exception: + logger.warning("风控定时扫描 Worker 轮次失败", exc_info=True) + if once: + raise + if once: + return + await asyncio.sleep(get_settings().risk_scan_poll_seconds) + finally: + await engine.dispose() + + +def main() -> None: + parser = argparse.ArgumentParser(description="奶龙风控规则定时扫描 Worker") + parser.add_argument("--once", action="store_true", help="执行一轮后退出") + parser.add_argument("--force", action="store_true", help="忽略间隔,立即执行一次") + args = parser.parse_args() + logging.basicConfig(level=logging.INFO) + try: + asyncio.run(serve(once=args.once, force=args.force)) + except KeyboardInterrupt: + pass + + +if __name__ == "__main__": + main() diff --git a/docs/21-风控业务第二版迁移清单.md b/docs/21-风控业务第二版迁移清单.md index 821e988..bd6803b 100644 --- a/docs/21-风控业务第二版迁移清单.md +++ b/docs/21-风控业务第二版迁移清单.md @@ -126,7 +126,7 @@ ### R6 定时扫描与后台任务 - [ ] 将定时扫描接入第二版 Worker 和租约机制。 -- [ ] 扫描开关、周期、重试和并发策略使用配置中心。 +- [ ] 扫描开关、周期、重试和并发策略使用环境变量配置。 - [ ] 手动扫描和定时扫描共用同一 Service。 - [ ] 默认关闭,不在 Web 进程启动后台线程。 @@ -955,3 +955,43 @@ tests/unit/api/test_risk_controller.py - R10.5 对话历史与长期留存继续暂缓,全部迁移完成后再处理。 结论:R11 业务迁移和质量门禁已完成。公共底座 3 项 MyPy 例外保持现状,Ruff、全量测试、数据库结构和迁移状态均已通过。 + +### R6.4 定时规则扫描补迁移 + +状态:`[x] 已完成,待确认` + +新增: + +- `app/service/risk_scan_schedule_config.py` +- `app/worker/risk_scan_scheduler.py` +- `tests/unit/service/test_risk_scan_schedule_config.py` +- `tests/unit/worker/test_risk_scan_scheduler.py` + +已实现: + +- 独立定时扫描 Worker,不在 Web 进程启动后台线程。 +- 从 `.env` 环境变量读取开关、周期、立即执行、重试次数和轮询间隔。 +- 默认关闭。 +- 使用 MySQL 咨询锁防止多进程重复扫描。 +- 手动和定时扫描共用同一 `RiskScanService`。 +- 扫描成功和失败写入系统审计。 +- 支持 `--once` 和 `--force` 本地验证参数。 + +环境变量: + +- `RISK_SCAN_SCHEDULE_ENABLED` +- `RISK_SCAN_INTERVAL_MINUTES` +- `RISK_SCAN_RUN_IMMEDIATELY` +- `RISK_SCAN_RETRY_LIMIT` +- `RISK_SCAN_POLL_SECONDS` + +验证结果:专项测试 `6 passed`,Ruff 通过,MyPy 通过,全量测试 `573 passed, 1 skipped`。 + +本地运行方式: + +```powershell +python -m app.worker.risk_scan_scheduler +python -m app.worker.risk_scan_scheduler --once --force +``` + +下一步:将定时扫描配置写入当前 active 配置发布,并同步推送代码。等待确认后执行。 diff --git a/docs/风控业务演示文档/15-模块验收与演示清单.md b/docs/风控业务演示文档/15-模块验收与演示清单.md index 97247f2..092a0eb 100644 --- a/docs/风控业务演示文档/15-模块验收与演示清单.md +++ b/docs/风控业务演示文档/15-模块验收与演示清单.md @@ -34,6 +34,9 @@ ### 规则扫描 - 手动扫描生成符合规则的预警。 +- 定时扫描配置开启后按周期执行。 +- 定时扫描默认关闭,多个 Worker 同时运行时不重复执行。 +- 定时扫描成功和失败写入系统审计。 - 重复扫描不重复生成同一交易和规则的预警。 - 多规则命中时正确合并。 @@ -87,4 +90,3 @@ - Agent 不执行正式处置。 - 所有关键动作都有审计。 - 外部依赖不可用时结构化功能可以降级运行。 - diff --git a/docs/风控业务演示文档/18-当前项目完成进度.md b/docs/风控业务演示文档/18-当前项目完成进度.md new file mode 100644 index 0000000..f54ff54 --- /dev/null +++ b/docs/风控业务演示文档/18-当前项目完成进度.md @@ -0,0 +1,177 @@ +# 当前项目完成进度 + +## 文档功能 + +本文档用于汇总风控业务模块截至当前日期的开发完成度、验证结果、剩余任务、阻塞项和下一阶段计划,作为项目汇报、演示准备和后续合并开发的统一进度基线。 + +## 进度基线 + +| 项目 | 当前状态 | +|---|---| +| 统计日期 | 2026-09-10 | +| 当前分支 | `RM2_develop` | +| 当前提交 | `a94d5c7` | +| 提交信息 | `feat: 迁移奶龙风控业务模块与演示文档` | +| 已推送分支 | `origin/RM2_develop` | +| 已合并分支 | `origin/qyqy_develop` | +| 私有前端 | `private_frontend/`,未提交、未推送 | + +## 总体完成度 + +| 范围 | 完成度 | 说明 | +|---|---:|---| +| 后端业务模块 | 96% | 主要业务功能和定时规则扫描调度均已完成 | +| 私有验证前端 | 90% | 可用于本地功能验证,但不作为公共正式前端 | +| 主项目正式前端 | 10% | 尚未按主项目设计系统和正式页面结构合并 | +| 主项目联调与验收 | 60% | 代码已合并,仍待主项目环境完整联调和正式前端接入 | +| 当前可演示能力 | 90% | 使用现有测试环境可演示完整后端业务和 Agent 流程 | + +整体判断: + +- 作为独立风控业务模块,已基本完成。 +- 作为主项目完整演示版本,正式前端和统一联调仍是主要剩余工作。 + +## 已完成事项 + +### 风控数据和查询 + +- 只读风控模型和 Repository。 +- 风险概览、预警列表、预警详情。 +- 八类证据查询。 +- 客户、产品、风险等级、规则和时间筛选。 +- 分页、客户数据范围和字段脱敏。 + +### 规则扫描和预警生成 + +- RW-003 大额快进快出。 +- RW-007 适当性错配。 +- RW-012 老年客户异常大额赎回。 +- RW-015 非正常时段小额操作。 +- RW-018 频繁交易初筛。 +- 同交易多规则合并和重复扫描防重。 +- 手动扫描接口。 + +### 人工处置和行为分 + +- 确认接收。 +- 进入调查。 +- 关闭误报并记录理由。 +- 完成结案并记录结论。 +- 升级处理。 +- 结案行为分扣减和审计。 + +### 证据、通知和日报 + +- 图片和文档证据归档。 +- 高风险通知记录。 +- 九段式日报。 +- 历史未闭环完整统计。 +- 日报流式生成。 +- 多邮箱校验和 dry-run 邮件发送。 + +### 奶龙风控智能助手 + +- 风险概览、预警查询和预警证据只读工具。 +- 模型自主选择工具和多轮调用。 +- 自然语言客户、产品、规则和时间筛选。 +- 完整客户、产品和规则汇总。 +- 误报、可放行、疑似误判和继续复核研判草案。 +- 结构化研判、回话术和工单摘要。 +- 非法工具、非法参数和协议残留防护。 +- 真实 Agent Run、Worker、结果查询、SSE 和审计验收。 + +### 文档 + +- 风控业务模块演示文档包。 +- 模块总览和主项目接入清单。 +- 规则、状态机、研判、接口、权限和数据边界。 +- Agent 工具和调用流程。 +- 证据、日报、行为分和审计说明。 +- 前端合并提示词与验收约束。 + +## 验证结果 + +| 检查项 | 结果 | +|---|---| +| 全量测试 | `567 passed, 1 skipped` | +| Ruff 静态检查 | 通过 | +| 风控专项测试 | 通过 | +| 真实 Agent Run 验收 | 三类业务对话通过 | +| 数据库结构审计 | 51 张业务表通过 | +| 约束审计 | 通过 | +| 迁移状态 | 与 head 对齐 | +| 依赖检查 | `pip check` 通过 | +| 依赖声明 | 已补充 `python-dotenv`、`python-multipart`、`greenlet`、`tzdata` | + +## 未完成和暂缓事项 + +### 对话历史与长期留存 + +状态:暂缓。 + +- 当前继续使用 MySQL 会话和消息存储。 +- 暂不改为 Redis-only。 +- 后续再做 Redis 活跃会话缓存和长期归档策略。 + +### 定时规则扫描 + +状态:已完成。 + +- 新增独立 `RiskScanSchedulerWorker`,不在 Web 进程启动后台线程。 +- 扫描开关、周期、是否立即执行和重试次数由 `.env` 环境变量控制。 +- 使用 MySQL 咨询锁防止多个 Worker 重复执行。 +- 手动扫描和定时扫描共用同一个 `RiskScanService`。 +- 扫描成功和失败写入系统审计。 +- 默认关闭,配置开启后才执行。 + +### 主项目正式前端 + +状态:未开始。 + +- 正式页面需要使用主项目现有导航、主题、组件、请求封装和权限模型。 +- `private_frontend` 只作为交互参考。 +- 需要补充桌面、移动端和权限场景验收。 + +### 公共底座质量例外 + +状态:确认不修改。 + +- 公共底座仍有 3 项既有 MyPy 类型问题。 +- 不影响当前业务运行。 +- 全量 MyPy 暂时不能完全通过。 + +### 环境依赖 + +- Milvus 未启动时语义记忆降级。 +- Redis 不可用时缓存和限流降级。 +- 模型不可用时使用模板或规则化降级。 +- `.env` 仍有 dotenv 解析警告。 + +## 当前阻塞 + +| 阻塞项 | 影响 | 处理方式 | +|---|---|---| +| 正式前端未合并 | 无法按主项目正式界面演示 | 按前端提示词文档执行合并 | +| 主项目完整联调未完成 | 跨模块权限、导航和接口仍需验证 | 在 qyqy_develop 环境联调 | +| 对话历史暂缓 | 跨轮长期记忆能力有限 | 迁移完成后单独实施 | +| 公共底座 MyPy 例外 | 全量类型检查无法完全通过 | 保持例外或由主项目后续处理 | + +## 下一阶段建议 + +1. 将 `RM2_develop` 与最新 `qyqy_develop` 保持同步。 +2. 将定时规则扫描的环境变量配置纳入主项目部署配置。 +3. 按 `17-前端合并提示词与验收约束.md` 合并正式前端。 +4. 使用主项目真实登录、账号、角色和客户归属完成联调。 +5. 执行桌面端、移动端、权限、降级和完整业务链路验收。 +6. 主项目稳定后再实施对话历史、Redis 缓存和长期留存。 + +## 完成判定 + +风控模块可以认为完成,需要同时满足: + +- 正式前端接入主项目。 +- 真实账号和 RBAC 验证通过。 +- 风控接口、页面和 Agent 完整可用。 +- 关键处置流程和审计可追溯。 +- 全量测试和结构审计持续通过。 +- 已知例外明确记录且不影响交付。 diff --git a/docs/风控业务演示文档/README.md b/docs/风控业务演示文档/README.md index 36d0c46..57663ae 100644 --- a/docs/风控业务演示文档/README.md +++ b/docs/风控业务演示文档/README.md @@ -38,6 +38,7 @@ | 15 | 模块验收与演示清单 | 说明演示前置条件和验收步骤 | | 16 | 已知限制与待办 | 说明当前限制、暂缓项和外部依赖 | | 17 | 前端合并提示词与验收约束 | 指导后续模型识别风控功能并合并前端 | +| 18 | 当前项目完成进度 | 汇总当前完成度、验证结果和剩余任务 | ## 推荐阅读顺序 @@ -47,3 +48,4 @@ 4. 09、10:理解奶龙风控智能助手和工具调用。 5. 11 至 16:理解证据、日报、审计、验收和限制。 6. 17:主项目前端合并时直接提供给模型。 +7. 18:项目汇报、进度同步和下一阶段安排。 diff --git a/tests/unit/service/test_risk_scan_schedule_config.py b/tests/unit/service/test_risk_scan_schedule_config.py new file mode 100644 index 0000000..5fe7686 --- /dev/null +++ b/tests/unit/service/test_risk_scan_schedule_config.py @@ -0,0 +1,41 @@ +from types import SimpleNamespace + +import pytest + +from app.service import risk_scan_schedule_config as config_module +from app.service.risk_scan_schedule_config import ( + RiskScanScheduleConfig, + load_risk_scan_schedule_config, +) + + +def test_risk_scan_schedule_config_defaults_to_disabled() -> None: + config = RiskScanScheduleConfig() + + assert config.enabled is False + assert config.interval_minutes == 5 + assert config.run_immediately is False + assert config.retry_limit == 2 + + +@pytest.mark.asyncio +async def test_risk_scan_schedule_config_reads_environment_settings( + monkeypatch: pytest.MonkeyPatch, +) -> None: + monkeypatch.setattr( + config_module, + "get_settings", + lambda: SimpleNamespace( + risk_scan_schedule_enabled=True, + risk_scan_interval_minutes=15, + risk_scan_run_immediately=True, + risk_scan_retry_limit=1, + ), + ) + + config = await load_risk_scan_schedule_config() + + assert config.enabled is True + assert config.interval_minutes == 15 + assert config.run_immediately is True + assert config.retry_limit == 1 diff --git a/tests/unit/worker/test_risk_scan_scheduler.py b/tests/unit/worker/test_risk_scan_scheduler.py new file mode 100644 index 0000000..9ea73af --- /dev/null +++ b/tests/unit/worker/test_risk_scan_scheduler.py @@ -0,0 +1,127 @@ +from contextlib import asynccontextmanager +from datetime import UTC, datetime + +import pytest + +from app.service.risk_scan_schedule_config import RiskScanScheduleConfig +from app.worker.risk_scan_scheduler import RiskScanSchedulerWorker + + +@asynccontextmanager +async def acquired_lock(): + yield True + + +@asynccontextmanager +async def busy_lock(): + yield False + + +@pytest.mark.asyncio +async def test_disabled_schedule_does_not_scan() -> None: + calls = 0 + + async def scan(): + nonlocal calls + calls += 1 + return {"created_count": 1} + + async def config_loader(): + return RiskScanScheduleConfig(enabled=False) + + worker = RiskScanSchedulerWorker( + config_loader=config_loader, + scan_executor=scan, + audit_writer=lambda *_args: None, # type: ignore[arg-type] + lock_factory=acquired_lock, + ) + + assert await worker.run_once(force=True) is False + assert calls == 0 + + +@pytest.mark.asyncio +async def test_force_scan_runs_once_and_writes_audit() -> None: + calls = 0 + audits: list[tuple[str, dict[str, object]]] = [] + + async def scan(): + nonlocal calls + calls += 1 + return {"created_count": 2, "high_risk_count": 1} + + async def config_loader(): + return RiskScanScheduleConfig(enabled=True, interval_minutes=5) + + async def audit(status, detail): + audits.append((status, detail)) + + worker = RiskScanSchedulerWorker( + config_loader=config_loader, + scan_executor=scan, + audit_writer=audit, + lock_factory=acquired_lock, + now=lambda: datetime(2026, 9, 11, 0, 0, tzinfo=UTC), + ) + + assert await worker.run_once(force=True) is True + assert calls == 1 + assert audits[0][0] == "succeeded" + assert audits[0][1]["created_count"] == 2 + + +@pytest.mark.asyncio +async def test_busy_lock_skips_scan() -> None: + calls = 0 + + async def scan(): + nonlocal calls + calls += 1 + return {"created_count": 1} + + async def config_loader(): + return RiskScanScheduleConfig(enabled=True, run_immediately=True) + + worker = RiskScanSchedulerWorker( + config_loader=config_loader, + scan_executor=scan, + audit_writer=lambda *_args: None, # type: ignore[arg-type] + lock_factory=busy_lock, + ) + + assert await worker.run_once() is False + assert calls == 0 + + +@pytest.mark.asyncio +async def test_scan_retries_before_success() -> None: + attempts = 0 + audits: list[str] = [] + + async def scan(): + nonlocal attempts + attempts += 1 + if attempts == 1: + raise RuntimeError("transient") + return {"created_count": 1} + + async def config_loader(): + return RiskScanScheduleConfig( + enabled=True, + run_immediately=True, + retry_limit=1, + ) + + async def audit(status, _detail): + audits.append(status) + + worker = RiskScanSchedulerWorker( + config_loader=config_loader, + scan_executor=scan, + audit_writer=audit, + lock_factory=acquired_lock, + ) + + assert await worker.run_once() is True + assert attempts == 2 + assert audits == ["succeeded"]