diff --git a/.github/workflows/execution-report-heartbeat.yml b/.github/workflows/execution-report-heartbeat.yml index d294167..5f2d560 100644 --- a/.github/workflows/execution-report-heartbeat.yml +++ b/.github/workflows/execution-report-heartbeat.yml @@ -44,7 +44,7 @@ jobs: RUNTIME_HEARTBEAT_ACCEPT_STAGES: ${{ vars.RUNTIME_HEARTBEAT_ACCEPT_STAGES }} RUNTIME_HEARTBEAT_REJECT_STAGES: ${{ vars.RUNTIME_HEARTBEAT_REJECT_STAGES }} RUNTIME_HEARTBEAT_MARKET_AWARE: ${{ vars.RUNTIME_HEARTBEAT_MARKET_AWARE || 'true' }} - RUNTIME_HEARTBEAT_NOTIFY_ON_SUCCESS: ${{ vars.RUNTIME_HEARTBEAT_NOTIFY_ON_SUCCESS || 'false' }} + RUNTIME_HEARTBEAT_NOTIFY_ON_SUCCESS: ${{ github.event_name != 'schedule' && vars.RUNTIME_HEARTBEAT_NOTIFY_ON_SUCCESS || 'false' }} RUNTIME_HEARTBEAT_MARKET_CALENDAR: ${{ vars.FIRSTRADE_MARKET_CALENDAR }} RUNTIME_HEARTBEAT_MARKET_TIMEZONE: ${{ vars.FIRSTRADE_MARKET_TIMEZONE }} RUNTIME_HEARTBEAT_PUBLICATION_GRACE_MINUTES: ${{ vars.RUNTIME_HEARTBEAT_PUBLICATION_GRACE_MINUTES || '30' }} diff --git a/scripts/execution_report_heartbeat.py b/scripts/execution_report_heartbeat.py index af9f9ec..f0743a4 100644 --- a/scripts/execution_report_heartbeat.py +++ b/scripts/execution_report_heartbeat.py @@ -14,6 +14,7 @@ from typing import Any from zoneinfo import ZoneInfo, ZoneInfoNotFoundError +from quant_platform_kit.common.execution_receipts import validate_execution_receipt from quant_platform_kit.common.operational_notification_localization import ( format_operational_alert, format_operational_heartbeat_status, @@ -184,7 +185,7 @@ def _heartbeat_skip_reason_for_schedule(since: dt.datetime, now: dt.datetime) -> ) -def _parse_timestamp(value: Any) -> dt.datetime | None: +def _parse_timestamp(value: Any, *, require_timezone: bool = False) -> dt.datetime | None: if not value: return None text = str(value).strip().replace("Z", "+00:00") @@ -193,6 +194,8 @@ def _parse_timestamp(value: Any) -> dt.datetime | None: except ValueError: return None if parsed.tzinfo is None: + if require_timezone: + return None parsed = parsed.replace(tzinfo=dt.timezone.utc) return parsed.astimezone(dt.timezone.utc) @@ -454,25 +457,31 @@ def _object_updated_at(entry: dict[str, Any]) -> dt.datetime | None: def _report_errors(payload: dict[str, Any]) -> list[Any]: - errors = payload.get("errors") - if isinstance(errors, list) and errors: - return errors - error_summary = payload.get("error_summary") - if isinstance(error_summary, dict): - nested = error_summary.get("errors") - if isinstance(nested, list) and nested: - return nested - if payload.get("error"): - return [payload.get("error")] - return [] + errors = [] + for scope in (payload, payload.get("summary"), payload.get("diagnostics")): + if not isinstance(scope, dict): + continue + for key in ("errors", "error", "plugin_error", "persistence_error", "report_persistence_error", "strategy_plugin_error", "strategy_plugin_alert_error", "strategy_run_persistence_error", "market_hours_check_error"): + value = scope.get(key) + if value: + errors.extend(value if isinstance(value, list) else [value]) + error_summary = scope.get("error_summary") + if isinstance(error_summary, dict) and error_summary.get("errors"): + value = error_summary["errors"] + errors.extend(value if isinstance(value, list) else [value]) + return errors + def _report_status(payload: dict[str, Any]) -> tuple[str, str]: - status = str(payload.get("status") or payload.get("summary", {}).get("status") or "").strip() - stage = str(payload.get("stage") or payload.get("summary", {}).get("stage") or "").strip() + summary = payload.get("summary") + summary = summary if isinstance(summary, dict) else {} + status = str(payload.get("status") or summary.get("status") or "").strip() + stage = str(payload.get("stage") or summary.get("stage") or "").strip() return status, stage + def _report_notification_failure(payload: dict[str, Any]) -> str: scopes = [payload] summary = payload.get("summary") @@ -589,6 +598,127 @@ def _is_accepted_report(payload: dict[str, Any]) -> tuple[bool, str]: +def _report_execution_status(payload: dict[str, Any]) -> str: + summary = payload.get("summary") + nested_status = summary.get("execution_status") if isinstance(summary, dict) else None + return str(payload.get("execution_status") or nested_status or "").strip() + + +def _market_closed_by_calendar(payload: dict[str, Any]) -> bool: + """Recheck closure across the run; a producer's fail-closed fallback is not proof.""" + try: + import pandas_market_calendars as mcal + timezone = ZoneInfo(str(payload.get("market_timezone") or "")) + started = _parse_timestamp(payload.get("started_at"), require_timezone=True) + finished = _parse_timestamp(payload.get("finished_at"), require_timezone=True) + if started is None or finished is None or finished < started: + return False + calendar = mcal.get_calendar(str(payload.get("market_calendar") or "")) + if str(calendar.tz) != str(timezone): + return False + schedule = calendar.schedule( + start_date=started.astimezone(timezone).date(), + end_date=finished.astimezone(timezone).date(), + ) + if len(schedule.index) == 0: + return True + for _day, session in schedule.iterrows(): + opened = _parse_timestamp(session.get("market_open"), require_timezone=True) + closed = _parse_timestamp(session.get("market_close"), require_timezone=True) + if opened is None or closed is None or closed <= opened: + return False + if started < closed and finished >= opened: + return False + return True + except Exception: # noqa: BLE001 - unreadable calendar must remain visible + return False + + +def _is_quiet_report(payload: dict[str, Any], *, now: dt.datetime | None = None) -> bool: + """Quiet a complete healthy receipt without duplicating the business notification.""" + if payload.get("errors") != [] or not isinstance(payload.get("errors"), list): + return False + if _report_errors(payload) or _report_notification_failure(payload): + return False + started = _parse_timestamp(payload.get("started_at"), require_timezone=True) + finished = _parse_timestamp(payload.get("finished_at"), require_timezone=True) + if started is None or finished is None or finished < started or (now is not None and finished > now + dt.timedelta(minutes=5)): + return False + summary = payload.get("summary") + diagnostics = payload.get("diagnostics") + if not isinstance(summary, dict) or not isinstance(diagnostics, dict): + return False + for scope in (payload, summary, diagnostics): + scope_status = str(scope.get("status") or "").strip().lower() + scope_stage = str(scope.get("stage") or "").strip().upper() + if scope_status and scope_status not in {"ok", "success", "completed", "no_action", "skipped"}: + return False + if scope_stage == "FUNDING_BLOCKED" or (scope_stage and scope_stage not in DEFAULT_ACCEPT_STAGES): + return False + if any(scope.get(key) for key in ( + "pending_reconciliation", "reconciliation_required", "execution_blocked", + "plugin_error", "persistence_error", "report_persistence_error", "notification_error", + )): + return False + for key in ("order_events_count", "orders_pending_count", "orders_filled_count", "orders_partially_filled_count"): + if key in scope and (type(scope[key]) is not int or scope[key] < 0): + return False + execution_status = str(scope.get("execution_status") or "").strip().lower() + if execution_status and execution_status not in {"no_op", "no_action", "completed", "dry_run_completed", "executed", "submitted", "broker_acknowledged", "partially_filled", "pending"}: + return False + if "execution_receipt" in payload: + try: + receipt = validate_execution_receipt(payload["execution_receipt"]) + except (ValueError, TypeError): + return False + if ( + receipt["outcome"] in {"reconciliation_required", "failed", "risk_blocked"} + or receipt["strategy_profile"] != payload.get("strategy_profile") + or receipt["platform"] != {"interactive_brokers": "ibkr", "charles_schwab": "schwab"}.get(payload.get("platform"), payload.get("platform")) + ): + return False + status, _stage = _report_status(payload) + if status.lower() == "skipped": + runtime_target = payload.get("runtime_target") + closure_payload = dict(payload) + for key in ("market", "market_calendar", "market_timezone"): + nested = runtime_target.get(key) if isinstance(runtime_target, dict) else None + if payload.get(key) and nested and payload[key] != nested: + return False + closure_payload[key] = payload.get(key) or nested + if diagnostics.get("skip_reason") != "market_closed": + return False + for scope in (payload, summary, diagnostics): + if any(scope.get(key) for key in ( + "orders_submitted", "submitted_orders", "orders_pending", "orders_filled", + "orders_partially_filled", "option_orders_submitted", "option_orders_pending", + "option_orders_filled", "option_orders_partially_filled", + )): + return False + if any(scope.get(key, 0) != 0 for key in ( + "order_events_count", "orders_pending_count", "orders_filled_count", "orders_partially_filled_count", + )): + return False + if str(scope.get("execution_status") or "").lower() not in {"", "no_op", "no_action"}: + return False + if "execution_receipt" in payload and receipt["outcome"] not in {"not_due", "no_action", "no_signal", "no_rebalance"}: + return False + if not closure_payload.get("market") or not closure_payload.get("market_calendar"): + return False + try: + ZoneInfo(str(closure_payload.get("market_timezone") or "")) + except (ValueError, ZoneInfoNotFoundError): + return False + return _market_closed_by_calendar(closure_payload) + return ( + status.lower() in {"ok", "success", "completed", "no_action"} + and str(summary.get("execution_status") or "").lower() in {"no_op", "no_action", "completed", "dry_run_completed", "executed", "submitted", "broker_acknowledged", "partially_filled", "pending"} + and type(summary.get("order_events_count")) is int + and summary["order_events_count"] >= 0 + ) + + + def _telegram_secret_project() -> str | None: return ( os.environ.get("RUNTIME_HEARTBEAT_GCP_PROJECT_ID") @@ -651,12 +781,19 @@ def _send_telegram(message: str) -> bool: return ok -def _send_normal_heartbeat(name: str, detail: str) -> int: - if not _env_bool("RUNTIME_HEARTBEAT_NOTIFY_ON_SUCCESS", False): +def _send_normal_heartbeat(name: str, detail: str, *, quiet_eligible: bool = True) -> int: + if quiet_eligible and not _env_bool("RUNTIME_HEARTBEAT_NOTIFY_ON_SUCCESS", False): return 0 - message = format_operational_heartbeat_status( - locale=_notification_locale(), name=name, detail=detail, - ) + if quiet_eligible: + message = format_operational_heartbeat_status( + locale=_notification_locale(), name=name, detail=detail, + ) + else: + message = format_operational_alert( + locale=_notification_locale(), alert_type="execution_report_heartbeat", name=name, + issues=["运行检查证据待核对" if _notification_locale() == "zh" else "Runtime check evidence needs review"], + technical_details=[detail], + ) print(message) return 0 if _send_telegram(message[:3900]) else 1 @@ -790,9 +927,18 @@ def main(now: dt.datetime | None = None) -> int: accepted = [] accepted_by_service: dict[str, tuple[str, dt.datetime, str]] = {} inspected = [] + quiet_by_service = {} + quiet_reports = [] + visible_report_details = [] + read_incomplete = False + latest_report_seen = set() + latest_services_seen = set() + unreadable_objects = [] + latest_report_issues = [] for uri, updated in sorted_objects[:max_reports]: payload = _cat_gcs_json(uri, project=project) if payload is None: + unreadable_objects.append(updated) inspected.append(f"- {updated.isoformat()} {uri} unreadable") continue matches, service_name, filter_reason = _payload_matches( @@ -801,6 +947,11 @@ def main(now: dt.datetime | None = None) -> int: required_targets=required_targets, ) if not matches: + raw_service = _payload_service_name(payload) + known_services = set(required_services) | target_services + if raw_service in known_services and raw_service not in latest_services_seen: + latest_services_seen.add(raw_service) + latest_report_issues.append(f"{raw_service}: conflicting report identity ({filter_reason})") inspected.append(f"- {updated.isoformat()} {uri} skipped {filter_reason}") continue if updated < due_since: @@ -814,36 +965,70 @@ def main(now: dt.datetime | None = None) -> int: ) continue ok, reason = _is_accepted_report(payload) + latest_services_seen.add(_payload_service_name(payload)) + is_latest_report = service_name not in latest_report_seen + if is_latest_report: + latest_report_seen.add(service_name) + if not ok: + latest_report_issues.append(f"{service_name or name}: {reason}") inspected.append(f"- {updated.isoformat()} {uri} {reason}") - if ok: + if ok and is_latest_report: + quiet = _is_quiet_report(payload, now=now) + freshness_detail = "" + if quiet and latest_due_at is not None: + finished_at = _parse_timestamp(payload.get("finished_at"), require_timezone=True) + if finished_at is None or finished_at < latest_due_at: + quiet = False + freshness_detail = "; payload finished_at predates latest due schedule" + if not quiet: + visible_report_details.append( + f"{service_name or name}: execution_status={_report_execution_status(payload) or 'unconfirmed'}; {reason}{freshness_detail}" + ) if required_keys: - accepted_by_service[service_name] = (uri, updated, reason) + quiet_by_service.setdefault(service_name, quiet) + accepted_by_service.setdefault(service_name, (uri, updated, reason)) else: accepted.append((uri, updated, reason)) + quiet_reports.append(quiet) + + read_incomplete = bool(unreadable_objects) + if ( + required_keys + and all(key in accepted_by_service and quiet_by_service.get(key) for key in required_keys) + and not latest_report_issues + ): + latest_complete_floor = min(accepted_by_service[key][1] for key in required_keys) + read_incomplete = any(updated >= latest_complete_floor for updated in unreadable_objects) if required_keys: missing = [key for key in required_keys if key not in accepted_by_service] - if not missing and not list_errors: + if not missing and not list_errors and not read_incomplete and not latest_report_issues: details = ", ".join( f"{required_labels[key]}@{accepted_by_service[key][1].isoformat()}" for key in required_keys ) + if visible_report_details: + details += "; " + "; ".join(visible_report_details) print(f"Execution report heartbeat OK for {name}: {details}") return _send_normal_heartbeat( - name, _notice("heartbeat_accepted_report", detail=details), + name, _notice("heartbeat_accepted_report", detail=details) if all(quiet_by_service[key] for key in required_keys) else details, + quiet_eligible=all(quiet_by_service[key] for key in required_keys), ) - if accepted and not list_errors: + if accepted and not list_errors and not read_incomplete and not latest_report_issues: uri, updated, reason = accepted[0] print( f"Execution report heartbeat OK for {name}: {reason}, updated={updated.isoformat()}, uri={uri}" ) return _send_normal_heartbeat( name, - _notice("heartbeat_accepted_report", detail=f"{reason}, {updated.isoformat()}"), + _notice("heartbeat_accepted_report", detail=f"{reason}, {updated.isoformat()}") if all(quiet_reports) else "; ".join(visible_report_details), + quiet_eligible=all(quiet_reports), ) issues = [] - technical_details = [] + technical_details = list(latest_report_issues) + if read_incomplete or latest_report_issues: + issues.append(_notice("heartbeat_no_acceptable_report", count=min(len(sorted_objects), max_reports))) if list_errors: issues.append(_notice("heartbeat_list_failed")) technical_details.extend(list_errors[:3]) diff --git a/tests/test_notification_quiet.py b/tests/test_notification_quiet.py new file mode 100644 index 0000000..4e73355 --- /dev/null +++ b/tests/test_notification_quiet.py @@ -0,0 +1,338 @@ +"""Offline receipt-only notification regression matrix; no trading mutations.""" +from __future__ import annotations +import copy +import datetime as dt +import socket +import subprocess +import urllib.request +from types import SimpleNamespace +import pytest +from scripts import execution_report_heartbeat as heartbeat +NOW = dt.datetime(2026, 10, 5, 13, 0, tzinfo=dt.timezone.utc) +TARGET = SimpleNamespace(service='mock-service', profile='mock_profile', scope='mock_scope') +@pytest.fixture(autouse=True) +def offline(monkeypatch): + def denied(*args, **kwargs): + raise AssertionError('No external operations in this suite') + monkeypatch.setattr(socket.socket, 'connect', denied) + monkeypatch.setattr(socket, 'create_connection', denied) + monkeypatch.setattr(subprocess, 'Popen', denied) + monkeypatch.setattr(urllib.request, 'urlopen', denied) + for key in list(__import__('os').environ): + if key.startswith(('RUNTIME_HEARTBEAT_', 'CLOUD_RUN_')): + monkeypatch.delenv(key, raising=False) +def report(): + return {'service_name': TARGET.service, 'strategy_profile': TARGET.profile, 'account_scope': TARGET.scope, + 'started_at': '2026-10-05T12:00:00Z', 'finished_at': '2026-10-05T12:01:00Z', + 'status': 'ok', 'dry_run': False, 'errors': [], + 'summary': {'execution_status': 'no_action', 'order_events_count': 0}, + 'diagnostics': {}, 'market': 'US', 'market_calendar': 'NYSE', 'market_timezone': 'America/New_York'} + +def mock_heartbeat(monkeypatch, payload): + monkeypatch.setattr(heartbeat, '_runtime_target_enabled', lambda: True) + monkeypatch.setattr(heartbeat, 'runtime_target_configuration_present', lambda env: False) + monkeypatch.setattr(heartbeat, 'load_runtime_targets', lambda env: []) + monkeypatch.setattr(heartbeat, '_hydrate_runtime_target_schedules', lambda targets, **kw: targets) + monkeypatch.setattr(heartbeat, 'filter_due_targets', lambda *args, **kw: ([], False)) + monkeypatch.setattr(heartbeat, '_heartbeat_skip_reason_for_schedule', lambda *args: None) + monkeypatch.setattr(heartbeat, '_load_required_services', lambda: []) + monkeypatch.setattr(heartbeat, '_report_globs', lambda *args: ['gs://mock/reports']) + monkeypatch.setattr(heartbeat, '_list_gcs_objects', lambda *args, **kw: [{'url': 'gs://mock/report.json', 'metadata': {'updated': NOW.isoformat()}}]) + monkeypatch.setattr(heartbeat, '_cat_gcs_json', lambda *args, **kw: payload) + sent = [] + monkeypatch.setattr(heartbeat, '_send_telegram', lambda msg, **kw: sent.append(msg) or True) + return sent + + +@pytest.mark.parametrize('flag', [None, 'false', 'true']) +def test_heartbeat_explicit_normal_preserves_manual_flag(monkeypatch, capsys, flag): + if flag is not None: + monkeypatch.setenv('RUNTIME_HEARTBEAT_NOTIFY_ON_SUCCESS', flag) + sent = mock_heartbeat(monkeypatch, report()) + assert heartbeat.main(NOW) == 0 + assert bool(sent) is (flag == 'true') + assert 'heartbeat OK' in capsys.readouterr().out + + +@pytest.mark.parametrize('execution_status', ['unknown', 'partial', 'pending_reconciliation', 'reconciliation_required']) +def test_accepted_business_or_unknown_report_is_visible_with_success_flag_off(monkeypatch, execution_status): + payload = report() + payload['summary']['execution_status'] = execution_status + sent = mock_heartbeat(monkeypatch, payload) + assert heartbeat.main(NOW) == 0 + assert len(sent) == 1 + assert payload['summary']['execution_status'] == execution_status + + +@pytest.mark.parametrize('change', [{'errors': None}, {'finished_at': None}, {'summary': {}}, {'summary': {'execution_status': 'no_op', 'order_events_count': '0'}}]) +def test_accepted_incomplete_report_is_visible(monkeypatch, change): + payload = report() + payload.update(change) + sent = mock_heartbeat(monkeypatch, payload) + assert heartbeat.main(NOW) == 0 + assert sent + + +def test_heartbeat_nested_error_does_not_hide_in_top_level_ok(monkeypatch): + payload = report() + payload['summary']['errors'] = ['plugin failed'] + sent = mock_heartbeat(monkeypatch, payload) + assert heartbeat.main(NOW) == 1 + assert sent + + + +@pytest.mark.parametrize('execution_status', ['executed', 'completed', 'submitted', 'broker_acknowledged', 'partially_filled', 'pending']) +def test_healthy_receipt_does_not_repeat_trade_notification(monkeypatch, execution_status): + payload = report() + payload['summary'].update(execution_status=execution_status, order_events_count=1) + sent = mock_heartbeat(monkeypatch, payload) + assert heartbeat.main(NOW) == 0 + assert not sent + assert payload['summary']['execution_status'] == execution_status + + +def test_newest_rejected_report_cannot_be_hidden_by_older_healthy_report(monkeypatch): + payload = report() + sent = mock_heartbeat(monkeypatch, payload) + monkeypatch.setattr(heartbeat, '_list_gcs_objects', lambda *args, **kw: [ + {'url': 'gs://mock/new.json', 'metadata': {'updated': NOW.isoformat()}}, + {'url': 'gs://mock/old.json', 'metadata': {'updated': (NOW - dt.timedelta(minutes=2)).isoformat()}}]) + def read(uri, **kw): + result = report() + if uri.endswith('new.json'): + result['summary']['errors'] = ['persistence failed'] + return result + monkeypatch.setattr(heartbeat, '_cat_gcs_json', read) + assert heartbeat.main(NOW) == 1 + assert sent + + +@pytest.mark.parametrize('calendar_result', ['holiday', 'closed', 'open', 'read_error', 'wrong_timezone']) +def test_closure_requires_calendar_evidence(monkeypatch, calendar_result): + from types import SimpleNamespace + import sys + class Calendar: + tz = 'UTC' if calendar_result == 'wrong_timezone' else 'America/New_York' + def schedule(self, **kw): + if calendar_result == 'read_error': + raise RuntimeError('calendar unavailable') + return SimpleNamespace(index=[] if calendar_result == 'holiday' else [1], iterrows=lambda: iter([(1, {'market_open': '2026-10-05T11:00:00Z' if calendar_result == 'open' else '2026-10-05T13:30:00Z', 'market_close': '2026-10-05T20:00:00Z'})])) + def open_at_time(self, schedule, when, **kw): + return calendar_result == 'open' + monkeypatch.setitem(sys.modules, 'pandas_market_calendars', SimpleNamespace(get_calendar=lambda name: Calendar())) + assert heartbeat._market_closed_by_calendar(report()) is (calendar_result in ['holiday', 'closed']) + + +def test_accepted_report_does_not_cover_incomplete_listing_or_read(monkeypatch): + sent = mock_heartbeat(monkeypatch, report()) + monkeypatch.setattr(heartbeat, '_report_globs', lambda *args: ['gs://mock/good', 'gs://mock/bad']) + def listing(glob, **kw): + if glob.endswith('bad'): + raise RuntimeError('list failed') + return [{'url': 'gs://mock/ok.json', 'metadata': {'updated': NOW.isoformat()}}] + monkeypatch.setattr(heartbeat, '_list_gcs_objects', listing) + assert heartbeat.main(NOW) == 1 + assert sent + + +def test_unscoped_mixed_services_unknown_is_visible(monkeypatch): + sent = mock_heartbeat(monkeypatch, report()) + monkeypatch.setattr(heartbeat, '_list_gcs_objects', lambda *args, **kw: [ + {'url': 'gs://mock/a.json', 'metadata': {'updated': NOW.isoformat()}}, + {'url': 'gs://mock/b.json', 'metadata': {'updated': (NOW - dt.timedelta(minutes=2)).isoformat()}}]) + def read(uri, **kw): + payload = report() + if uri.endswith('b.json'): + payload['service_name'] = 'other-service' + payload['summary']['execution_status'] = 'unknown' + return payload + monkeypatch.setattr(heartbeat, '_cat_gcs_json', read) + assert heartbeat.main(NOW) == 0 + assert len(sent) == 1 + assert 'unknown' in sent[0] + + +def test_conflicting_latest_identity_cannot_be_hidden_by_old_healthy(monkeypatch): + sent = mock_heartbeat(monkeypatch, report()) + monkeypatch.setenv('RUNTIME_HEARTBEAT_ACCOUNT_SCOPE', TARGET.scope) + monkeypatch.setattr(heartbeat, '_load_required_services', lambda: [TARGET.service]) + monkeypatch.setattr(heartbeat, '_list_gcs_objects', lambda *args, **kw: [ + {'url': 'gs://mock/new.json', 'metadata': {'updated': NOW.isoformat()}}, + {'url': 'gs://mock/old.json', 'metadata': {'updated': (NOW - dt.timedelta(minutes=2)).isoformat()}}]) + def read(uri, **kw): + payload = report() + if uri.endswith('new.json'): + payload['account_scope'] = 'wrong-scope' + return payload + monkeypatch.setattr(heartbeat, '_cat_gcs_json', read) + assert heartbeat.main(NOW) == 1 + assert len(sent) == 1 + assert 'conflicting report identity' in sent[0] + + +@pytest.mark.parametrize('value', [True, 1, 'bad']) +def test_malformed_nested_error_remains_alertable(monkeypatch, value): + payload = report() + payload['error_summary'] = {'errors': value} + sent = mock_heartbeat(monkeypatch, payload) + assert heartbeat.main(NOW) == 1 + assert sent + + +@pytest.mark.parametrize('value', [None, [], 'bad']) +def test_incomplete_summary_does_not_crash_before_visibility(monkeypatch, value): + payload = report() + payload['summary'] = value + sent = mock_heartbeat(monkeypatch, payload) + assert heartbeat.main(NOW) == 0 + assert sent + + +@pytest.mark.parametrize('field', ['strategy_plugin_error', 'strategy_plugin_alert_error', 'strategy_run_persistence_error', 'market_hours_check_error']) +def test_actual_producer_error_field_stays_visible(monkeypatch, field): + payload = report() + payload['diagnostics'][field] = 'mock failure' + sent = mock_heartbeat(monkeypatch, payload) + assert heartbeat.main(NOW) == 1 + assert sent + + +@pytest.mark.parametrize('outcome', ['reconciliation_required', 'failed', 'risk_blocked', 'filled', 'submitted']) +def test_receipt_fact_is_preserved_and_controls_health_disposition(monkeypatch, outcome): + from quant_platform_kit.common.execution_receipts import build_execution_receipt + payload = report() + payload['platform'] = 'firstrade' + payload['execution_receipt'] = build_execution_receipt(platform='firstrade', strategy_profile=TARGET.profile, strategy_revision='a' * 40, execution_mode='paper', outcome=outcome, observed_at=NOW, **({'broker_confirmation': 'not_observed'} if outcome == 'failed' else {})) + original = copy.deepcopy(payload['execution_receipt']) + sent = mock_heartbeat(monkeypatch, payload) + assert heartbeat.main(NOW) == 0 + assert bool(sent) is (outcome in ['reconciliation_required', 'failed', 'risk_blocked']) + assert payload['execution_receipt'] == original + + +@pytest.mark.parametrize('summary', [{'execution_status': 'submitted', 'order_events_count': 1}, {'execution_status': 'pending', 'order_events_count': 0}, {'orders_pending_count': 1}, {'orders_submitted': ['mock']}]) +def test_conflicting_closed_report_stays_visible(monkeypatch, summary): + payload = report() + payload.update(status='skipped', summary=summary, diagnostics={'skip_reason': 'market_closed'}) + sent = mock_heartbeat(monkeypatch, payload) + monkeypatch.setattr(heartbeat, '_market_closed_by_calendar', lambda payload: True) + assert heartbeat.main(NOW) == 0 + assert sent + + +def test_newest_healthy_report_restores_quiet_after_old_unknown(monkeypatch): + sent = mock_heartbeat(monkeypatch, report()) + monkeypatch.setattr(heartbeat, '_list_gcs_objects', lambda *args, **kw: [ + {'url': 'gs://mock/new.json', 'metadata': {'updated': NOW.isoformat()}}, + {'url': 'gs://mock/old.json', 'metadata': {'updated': (NOW - dt.timedelta(minutes=2)).isoformat()}}]) + def read(uri, **kw): + payload = report() + if uri.endswith('old.json'): + payload['summary']['execution_status'] = 'unknown' + return payload + monkeypatch.setattr(heartbeat, '_cat_gcs_json', read) + assert heartbeat.main(NOW) == 0 + assert not sent + + +def test_funding_guard_is_visible_even_with_top_level_ok(monkeypatch): + payload = report() + payload['stage'] = 'FUNDING_BLOCKED' + sent = mock_heartbeat(monkeypatch, payload) + assert heartbeat.main(NOW) == 0 + assert sent + + +def test_missing_timezone_evidence_is_visible(monkeypatch): + payload = report() + payload['started_at'] = '2026-10-05T12:00:00' + sent = mock_heartbeat(monkeypatch, payload) + assert heartbeat.main(NOW) == 0 + assert sent + + +def test_closure_crossing_open_cannot_be_quiet(monkeypatch): + import sys + from types import SimpleNamespace + class Calendar: + tz = 'America/New_York' + def schedule(self, **kw): + return SimpleNamespace(index=[1], iterrows=lambda: iter([(1, {'market_open': '2026-10-05T13:30:00Z', 'market_close': '2026-10-05T20:00:00Z'})])) + monkeypatch.setitem(sys.modules, 'pandas_market_calendars', SimpleNamespace(get_calendar=lambda name: Calendar())) + payload = report() + payload['finished_at'] = '2026-10-05T13:31:00Z' + assert not heartbeat._market_closed_by_calendar(payload) + + +def test_target_qualified_current_health_is_not_overridden_by_old_wrong_scope(monkeypatch): + sent = mock_heartbeat(monkeypatch, report()) + target = {'service': TARGET.service, 'strategy_profile': TARGET.profile, 'account_scope': TARGET.scope} + monkeypatch.setattr(heartbeat, 'load_runtime_targets', lambda env: [target]) + monkeypatch.setattr(heartbeat, 'filter_due_targets', lambda *args, **kw: ([target], True)) + monkeypatch.setattr(heartbeat, '_list_gcs_objects', lambda *args, **kw: [ + {'url': 'gs://mock/new.json', 'metadata': {'updated': NOW.isoformat()}}, + {'url': 'gs://mock/old.json', 'metadata': {'updated': (NOW - dt.timedelta(minutes=2)).isoformat()}}]) + def read(uri, **kw): + payload = report() + if uri.endswith('old.json'): + payload['account_scope'] = 'wrong-scope' + return payload + monkeypatch.setattr(heartbeat, '_cat_gcs_json', read) + assert heartbeat.main(NOW) == 0 + assert not sent + + +@pytest.mark.parametrize('unreadable_is_newer', [False, True]) +def test_only_current_window_unreadable_can_override_complete_exact_targets(monkeypatch, unreadable_is_newer): + sent = mock_heartbeat(monkeypatch, report()) + target = {'service': TARGET.service, 'strategy_profile': TARGET.profile, 'account_scope': TARGET.scope} + monkeypatch.setattr(heartbeat, 'load_runtime_targets', lambda env: [target]) + monkeypatch.setattr(heartbeat, 'filter_due_targets', lambda *args, **kw: ([target], True)) + bad_time = NOW if unreadable_is_newer else NOW - dt.timedelta(minutes=4) + monkeypatch.setattr(heartbeat, '_list_gcs_objects', lambda *args, **kw: [ + {'url': 'gs://mock/good.json', 'metadata': {'updated': (NOW - dt.timedelta(minutes=2)).isoformat()}}, + {'url': 'gs://mock/bad.json', 'metadata': {'updated': bad_time.isoformat()}}]) + monkeypatch.setattr(heartbeat, '_cat_gcs_json', lambda uri, **kw: None if uri.endswith('bad.json') else report()) + assert heartbeat.main(NOW) == (1 if unreadable_is_newer else 0) + assert bool(sent) is unreadable_is_newer + + +@pytest.mark.parametrize('kind', ['unknown', 'incomplete', 'reconciliation_required', 'persistence_error']) +def test_unconfirmed_check_message_never_claims_normal_or_trade_failure(monkeypatch, kind): + payload = report() + if kind == 'incomplete': + payload['finished_at'] = None + elif kind == 'persistence_error': + payload['diagnostics']['strategy_run_persistence_error'] = 'mock failure' + else: + payload['summary']['execution_status'] = kind + sent = mock_heartbeat(monkeypatch, payload) + assert heartbeat.main(NOW) == (1 if kind == 'persistence_error' else 0) + assert sent + assert heartbeat._notice('status_normal') not in sent[0] + assert '✅' not in sent[0] + assert '交易失败' not in sent[0] + + +def test_confirmed_manual_opt_in_uses_existing_normal_formatter(monkeypatch): + monkeypatch.setenv('RUNTIME_HEARTBEAT_NOTIFY_ON_SUCCESS', 'true') + sent = mock_heartbeat(monkeypatch, report()) + assert heartbeat.main(NOW) == 0 + assert heartbeat._notice('status_normal') in sent[0] + + +@pytest.mark.parametrize('due_at', [None, NOW - dt.timedelta(hours=2), NOW - dt.timedelta(minutes=30)]) +def test_reuploaded_old_payload_is_visible_only_against_existing_precise_due(monkeypatch, due_at): + sent = mock_heartbeat(monkeypatch, report()) + target = {'service': TARGET.service, 'strategy_profile': TARGET.profile, 'account_scope': TARGET.scope} + if due_at is not None: + target['_heartbeat_latest_due_at'] = due_at + monkeypatch.setattr(heartbeat, 'load_runtime_targets', lambda env: [target]) + monkeypatch.setattr(heartbeat, 'filter_due_targets', lambda *args, **kw: ([target], True)) + assert heartbeat.main(NOW) == 0 + stale = due_at is not None and due_at > dt.datetime(2026, 10, 5, 12, 1, tzinfo=dt.timezone.utc) + assert bool(sent) is stale + if stale: + assert heartbeat._notice('status_normal') not in sent[0] + assert 'predates latest due' in sent[0]