From ab4e4bafbce5972d8d7633078e809fd8847ca65f Mon Sep 17 00:00:00 2001 From: Pigbibi <20649888+Pigbibi@users.noreply.github.com> Date: Fri, 31 Jul 2026 10:22:42 +0800 Subject: [PATCH] fix: route strategy health to optimization watcher Co-Authored-By: Codex --- docs/ai_autonomy_architecture.md | 12 ++ ops/quant-monitor/AGENTS.md | 14 +- ops/quant-monitor/README.md | 12 ++ ops/quant-monitor/scripts/health_cycle.py | 194 +++++++++--------- .../systemd/codex-quant.service.example | 1 - .../tests/test_monitor_fail_closed.py | 83 ++++---- scripts/consume_daily_briefing.py | 22 +- scripts/run_strategy_optimization_watcher.py | 53 ++++- service/briefing_consumer.py | 88 +++++++- service/briefing_dispatch.py | 110 ++++++++-- service/strategy_watch.py | 73 ++++++- tests/test_briefing_dispatch.py | 78 ++++++- tests/test_briefing_model_router.py | 19 +- tests/test_consume_daily_briefing.py | 69 +++++++ .../test_run_strategy_optimization_watcher.py | 35 +++- tests/test_strategy_watch.py | 34 ++- 16 files changed, 700 insertions(+), 197 deletions(-) create mode 100644 tests/test_consume_daily_briefing.py diff --git a/docs/ai_autonomy_architecture.md b/docs/ai_autonomy_architecture.md index 70044706..d64f6882 100644 --- a/docs/ai_autonomy_architecture.md +++ b/docs/ai_autonomy_architecture.md @@ -278,6 +278,18 @@ GitHub Codex App 是唯一 AI PR reviewer。AIAuditBridge 只保留月度审计 这样可以把策略优化先收敛为可审计、可回放的建议流,再逐步扩展到受控执行面。 +Quant Monitor 的具体路由遵循同一边界: + +- lifecycle score、status 或 drift 越过阈值,只生成 monitoring evidence; +- evidence 写入对应策略仓的去重 optimization issue,供 AI 诊断和有界、no-order + 实验规划; +- 监控不因分数低或漂移高直接发送 Telegram; +- 数据/可信工件不可用、circuit breaker 或 optimization issue 记录失败,才升级为 + 人工运维告警。 + +“监测到劣化”只代表产生优化触发证据,不代表已授权启动优化周期,更不代表可以影响 +broker、订单、资金或 live deployment。 + --- ## 4. 可落地的阶段性改造计划 diff --git a/ops/quant-monitor/AGENTS.md b/ops/quant-monitor/AGENTS.md index 329cb8ff..66cf5cb1 100644 --- a/ops/quant-monitor/AGENTS.md +++ b/ops/quant-monitor/AGENTS.md @@ -13,7 +13,7 @@ VPS Codex 定时监控(`codex-quant.timer` 每 30 分钟)+ 收盘简报(`c | `QUANT_PLATFORM_KIT_ROOT` | `~/Projects/QuantPlatformKit` | | `GLOBAL_TELEGRAM_CHAT_ID` | systemd 注入,勿提交 git | | `GH_TOKEN` | `gh` 拉仓 + 开 Issue | -| `QSL_MONITOR_ISSUE_OWNER` | 3σ 漂移 Issue @ 的用户(默认 `Pigbibi`) | +| `QSL_GITHUB_REPO` | 非策略类 briefing Issue 的默认仓库;策略证据按 domain 写入对应策略仓 | 凭证:`scripts/load_telegram_env.sh` 从 GCP `quant-sentinel-telegram-bot-token` 加载。 @@ -21,15 +21,16 @@ VPS Codex 定时监控(`codex-quant.timer` 每 30 分钟)+ 收盘简报(`c 1. `sync_strategy_repos.sh` — `git pull` 四策略仓 + QPK 2. `health_cycle.py` — `build_dashboard` + `run_drift_detection` -3. lifecycle `overall_score < 60` → Telegram 量化哨兵(研究/回测健康,不代表平台执行故障) -4. drift ≥ 0.50(~2σ)→ `create_issues_for_domain` 开 Issue(不 @) -5. drift ≥ 0.75(~3σ)→ Telegram + Issue **@owner** +3. lifecycle `overall_score < 60` 或 drift ≥ 0.50 → 生成 monitoring evidence +4. evidence → 对应策略仓的去重、issue-only AI optimization proposal +5. 分数或漂移本身不发 Telegram;数据/工件不可用或 Issue 记录失败才通知人工 ## 每日收盘后(daily_briefing_pipeline.sh) 1. `daily_briefing_builder.py` → `data/daily-reports/YYYY-MM-DD/.json` 2. `AIAuditBridge/scripts/consume_daily_briefing.py --dispatch` -3. 正常 → quiet;review → Issue;critical → Telegram +3. 正常 → quiet;review/critical → issue-only AI optimization proposal +4. data unavailable、circuit breaker 或 proposal 记录失败 → Telegram ## 部署 @@ -41,4 +42,5 @@ bash ops/quant-monitor/scripts/deploy_to_vps.sh - 不要手填 token/chat id 到仓库 - 报警只走量化哨兵 bot -- 日报默认不通知人,除非 `briefing_consumer` 判定 critical +- 策略健康证据只进入可审计的 issue-only 优化队列,不自动改策略、参数、仓位或部署 +- 策略劣化记录成功后不通知人;只对数据/运行风险和记录失败 fail-closed 通知 diff --git a/ops/quant-monitor/README.md b/ops/quant-monitor/README.md index f7099dbf..e66c73d3 100644 --- a/ops/quant-monitor/README.md +++ b/ops/quant-monitor/README.md @@ -27,6 +27,18 @@ contract 和大小限制校验后原子切换;代码仓库与 lifecycle 数据 Token 从 GCP Secret `quant-sentinel-telegram-bot-token` 加载;**不要**把 token 或 chat id 写进 git。 +告警路由: + +| 事件 | 处置 | +|------|------| +| lifecycle score / drift 劣化 | 写入对应策略仓的去重、issue-only AI optimization proposal | +| 数据或可信工件不可用 | Telegram | +| circuit breaker / runtime risk | Telegram | +| optimization proposal 记录失败 | Telegram | + +监控证据只触发研究审查,不自动修改策略代码、live 参数、仓位、风险预算,不自动 +merge 或 deploy。成功记录策略劣化后,monitor 正常结束,不再重复通知人工。 + | 变量 | 说明 | |------|------| | `QUANT_SENTINEL_TELEGRAM_SECRET_NAME` | 默认 `quant-sentinel-telegram-bot-token` | diff --git a/ops/quant-monitor/scripts/health_cycle.py b/ops/quant-monitor/scripts/health_cycle.py index 47004ae2..09a891e3 100755 --- a/ops/quant-monitor/scripts/health_cycle.py +++ b/ops/quant-monitor/scripts/health_cycle.py @@ -1,5 +1,5 @@ #!/usr/bin/env python3 -"""VPS health cycle — roadmap task 7 (scores, drift, issues, Telegram).""" +"""VPS health cycle — route strategy evidence to issue-only optimization monitoring.""" from __future__ import annotations @@ -7,7 +7,6 @@ import json import os import re -import subprocess import sys import tempfile from datetime import datetime, timedelta, timezone @@ -151,25 +150,10 @@ def _alert_fingerprint(lines: list[str]) -> str: return hashlib.sha256(payload.encode("utf-8")).hexdigest() -def _strategy_health_alert(row: dict[str, Any]) -> tuple[str, str] | None: - try: - score = float(row.get("overall_score")) - except (TypeError, ValueError): - return None - if score >= SCORE_ALERT: - return None - profile = str(row.get("strategy_profile") or "?") - domain = str(row.get("domain") or "?") - return ( - f"[{domain}] {profile}: lifecycle_health_score={score:.1f}", - f"strategy_health_below_{SCORE_ALERT:g}:{domain}:{profile}", - ) - - def _build_alert_body(lines: list[str]) -> str: return ( - "🚨 quant-monitor strategy_lifecycle\n" - "• scope: research/backtest lifecycle health; not platform runtime execution\n" + "🚨 quant-monitor operational\n" + "• scope: data/evidence or optimization-record delivery failure\n" + "\n".join(f"• {line}" for line in lines) ) @@ -213,17 +197,83 @@ def _clear_alert(root: Path) -> None: pass -def _create_issues_for_available_domains( +def _build_monitoring_findings( + strategies: list[dict[str, Any]], drift_results: dict[str, list[Any]], - create_issues_for_domain, - *, - domains=DOMAINS, -) -> list[dict[str, Any]]: - issue_results: list[dict[str, Any]] = [] - for domain in domains: - if domain in drift_results: - issue_results.extend(create_issues_for_domain(domain)) - return issue_results +) -> list[Any]: + from service.strategy_watch import build_strategy_monitoring_finding + + records: dict[tuple[str, str], dict[str, Any]] = {} + metric_keys = ( + "overall_score", + "performance_score", + "risk_score", + "decay_score", + "stability_score", + "operational_score", + "status", + "as_of", + ) + for row in strategies: + try: + score = float(row.get("overall_score")) + except (TypeError, ValueError): + continue + if score >= SCORE_ALERT: + continue + domain = str(row.get("domain") or "").strip() + profile = str(row.get("strategy_profile") or "").strip() + if not domain or not profile: + continue + record = records.setdefault( + (domain, profile), + {"metrics": {}, "signals": [], "severity": "medium", "generated_at": ""}, + ) + record["metrics"].update({key: row[key] for key in metric_keys if key in row}) + record["signals"].append( + { + "metric": "overall_score", + "reason": f"overall_score={score:.1f} is below monitoring threshold {SCORE_ALERT:.1f}", + } + ) + record["generated_at"] = str(row.get("as_of") or "") + if str(row.get("status") or "").lower() == "critical" or score <= 40.0: + record["severity"] = "high" + + for domain, drifts in drift_results.items(): + for drift in drifts: + score = float(drift.drift_score or 0.0) + if score < DRIFT_REVIEW: + continue + profile = str(drift.strategy_profile or "").strip() + if not profile: + continue + record = records.setdefault( + (domain, profile), + {"metrics": {}, "signals": [], "severity": "medium", "generated_at": ""}, + ) + record["metrics"]["drift_score"] = score + record["signals"].append( + { + "metric": "drift_score", + "reason": f"drift_score={score:.2f} exceeds monitoring threshold {DRIFT_REVIEW:.2f}", + } + ) + if score >= DRIFT_CRITICAL: + record["severity"] = "high" + + return [ + build_strategy_monitoring_finding( + domain=domain, + profile=profile, + severity=record["severity"], + metrics=record["metrics"], + signals=record["signals"], + source="quant-monitor/health-cycle", + generated_at=record["generated_at"], + ) + for (domain, profile), record in sorted(records.items()) + ] def _send_telegram(text: str) -> bool: @@ -239,46 +289,16 @@ def _send_telegram(text: str) -> bool: return False -def _create_owner_issue(*, title: str, body: str) -> str | None: - repo = (os.environ.get("QSL_GITHUB_REPO") or "QuantStrategyLab/CnEquityStrategies").strip() - owner = (os.environ.get("QSL_MONITOR_ISSUE_OWNER") or "Pigbibi").strip() - full_body = f"{body}\n\ncc @{owner}" - try: - out = subprocess.check_output( - [ - "gh", - "issue", - "create", - "--repo", - repo, - "--title", - title, - "--body", - full_body, - "--label", - "monitoring", - "--label", - "drift-critical", - ], - text=True, - stderr=subprocess.STDOUT, - env={**os.environ, "GH_PROMPT": "disabled"}, - ) - return out.strip() - except Exception: - return None - - def main() -> int: root = Path(os.environ.get("QUANT_MONITOR_ROOT") or Path(__file__).resolve().parents[1]) out_dir = root / "data" / "health" dash_dir = out_dir / "dashboard" out_dir.mkdir(parents=True, exist_ok=True) - from quant_platform_kit.strategy_lifecycle.codex_integration import create_issues_for_domain from quant_platform_kit.strategy_lifecycle.drift_detector import run_drift_detection from quant_platform_kit.strategy_lifecycle.health_dashboard import build_dashboard from quant_platform_kit.strategy_lifecycle.performance_monitor import run_monitor + from scripts.run_strategy_optimization_watcher import dispatch_strategy_watch_findings ready_domains, artifact_errors = _load_lifecycle_artifact_status(root) snapshot_results, drift_results, lifecycle_errors = _refresh_and_collect_drift( @@ -338,41 +358,14 @@ def main() -> int: elif not json_path.is_file(): collector_payload_invalid = True - telegram_lines: list[str] = [] - critical_lines: list[str] = [] alert_identities: list[str] = [] - - for row in strategies: - alert = _strategy_health_alert(row) - if alert: - line, identity = alert - telegram_lines.append(line) - alert_identities.append(identity) - - for domain in DOMAINS: - drifts = drift_results.get(domain, []) - for drift in drifts: - score = float(drift.drift_score or 0.0) - label = f"[{domain}] {drift.strategy_profile}: drift_score={score:.2f}" - if score >= DRIFT_CRITICAL: - critical_lines.append(label) - alert_identities.append( - f"critical_drift:{domain}:{drift.strategy_profile}" - ) - elif score >= DRIFT_REVIEW: - pass # tracked via create_issues_for_domain below - - issue_results = _create_issues_for_available_domains( - drift_results, - create_issues_for_domain, + monitoring_findings = _build_monitoring_findings(strategies, drift_results) + optimization_watch = dispatch_strategy_watch_findings( + monitoring_findings, + dry_run=False, + comment_existing=False, ) - for line in critical_lines: - _create_owner_issue( - title=f"[monitor] critical drift — {line}", - body=f"Quant-monitor detected critical drift.\n\n- {line}", - ) - data_error_lines: list[str] = [] for error in data_errors: data_error_lines.append( @@ -384,7 +377,16 @@ def main() -> int: if collector_payload_invalid: data_error_lines.append("[collector] dashboard_data_unavailable") alert_identities.append("data_error:collector:dashboard_data_unavailable") - notify_lines = telegram_lines + critical_lines + data_error_lines + optimization_error_lines: list[str] = [] + optimization_errors = int(optimization_watch.get("errors") or 0) + if optimization_errors: + optimization_error_lines.append( + f"[optimization] issue_record_failed ({optimization_errors})" + ) + alert_identities.append( + f"optimization_error:issue_record_failed:{optimization_errors}" + ) + notify_lines = data_error_lines + optimization_error_lines telegram_sent = False duplicate_alert_suppressed = False if notify_lines: @@ -407,7 +409,11 @@ def main() -> int: "duplicate_alert_suppressed": duplicate_alert_suppressed, "data_errors": data_errors, "snapshot_count": sum(len(rows) for rows in snapshot_results.values()), - "issues_created": len([r for r in issue_results if r.get("issue_url")]), + "optimization_findings": len(monitoring_findings), + "optimization_issues_created": len( + [result for result in optimization_watch.get("issues", []) if result.get("created")] + ), + "optimization_issue_errors": optimization_errors, "ok": not notify_lines and not collector_payload_invalid, "collector_payload_valid": not collector_payload_invalid, "snapshot_data_status": normalized_payload.get("data_status"), diff --git a/ops/quant-monitor/systemd/codex-quant.service.example b/ops/quant-monitor/systemd/codex-quant.service.example index dc71926b..0a0edd98 100644 --- a/ops/quant-monitor/systemd/codex-quant.service.example +++ b/ops/quant-monitor/systemd/codex-quant.service.example @@ -16,7 +16,6 @@ Environment=QUANT_SENTINEL_GCP_PROJECT=firstradequant # Set on VPS only — do not commit real chat id to git: Environment=GLOBAL_TELEGRAM_CHAT_ID= Environment=QSL_GITHUB_REPO=QuantStrategyLab/CnEquityStrategies -Environment=QSL_MONITOR_ISSUE_OWNER=Pigbibi RuntimeDirectory=quant-monitor RuntimeDirectoryMode=0750 ExecStartPre=/bin/bash /home/ubuntu/Projects/AIAuditBridge/ops/quant-monitor/scripts/load_telegram_env.sh /run/quant-monitor/telegram.env diff --git a/ops/quant-monitor/tests/test_monitor_fail_closed.py b/ops/quant-monitor/tests/test_monitor_fail_closed.py index 6bac09d6..b19e981c 100644 --- a/ops/quant-monitor/tests/test_monitor_fail_closed.py +++ b/ops/quant-monitor/tests/test_monitor_fail_closed.py @@ -180,43 +180,54 @@ def test_health_cycle_alert_fingerprint_is_deduplicated_until_recovery(self) -> HEALTH_CYCLE._clear_alert(root) self.assertFalse(HEALTH_CYCLE._is_duplicate_alert(root, fingerprint)) - def test_health_cycle_score_changes_keep_the_same_incident_fingerprint(self) -> None: - first_line, first_identity = HEALTH_CYCLE._strategy_health_alert( - { - "domain": "us_equity", - "strategy_profile": "example", - "overall_score": 59.9, - } - ) - second_line, second_identity = HEALTH_CYCLE._strategy_health_alert( - { - "domain": "us_equity", - "strategy_profile": "example", - "overall_score": 59.8, - } + def test_health_cycle_builds_issue_only_monitoring_finding(self) -> None: + findings = HEALTH_CYCLE._build_monitoring_findings( + [ + { + "domain": "us_equity", + "strategy_profile": "global_etf_rotation", + "status": "critical", + "overall_score": 14.2, + "performance_score": 0.0, + } + ], + {}, ) - self.assertNotEqual(first_line, second_line) - self.assertEqual( - HEALTH_CYCLE._alert_fingerprint([first_identity]), - HEALTH_CYCLE._alert_fingerprint([second_identity]), + self.assertEqual(len(findings), 1) + self.assertEqual(findings[0].snapshot.repo, "QuantStrategyLab/UsEquityStrategies") + self.assertEqual(findings[0].finding_type, "monitoring_trigger") + self.assertEqual(findings[0].severity, "high") + + def test_health_cycle_merges_drift_into_strategy_monitoring_finding(self) -> None: + drift = types.SimpleNamespace( + strategy_profile="example", + drift_score=0.8, ) - def test_health_cycle_labels_strategy_scores_as_lifecycle_not_runtime(self) -> None: - line, _identity = HEALTH_CYCLE._strategy_health_alert( - { - "domain": "us_equity", - "strategy_profile": "example", - "overall_score": 59.9, - } + findings = HEALTH_CYCLE._build_monitoring_findings( + [ + { + "domain": "crypto", + "strategy_profile": "example", + "status": "review", + "overall_score": 45.0, + } + ], + {"crypto": [drift]}, ) - body = HEALTH_CYCLE._build_alert_body([line]) + self.assertEqual(len(findings), 1) + self.assertEqual(findings[0].snapshot.current_metrics["drift_score"], 0.8) + self.assertEqual(findings[0].severity, "high") + self.assertEqual(len(findings[0].signals), 2) + + def test_health_cycle_telegram_body_is_operational_only(self) -> None: + body = HEALTH_CYCLE._build_alert_body(["[collector] dashboard_data_unavailable"]) - self.assertIn("quant-monitor strategy_lifecycle", body) - self.assertIn("research/backtest lifecycle health", body) - self.assertIn("not platform runtime execution", body) - self.assertIn("lifecycle_health_score=59.9", body) + self.assertIn("quant-monitor operational", body) + self.assertIn("data/evidence or optimization-record delivery", body) + self.assertNotIn("strategy_lifecycle", body) def test_health_cycle_non_object_alert_state_is_a_cache_miss(self) -> None: with tempfile.TemporaryDirectory() as tmp: @@ -227,18 +238,6 @@ def test_health_cycle_non_object_alert_state_is_a_cache_miss(self) -> None: state_path.write_text(json.dumps(payload), encoding="utf-8") self.assertFalse(HEALTH_CYCLE._is_duplicate_alert(root, "fingerprint")) - def test_health_cycle_skips_issue_creation_when_drift_is_unavailable(self) -> None: - created_for: list[str] = [] - - results = HEALTH_CYCLE._create_issues_for_available_domains( - {"us_equity": []}, - lambda domain: created_for.append(domain) or [], - domains=("us_equity", "crypto"), - ) - - self.assertEqual(results, []) - self.assertEqual(created_for, ["us_equity"]) - def test_daily_briefing_marks_missing_dashboard_unavailable(self) -> None: with tempfile.TemporaryDirectory() as tmp: root = Path(tmp) diff --git a/scripts/consume_daily_briefing.py b/scripts/consume_daily_briefing.py index fbf96138..43807720 100644 --- a/scripts/consume_daily_briefing.py +++ b/scripts/consume_daily_briefing.py @@ -14,6 +14,17 @@ from service.dual_review_orchestrator import orchestrate_from_payload +def _dispatch_failed(summary: object) -> bool: + if not isinstance(summary, dict): + return True + if summary.get("errors"): + return True + optimization_watch = summary.get("optimization_watch") + return isinstance(optimization_watch, dict) and int( + optimization_watch.get("errors") or 0 + ) > 0 + + def main(argv: list[str] | None = None) -> int: parser = argparse.ArgumentParser(description="Consume quant-monitor daily briefing reports.") parser.add_argument( @@ -56,7 +67,16 @@ def main(argv: list[str] | None = None) -> int: dual_results.append(entry) payload["dual_review"] = summarize_dual_review_runs(dual_results) print(json.dumps(payload, ensure_ascii=False, indent=2)) - exit_code = 0 if result.action.value == "quiet" else 2 + if result.action.value == "quiet": + exit_code = 0 + elif ( + result.action.value == "github_issue" + and args.dispatch + and not _dispatch_failed(payload.get("dispatch")) + ): + exit_code = 0 + else: + exit_code = 2 if args.dual_review and payload.get("dual_review", {}).get("disagreements"): exit_code = 2 return exit_code diff --git a/scripts/run_strategy_optimization_watcher.py b/scripts/run_strategy_optimization_watcher.py index bb0696dd..afe012f3 100755 --- a/scripts/run_strategy_optimization_watcher.py +++ b/scripts/run_strategy_optimization_watcher.py @@ -16,7 +16,13 @@ if str(ROOT) not in sys.path: sys.path.insert(0, str(ROOT)) -from service.strategy_watch import evaluate_strategy_watch, finding_to_automation_task, issue_for_task, watcher_issue_key # noqa: E402 +from service.strategy_watch import ( # noqa: E402 + StrategyWatchFinding, + evaluate_strategy_watch, + finding_to_automation_task, + issue_for_task, + watcher_issue_key, +) REPO_RE = re.compile(r"^[A-Za-z0-9_.-]+/[A-Za-z0-9_.-]+$") @@ -185,21 +191,18 @@ def _payload_for_source_repo(payload: dict[str, Any], source_repo: str) -> dict[ return normalized -def run_watcher( - payload: dict[str, Any], +def dispatch_strategy_watch_findings( + findings: list[StrategyWatchFinding], *, source_repo: str = "", dry_run: bool = True, + comment_existing: bool = True, create_issue: Callable[[str, str, str], str] = create_github_issue, comment_issue: Callable[[str, str, str], str] = comment_github_issue, list_issues: Callable[[str], dict[str, str]] = list_open_issue_urls, ) -> dict[str, Any]: - if not dry_run and not source_repo: - raise ValueError("source_repo is required for non-dry-run strategy watcher runs") if source_repo and not REPO_RE.fullmatch(source_repo): raise ValueError("source_repo must be in owner/name form") - watch_payload = _payload_for_source_repo(payload, source_repo) - findings = evaluate_strategy_watch(watch_payload) issues: list[dict[str, Any]] = [] open_issue_cache: dict[str, dict[str, str]] = {} for finding in findings: @@ -207,6 +210,8 @@ def run_watcher( issue = issue_for_task(task) issue_key = watcher_issue_key(task) repo = source_repo or finding.snapshot.repo + if not REPO_RE.fullmatch(repo): + raise ValueError("finding repository must be in owner/name form") issue_result: dict[str, Any] = { "repo": repo, "title": issue["title"], @@ -223,9 +228,12 @@ def run_watcher( existing_url = open_issue_cache[repo].get(issue_key, "") if existing_url: issue_result["existing_url"] = existing_url - issue_result["comment_url"] = comment_issue(repo, existing_url, issue["body"]) - issue_result["commented"] = True - issue_result["skipped_reason"] = "open issue already exists; appended watcher update" + if comment_existing: + issue_result["comment_url"] = comment_issue(repo, existing_url, issue["body"]) + issue_result["commented"] = True + issue_result["skipped_reason"] = "open issue already exists; appended watcher update" + else: + issue_result["skipped_reason"] = "open issue already records this strategy" else: issue_result["url"] = create_issue(repo, issue["title"], issue["body"]) open_issue_cache[repo][issue_key] = str(issue_result["url"]) @@ -243,6 +251,31 @@ def run_watcher( } +def run_watcher( + payload: dict[str, Any], + *, + source_repo: str = "", + dry_run: bool = True, + create_issue: Callable[[str, str, str], str] = create_github_issue, + comment_issue: Callable[[str, str, str], str] = comment_github_issue, + list_issues: Callable[[str], dict[str, str]] = list_open_issue_urls, +) -> dict[str, Any]: + if not dry_run and not source_repo: + raise ValueError("source_repo is required for non-dry-run strategy watcher runs") + if source_repo and not REPO_RE.fullmatch(source_repo): + raise ValueError("source_repo must be in owner/name form") + watch_payload = _payload_for_source_repo(payload, source_repo) + findings = evaluate_strategy_watch(watch_payload) + return dispatch_strategy_watch_findings( + findings, + source_repo=source_repo, + dry_run=dry_run, + create_issue=create_issue, + comment_issue=comment_issue, + list_issues=list_issues, + ) + + def main() -> int: try: input_path = resolve_input_path( diff --git a/service/briefing_consumer.py b/service/briefing_consumer.py index 0b27daf4..8e6b6089 100644 --- a/service/briefing_consumer.py +++ b/service/briefing_consumer.py @@ -1,9 +1,9 @@ -"""Consume quant-monitor daily briefing JSON and classify alert severity. +"""Consume quant-monitor daily briefing JSON and classify routing. Roadmap task 10b: - all normal → quiet -- deviation ~2σ (review / elevated drift) → github_issue -- deviation >3σ / circuit breaker / critical → telegram +- strategy health / drift degradation → issue-only optimization monitor +- data unavailable / circuit breaker → telegram """ from __future__ import annotations @@ -48,6 +48,10 @@ class BriefingFinding: reason: str strategy_profile: str = "" domain: str = "" + kind: str = "monitoring" + severity: str = "medium" + metrics: dict[str, Any] = field(default_factory=dict) + signals: list[dict[str, Any]] = field(default_factory=list) def to_dict(self) -> dict[str, Any]: return { @@ -56,6 +60,10 @@ def to_dict(self) -> dict[str, Any]: "reason": self.reason, "strategy_profile": self.strategy_profile, "domain": self.domain, + "kind": self.kind, + "severity": self.severity, + "metrics": self.metrics, + "signals": self.signals, } @@ -114,29 +122,61 @@ def _classify_strategy( level = BriefingAction.QUIET reasons: list[str] = [] + signals: list[dict[str, Any]] = [] + severity = "medium" + kind = "strategy_monitoring" if status == "critical": - level = BriefingAction.TELEGRAM + level = BriefingAction.GITHUB_ISSUE + severity = "high" reasons.append("status=critical") + signals.append({"metric": "status", "reason": "status=critical"}) elif status == "review": level = _max_level(level, BriefingAction.GITHUB_ISSUE) reasons.append("status=review") + signals.append({"metric": "status", "reason": "status=review"}) if drift_score is not None: if drift_score >= _DRIFT_CRITICAL: - level = BriefingAction.TELEGRAM + level = _max_level(level, BriefingAction.GITHUB_ISSUE) + severity = "high" reasons.append(f"drift_score={drift_score:.2f}") + signals.append( + { + "metric": "drift_score", + "reason": f"drift_score={drift_score:.2f} exceeds {_DRIFT_CRITICAL:.2f}", + } + ) elif drift_score >= _DRIFT_WARN: level = _max_level(level, BriefingAction.GITHUB_ISSUE) reasons.append(f"drift_score={drift_score:.2f}") + signals.append( + { + "metric": "drift_score", + "reason": f"drift_score={drift_score:.2f} exceeds {_DRIFT_WARN:.2f}", + } + ) if overall_score is not None: if overall_score <= _SCORE_CRITICAL: - level = BriefingAction.TELEGRAM + level = _max_level(level, BriefingAction.GITHUB_ISSUE) + severity = "high" reasons.append(f"overall_score={overall_score:.1f}") + signals.append( + { + "metric": "overall_score", + "reason": f"overall_score={overall_score:.1f} is below {_SCORE_CRITICAL:.1f}", + } + ) elif overall_score <= _SCORE_REVIEW: level = _max_level(level, BriefingAction.GITHUB_ISSUE) reasons.append(f"overall_score={overall_score:.1f}") + signals.append( + { + "metric": "overall_score", + "reason": f"overall_score={overall_score:.1f} is below {_SCORE_REVIEW:.1f}", + } + ) flags = strategy.get("risk_flags") or strategy.get("alerts") or () if isinstance(flags, Mapping): @@ -145,16 +185,34 @@ def _classify_strategy( text = str(flag).lower() if any(keyword in text for keyword in _CIRCUIT_KEYWORDS): level = BriefingAction.TELEGRAM + kind = "runtime_risk" + severity = "high" reasons.append(f"flag={flag}") + signals.append({"metric": "risk_flag", "reason": f"flag={flag}"}) if level == BriefingAction.QUIET: return None + metric_keys = ( + "overall_score", + "performance_score", + "risk_score", + "decay_score", + "stability_score", + "operational_score", + "drift_score", + "status", + "as_of", + ) return BriefingFinding( source=source, level=level, reason="; ".join(reasons) or "anomaly", strategy_profile=profile, domain=str(strategy.get("domain") or domain), + kind=kind, + severity=severity, + metrics={key: strategy[key] for key in metric_keys if key in strategy}, + signals=signals, ) @@ -174,7 +232,15 @@ def _classify_report_payload( else BriefingAction.GITHUB_ISSUE ) findings.append( - BriefingFinding(source=source, level=level, reason=error, domain=str(payload.get("domain") or "")) + BriefingFinding( + source=source, + level=level, + reason=error, + domain=str(payload.get("domain") or ""), + kind="data_unavailable" if data_unavailable else "data_quality", + severity="high" if level == BriefingAction.TELEGRAM else "medium", + signals=[{"metric": "data_status", "reason": error}], + ) ) return findings @@ -189,16 +255,18 @@ def _classify_report_payload( findings.append(finding) summary = payload.get("summary") - if isinstance(summary, Mapping): + if isinstance(summary, Mapping) and not strategies: critical = int(summary.get("critical") or 0) review = int(summary.get("review") or 0) if critical > 0: findings.append( BriefingFinding( source=source, - level=BriefingAction.TELEGRAM, + level=BriefingAction.GITHUB_ISSUE, reason=f"summary critical={critical}", domain=domain, + kind="strategy_monitoring_summary", + severity="high", ) ) elif review > 0: @@ -208,6 +276,7 @@ def _classify_report_payload( level=BriefingAction.GITHUB_ISSUE, reason=f"summary review={review}", domain=domain, + kind="strategy_monitoring_summary", ) ) @@ -236,6 +305,7 @@ def consume_briefing_dir(report_dir: str | Path, *, day: str = "") -> BriefingCo source=file_path.name, level=BriefingAction.GITHUB_ISSUE, reason=f"invalid_json: {exc}", + kind="data_quality", ) ) continue diff --git a/service/briefing_dispatch.py b/service/briefing_dispatch.py index 2eb215c3..5f031c4e 100644 --- a/service/briefing_dispatch.py +++ b/service/briefing_dispatch.py @@ -10,7 +10,9 @@ import urllib.request from typing import Any -from service.briefing_consumer import BriefingAction, BriefingConsumptionResult +from scripts.run_strategy_optimization_watcher import dispatch_strategy_watch_findings +from service.briefing_consumer import BriefingAction, BriefingConsumptionResult, BriefingFinding +from service.strategy_watch import StrategyWatchFinding, build_strategy_monitoring_finding _REPOSITORY_RE = re.compile( r"[A-Za-z0-9][A-Za-z0-9_-]*(?:\.[A-Za-z0-9_-]+)*/" @@ -63,14 +65,17 @@ def _format_telegram_body(result: BriefingConsumptionResult) -> str: return "\n".join(lines) -def _format_github_body(result: BriefingConsumptionResult) -> str: +def _format_github_body( + result: BriefingConsumptionResult, + findings: list[BriefingFinding] | None = None, +) -> str: lines = [ f"## Daily briefing alerts ({result.day})", "", f"Report dir: `{result.report_dir}`", "", ] - for finding in result.findings: + for finding in result.findings if findings is None else findings: if finding.level != BriefingAction.GITHUB_ISSUE: continue lines.append( @@ -102,7 +107,7 @@ def send_telegram_alert(*, text: str, token: str, chat_ids: tuple[str, ...]) -> return ok -def create_github_issue(*, title: str, body: str, labels: tuple[str, ...] = ("briefing", "monitoring")) -> str | None: +def create_github_issue(*, title: str, body: str, labels: tuple[str, ...] = ()) -> str | None: repo = str( os.environ.get("QSL_GITHUB_REPO") or os.environ.get("GITHUB_REPOSITORY") @@ -138,6 +143,31 @@ def shutil_which(name: str) -> str | None: return which(name) +def _strategy_monitoring_findings( + result: BriefingConsumptionResult, +) -> list[StrategyWatchFinding]: + findings: list[StrategyWatchFinding] = [] + for finding in result.findings: + if ( + finding.level != BriefingAction.GITHUB_ISSUE + or finding.kind != "strategy_monitoring" + or not finding.strategy_profile + ): + continue + findings.append( + build_strategy_monitoring_finding( + domain=finding.domain, + profile=finding.strategy_profile, + severity=finding.severity, + metrics=finding.metrics, + signals=finding.signals or [{"metric": "briefing", "reason": finding.reason}], + source=f"quant-monitor/daily-briefing/{finding.source}", + generated_at=str(finding.metrics.get("as_of") or result.day), + ) + ) + return findings + + def dispatch_briefing_result( result: BriefingConsumptionResult, *, @@ -148,6 +178,9 @@ def dispatch_briefing_result( "action": result.action.value, "telegram_sent": False, "github_issue": None, + "optimization_watch": None, + "operational_fallback_sent": False, + "errors": [], "skipped": [], } @@ -155,26 +188,67 @@ def dispatch_briefing_result( summary["skipped"].append("quiet") return summary - if result.action == BriefingAction.TELEGRAM: - token = _telegram_token() - chat_ids = _telegram_chat_ids() - text = _format_telegram_body(result) - if dry_run: - summary["telegram_dry_run"] = text - return summary - if token and chat_ids: - summary["telegram_sent"] = send_telegram_alert(text=text, token=token, chat_ids=chat_ids) - else: - summary["skipped"].append("telegram_missing_env") - if result.action in {BriefingAction.GITHUB_ISSUE, BriefingAction.TELEGRAM}: - github_findings = [f for f in result.findings if f.level == BriefingAction.GITHUB_ISSUE] + try: + optimization_findings = _strategy_monitoring_findings(result) + if optimization_findings: + summary["optimization_watch"] = dispatch_strategy_watch_findings( + optimization_findings, + dry_run=dry_run, + comment_existing=False, + ) + if int(summary["optimization_watch"].get("errors") or 0): + summary["errors"].append("optimization_record_failed") + except (OSError, RuntimeError, ValueError, subprocess.CalledProcessError) as exc: + summary["optimization_watch"] = { + "status": "error", + "errors": 1, + "error_type": type(exc).__name__, + } + summary["errors"].append("optimization_record_failed") + github_findings = [ + finding + for finding in result.findings + if finding.level == BriefingAction.GITHUB_ISSUE + and not (finding.kind == "strategy_monitoring" and finding.strategy_profile) + ] if github_findings: title = f"[briefing] {result.day} — {len(github_findings)} review-level alert(s)" - body = _format_github_body(result) + body = _format_github_body(result, github_findings) if dry_run: summary["github_dry_run"] = {"title": title, "body": body} else: summary["github_issue"] = create_github_issue(title=title, body=body) + if not summary["github_issue"]: + summary["errors"].append("github_issue_record_failed") + + record_failed = any( + error in {"optimization_record_failed", "github_issue_record_failed"} + for error in summary["errors"] + ) + if result.action == BriefingAction.TELEGRAM or record_failed: + if result.action == BriefingAction.TELEGRAM: + text = _format_telegram_body(result) + else: + text = ( + f"🚨 量化哨兵 operational ({result.day})\n\n" + "• optimization-record delivery failure; manual review required" + ) + if record_failed and result.action == BriefingAction.TELEGRAM: + text += "\n• optimization-record delivery failure; manual review required" + if dry_run: + summary["telegram_dry_run"] = text + else: + token = _telegram_token() + chat_ids = _telegram_chat_ids() + if token and chat_ids: + sent = send_telegram_alert(text=text, token=token, chat_ids=chat_ids) + summary["telegram_sent"] = sent + summary["operational_fallback_sent"] = bool(record_failed and sent) + if not sent: + summary["errors"].append("telegram_delivery_failed") + else: + summary["skipped"].append("telegram_missing_env") + summary["errors"].append("telegram_missing_env") return summary diff --git a/service/strategy_watch.py b/service/strategy_watch.py index c320080c..90616004 100644 --- a/service/strategy_watch.py +++ b/service/strategy_watch.py @@ -19,6 +19,15 @@ METRICS_KIND_PERFORMANCE = "performance" METRICS_KIND_OPERATIONAL = "operational_quality" REQUIRED_PERFORMANCE_METRICS = ("sharpe", "cagr", "calmar", "win_rate", "max_dd") +MONITORING_SCHEMA_VERSION = "strategy_monitoring_evidence.v1" +METRICS_KIND_MONITORING = "monitoring_evidence" +MONITORING_FINDING_TYPE = "monitoring_trigger" +STRATEGY_REPOSITORY_BY_DOMAIN = { + "cn_equity": "QuantStrategyLab/CnEquityStrategies", + "hk_equity": "QuantStrategyLab/HkEquityStrategies", + "us_equity": "QuantStrategyLab/UsEquityStrategies", + "crypto": "QuantStrategyLab/CryptoStrategies", +} def _dict_payload(value: Any) -> dict[str, Any]: @@ -96,6 +105,41 @@ def to_dict(self) -> dict[str, Any]: } +def build_strategy_monitoring_finding( + *, + domain: str, + profile: str, + severity: str, + metrics: dict[str, Any], + signals: list[dict[str, Any]], + source: str, + generated_at: str = "", + repo: str = "", +) -> StrategyWatchFinding: + """Build a pre-classified, issue-only finding from trusted monitor evidence.""" + normalized_domain = str(domain or "").strip() + normalized_profile = str(profile or "").strip() + resolved_repo = str(repo or STRATEGY_REPOSITORY_BY_DOMAIN.get(normalized_domain) or "").strip() + if not normalized_domain or not normalized_profile: + raise ValueError("strategy monitoring finding requires domain and profile") + if not resolved_repo: + raise ValueError(f"no strategy repository is configured for domain={normalized_domain!r}") + return StrategyWatchFinding( + snapshot=StrategyWatchSnapshot( + repo=resolved_repo, + profile=normalized_profile, + schema_version=MONITORING_SCHEMA_VERSION, + metrics_kind=METRICS_KIND_MONITORING, + current_metrics=dict(metrics), + source=str(source or "").strip(), + generated_at=str(generated_at or "").strip(), + ), + severity="high" if str(severity).strip().lower() == "high" else "medium", + signals=[dict(signal) for signal in signals], + finding_type=MONITORING_FINDING_TYPE, + ) + + def _snapshots_from_payload(payload: dict[str, Any]) -> list[StrategyWatchSnapshot]: default_repo = str(payload.get("repo") or payload.get("repository") or "").strip() default_schema_version = str(payload.get("schema_version") or "").strip() @@ -262,9 +306,24 @@ def finding_to_automation_task(finding: StrategyWatchFinding) -> AutomationTask: event_key = finding_event_key(finding) signal_reasons = [str(signal.get("reason") or signal.get("metric") or "metric degraded") for signal in finding.signals] finding_type = str(finding.finding_type or "metric_degradation") + if finding_type == "data_quality": + trigger_kind = "strategy_metrics_contract_invalid" + evidence_summary = "Strategy metrics payload failed watcher contract validation." + rationale = ( + "Open a data-quality issue so the source repo publishes " + "strategy_performance.v2 before optimization automation runs again." + ) + elif finding_type == MONITORING_FINDING_TYPE: + trigger_kind = "strategy_monitoring_trigger" + evidence_summary = "Strategy monitoring evidence crossed a research-review threshold." + rationale = "Open a research optimization issue for AI diagnosis and bounded, no-order experiment planning." + else: + trigger_kind = "strategy_metric_degradation" + evidence_summary = "Deterministic strategy metrics crossed degradation thresholds." + rationale = "Open a research optimization issue for AI diagnosis and sandbox experiment planning." trigger = TriggerRecord( source="strategy_optimization_watcher", - kind="strategy_metrics_contract_invalid" if finding_type == "data_quality" else "strategy_metric_degradation", + kind=trigger_kind, severity=finding.severity, reason="; ".join(signal_reasons) or ("strategy metrics contract invalid" if finding_type == "data_quality" else "strategy metrics degraded"), subject=finding.snapshot.subject(), @@ -272,11 +331,7 @@ def finding_to_automation_task(finding: StrategyWatchFinding) -> AutomationTask: evidence=signal_reasons, ) evidence = EvidenceBundle( - summary=( - "Strategy metrics payload failed watcher contract validation." - if finding_type == "data_quality" - else "Deterministic strategy metrics crossed degradation thresholds." - ), + summary=evidence_summary, artifacts=[finding.snapshot.source] if finding.snapshot.source else [], metrics={ "current": finding.snapshot.current_metrics, @@ -291,11 +346,7 @@ def finding_to_automation_task(finding: StrategyWatchFinding) -> AutomationTask: action=ISSUE_ONLY_ACTION, lane=lane, target=finding.snapshot.repo, - rationale=( - "Open a data-quality issue so the source repo publishes strategy_performance.v2 before optimization automation runs again." - if finding_type == "data_quality" - else "Open a research optimization issue for AI diagnosis and sandbox experiment planning." - ), + rationale=rationale, requires_human_review=True, metadata={"profile": finding.snapshot.profile, "plugin": finding.snapshot.plugin, "event_key": event_key, "finding_type": finding_type}, ) diff --git a/tests/test_briefing_dispatch.py b/tests/test_briefing_dispatch.py index 53a42d4c..6dd4f925 100644 --- a/tests/test_briefing_dispatch.py +++ b/tests/test_briefing_dispatch.py @@ -4,7 +4,12 @@ import unittest from unittest.mock import patch -from service.briefing_consumer import BriefingAction, BriefingConsumptionResult, BriefingFinding +from service.briefing_consumer import ( + BriefingAction, + BriefingConsumptionResult, + BriefingFinding, + consume_briefing_report, +) from service.briefing_dispatch import create_github_issue, dispatch_briefing_result, send_telegram_alert @@ -33,6 +38,77 @@ def test_dispatch_telegram_dry_run(self) -> None: self.assertIn("telegram_dry_run", summary) self.assertIn("demo", summary["telegram_dry_run"]) + @patch("service.briefing_dispatch.dispatch_strategy_watch_findings") + def test_strategy_health_dispatches_to_issue_only_watcher(self, dispatch_findings) -> None: + dispatch_findings.return_value = { + "status": "ok", + "findings": 1, + "issues": [{"repo": "QuantStrategyLab/UsEquityStrategies", "created": True}], + "errors": 0, + } + findings = consume_briefing_report( + { + "domain": "us_equity", + "strategies": [ + { + "strategy_profile": "global_etf_rotation", + "status": "critical", + "overall_score": 14.2, + "performance_score": 0.0, + } + ], + } + ) + result = BriefingConsumptionResult(day="2026-07-30", report_dir="/tmp", findings=findings) + + summary = dispatch_briefing_result(result) + + self.assertEqual(summary["action"], "github_issue") + self.assertFalse(summary["telegram_sent"]) + self.assertEqual(summary["optimization_watch"]["findings"], 1) + dispatched = dispatch_findings.call_args.args[0] + self.assertEqual(dispatched[0].snapshot.repo, "QuantStrategyLab/UsEquityStrategies") + self.assertEqual(dispatched[0].finding_type, "monitoring_trigger") + self.assertEqual(summary["errors"], []) + + @patch("service.briefing_dispatch.send_telegram_alert", return_value=True) + @patch("service.briefing_dispatch.dispatch_strategy_watch_findings") + def test_strategy_record_failure_falls_back_to_operational_telegram( + self, + dispatch_findings, + send_telegram, + ) -> None: + dispatch_findings.return_value = { + "status": "partial_error", + "findings": 1, + "issues": [{"error": "record failed"}], + "errors": 1, + } + findings = consume_briefing_report( + { + "domain": "crypto", + "strategies": [ + { + "strategy_profile": "crypto_live_pool_rotation", + "status": "critical", + "overall_score": 27.7, + } + ], + } + ) + result = BriefingConsumptionResult(day="2026-07-30", report_dir="/tmp", findings=findings) + + with patch.dict( + os.environ, + {"TELEGRAM_TOKEN": "token", "GLOBAL_TELEGRAM_CHAT_ID": "123"}, + clear=True, + ): + summary = dispatch_briefing_result(result) + + self.assertIn("optimization_record_failed", summary["errors"]) + self.assertTrue(summary["operational_fallback_sent"]) + self.assertIn("optimization-record delivery failure", send_telegram.call_args.kwargs["text"]) + @patch("service.briefing_dispatch.urllib.request.urlopen") def test_send_telegram_alert_success(self, mock_urlopen) -> None: class _Resp: diff --git a/tests/test_briefing_model_router.py b/tests/test_briefing_model_router.py index 272cb309..5cbf5b0f 100644 --- a/tests/test_briefing_model_router.py +++ b/tests/test_briefing_model_router.py @@ -87,7 +87,7 @@ def test_github_issue_for_review_status(self) -> None: self.assertEqual(len(findings), 1) self.assertEqual(findings[0].level, BriefingAction.GITHUB_ISSUE) - def test_telegram_for_critical_drift(self) -> None: + def test_optimization_issue_for_critical_drift(self) -> None: findings = consume_briefing_report( { "strategies": [ @@ -95,6 +95,21 @@ def test_telegram_for_critical_drift(self) -> None: ], } ) + self.assertEqual(findings[0].level, BriefingAction.GITHUB_ISSUE) + + def test_telegram_for_circuit_breaker(self) -> None: + findings = consume_briefing_report( + { + "strategies": [ + { + "strategy_profile": "demo", + "status": "critical", + "overall_score": 20, + "risk_flags": ["circuit_breaker"], + }, + ], + } + ) self.assertEqual(findings[0].level, BriefingAction.TELEGRAM) def test_telegram_when_briefing_data_is_unavailable(self) -> None: @@ -127,7 +142,7 @@ def test_consume_briefing_dir_reads_files(self) -> None: encoding="utf-8", ) result = consume_briefing_dir(report_dir) - self.assertEqual(result.action, BriefingAction.TELEGRAM) + self.assertEqual(result.action, BriefingAction.GITHUB_ISSUE) self.assertEqual(len(result.findings), 1) diff --git a/tests/test_consume_daily_briefing.py b/tests/test_consume_daily_briefing.py new file mode 100644 index 00000000..533406af --- /dev/null +++ b/tests/test_consume_daily_briefing.py @@ -0,0 +1,69 @@ +from __future__ import annotations + +import json +from pathlib import Path +import tempfile +from unittest.mock import patch + +from scripts.consume_daily_briefing import main + + +def _write_critical_report(report_dir: Path) -> None: + (report_dir / "us_equity.json").write_text( + json.dumps( + { + "domain": "us_equity", + "ok": True, + "strategies": [ + { + "strategy_profile": "global_etf_rotation", + "status": "critical", + "overall_score": 14.2, + } + ], + } + ), + encoding="utf-8", + ) + + +def test_successful_optimization_record_exits_zero() -> None: + with tempfile.TemporaryDirectory() as tmp: + report_dir = Path(tmp) + _write_critical_report(report_dir) + dispatch_summary = { + "action": "github_issue", + "optimization_watch": {"status": "ok", "errors": 0}, + "errors": [], + } + + with patch( + "scripts.consume_daily_briefing.dispatch_briefing_result", + return_value=dispatch_summary, + ): + assert main(["--report-dir", str(report_dir), "--dispatch"]) == 0 + + +def test_optimization_record_failure_exits_nonzero() -> None: + with tempfile.TemporaryDirectory() as tmp: + report_dir = Path(tmp) + _write_critical_report(report_dir) + dispatch_summary = { + "action": "github_issue", + "optimization_watch": {"status": "partial_error", "errors": 1}, + "errors": ["optimization_record_failed"], + } + + with patch( + "scripts.consume_daily_briefing.dispatch_briefing_result", + return_value=dispatch_summary, + ): + assert main(["--report-dir", str(report_dir), "--dispatch"]) == 2 + + +def test_undispatched_optimization_finding_exits_nonzero() -> None: + with tempfile.TemporaryDirectory() as tmp: + report_dir = Path(tmp) + _write_critical_report(report_dir) + + assert main(["--report-dir", str(report_dir)]) == 2 diff --git a/tests/test_run_strategy_optimization_watcher.py b/tests/test_run_strategy_optimization_watcher.py index 50aa576a..9146d775 100644 --- a/tests/test_run_strategy_optimization_watcher.py +++ b/tests/test_run_strategy_optimization_watcher.py @@ -6,7 +6,14 @@ import unittest from unittest.mock import patch -from scripts.run_strategy_optimization_watcher import list_open_issue_urls, parse_bool, resolve_input_path, run_watcher +from scripts.run_strategy_optimization_watcher import ( + dispatch_strategy_watch_findings, + list_open_issue_urls, + parse_bool, + resolve_input_path, + run_watcher, +) +from service.strategy_watch import build_strategy_monitoring_finding, finding_to_automation_task, watcher_issue_key def _performance_payload(*, repo: str = "QuantStrategyLab/TestStrategies", profile: str = "live", sharpe: float = 0.5) -> dict[str, object]: @@ -21,6 +28,32 @@ def _performance_payload(*, repo: str = "QuantStrategyLab/TestStrategies", profi class RunStrategyOptimizationWatcherTest(unittest.TestCase): + def test_monitoring_dispatch_does_not_repeat_existing_issue(self) -> None: + finding = build_strategy_monitoring_finding( + domain="crypto", + profile="crypto_live_pool_rotation", + severity="high", + metrics={"overall_score": 27.7}, + signals=[{"metric": "overall_score", "reason": "overall_score=27.7"}], + source="quant-monitor/health-cycle", + ) + issue_key = watcher_issue_key(finding_to_automation_task(finding)) + comment_calls: list[tuple[str, str, str]] = [] + + result = dispatch_strategy_watch_findings( + [finding], + dry_run=False, + comment_existing=False, + create_issue=lambda repo, title, body: "https://example.test/new", + comment_issue=lambda repo, url, body: comment_calls.append((repo, url, body)) or "", + list_issues=lambda repo: {issue_key: "https://example.test/existing"}, + ) + + self.assertEqual(result["errors"], 0) + self.assertEqual(result["issues"][0]["existing_url"], "https://example.test/existing") + self.assertEqual(result["issues"][0]["skipped_reason"], "open issue already records this strategy") + self.assertEqual(comment_calls, []) + def test_dry_run_does_not_create_issue(self) -> None: calls: list[tuple[str, str, str]] = [] diff --git a/tests/test_strategy_watch.py b/tests/test_strategy_watch.py index 69aa7566..68b94778 100644 --- a/tests/test_strategy_watch.py +++ b/tests/test_strategy_watch.py @@ -2,10 +2,42 @@ import unittest -from service.strategy_watch import evaluate_strategy_watch, finding_to_automation_task, issue_for_task, watcher_issue_key +from service.strategy_watch import ( + build_strategy_monitoring_finding, + evaluate_strategy_watch, + finding_to_automation_task, + issue_for_task, + watcher_issue_key, +) class StrategyWatchTest(unittest.TestCase): + def test_monitoring_trigger_becomes_issue_only_optimization_record(self) -> None: + finding = build_strategy_monitoring_finding( + domain="us_equity", + profile="global_etf_rotation", + severity="high", + metrics={"overall_score": 14.2, "performance_score": 0.0}, + signals=[ + { + "metric": "overall_score", + "reason": "overall_score=14.2 is below monitoring threshold 60.0", + } + ], + source="quant-monitor/daily-briefing", + generated_at="2026-07-31T00:00:00Z", + ) + + task = finding_to_automation_task(finding) + payload = task.to_dict() + + self.assertEqual(finding.snapshot.repo, "QuantStrategyLab/UsEquityStrategies") + self.assertEqual(finding.finding_type, "monitoring_trigger") + self.assertEqual(payload["trigger"]["kind"], "strategy_monitoring_trigger") + self.assertEqual(payload["proposed_action"]["action"], "open_issue") + self.assertFalse(payload["gate_decision"]["metadata"]["live_impact_allowed"]) + self.assertIn("bounded, no-order", payload["proposed_action"]["rationale"]) + def test_degraded_snapshot_becomes_issue_only_task(self) -> None: findings = evaluate_strategy_watch( {