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

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion .github/workflows/execution-report-heartbeat.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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' }}
Expand Down
237 changes: 211 additions & 26 deletions scripts/execution_report_heartbeat.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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")
Expand All @@ -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)

Expand Down Expand Up @@ -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")
Expand Down Expand Up @@ -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")
Expand Down Expand Up @@ -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

Expand Down Expand Up @@ -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(
Expand All @@ -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:
Expand All @@ -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])
Expand Down
Loading
Loading