From 13a6a962201ee1e2a56d20b5a3ca1325edfd836a Mon Sep 17 00:00:00 2001 From: Pigbibi <20649888+Pigbibi@users.noreply.github.com> Date: Wed, 7 Oct 2026 22:19:39 +0800 Subject: [PATCH 1/3] =?UTF-8?q?=E5=85=B1=E4=BA=AB=E8=BF=90=E8=A1=8C?= =?UTF-8?q?=E5=BF=83=E8=B7=B3=E8=A7=84=E5=88=99=E4=B8=8E=E8=A1=8C=E4=B8=BA?= =?UTF-8?q?=E6=B5=8B=E8=AF=95?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-Authored-By: Codex --- .../common/runtime_heartbeat_policy.py | 688 ++++++++++++++++++ tests/test_runtime_heartbeat_policy.py | 436 +++++++++++ 2 files changed, 1124 insertions(+) create mode 100644 src/quant_platform_kit/common/runtime_heartbeat_policy.py create mode 100644 tests/test_runtime_heartbeat_policy.py diff --git a/src/quant_platform_kit/common/runtime_heartbeat_policy.py b/src/quant_platform_kit/common/runtime_heartbeat_policy.py new file mode 100644 index 0000000..d2c1703 --- /dev/null +++ b/src/quant_platform_kit/common/runtime_heartbeat_policy.py @@ -0,0 +1,688 @@ +"""Per-runtime-target schedule and market-session policy for heartbeat checks.""" + +from __future__ import annotations + +import datetime as dt +import json +import sys +from collections.abc import Callable, Mapping +from typing import Any +from zoneinfo import ZoneInfo + + +SessionDatesLoader = Callable[..., set[dt.date]] +WarningLogger = Callable[[str], None] +ProfileResolver = Callable[[str], str] + +_LATEST_DUE_AT_KEY = "_heartbeat_latest_due_at" +_MARKET_DEFAULTS = { + "US": ("NYSE", "America/New_York"), + "HK": ("XHKG", "Asia/Hong_Kong"), + "CN": ("SSE", "Asia/Shanghai"), + "SG": ("XSES", "Asia/Singapore"), + "CRYPTO": ("24/7", "UTC"), +} +_TIMEZONE_MARKETS = { + "America/New_York": "US", + "Asia/Hong_Kong": "HK", + "Asia/Shanghai": "CN", + "Asia/Singapore": "SG", +} +_PLATFORM_MARKETS = { + "schwab": "US", + "firstrade": "US", + "qmt": "CN", + "binance": "CRYPTO", +} + + +def _split_values(raw: str | None) -> list[str]: + if not raw: + return [] + values = str(raw).replace(";", ",").replace("\n", ",").split(",") + return [value.strip() for value in values if value.strip()] + + +def _enabled(value: Any, *, default: bool = True) -> bool: + if value is None or not str(value).strip(): + return default + return str(value).strip().lower() not in {"0", "false", "no", "n", "off"} + + +def runtime_target_permits_standard_execution(runtime_target: Mapping[str, Any]) -> bool: + """Return whether the generic execution heartbeat should require a receipt.""" + + continuity = _mapping(runtime_target.get("live_continuity")) + state = str(continuity.get("state") or "").strip().upper() + return not state or state in {"ACTIVE_LKG", "ROLLBACK_LKG"} + + +def _mapping(value: Any) -> Mapping[str, Any]: + return value if isinstance(value, Mapping) else {} + + +def _target_field( + item: Mapping[str, Any], + defaults: Mapping[str, Any], + *names: str, +) -> Any: + for source in ( + item, + _mapping(item.get("env")), + defaults, + _mapping(defaults.get("env")), + ): + for name in names: + if name in source: + return source[name] + return None + + +def _runtime_target( + item: Mapping[str, Any], + defaults: Mapping[str, Any], +) -> dict[str, Any]: + value = _target_field(item, defaults, "runtime_target", "runtime_target_json") + if isinstance(value, str): + try: + value = json.loads(value) + except json.JSONDecodeError as exc: + raise ValueError(f"runtime target JSON is invalid: {exc}") from exc + if isinstance(value, Mapping): + return dict(value) + return dict(item) + + +def _first_value(sources: list[Mapping[str, Any]], keys: tuple[str, ...]) -> str: + for source in sources: + for key in keys: + value = source.get(key) + if value is not None and str(value).strip(): + return str(value).strip() + return "" + + +def _suffix_value(sources: list[Mapping[str, Any]], suffix: str) -> str: + for source in sources: + for key, value in source.items(): + if str(key).upper().endswith(suffix) and value is not None and str(value).strip(): + return str(value).strip() + return "" + + +def _service_values( + item: Mapping[str, Any], + runtime_target: Mapping[str, Any], + defaults: Mapping[str, Any], + environ: Mapping[str, str], +) -> list[str]: + value = _target_field( + item, + defaults, + "service", + "service_name", + "cloud_run_service", + ) + if value is None: + value = _first_value( + [runtime_target], + ("service", "service_name", "cloud_run_service"), + ) + if value: + return _split_values(str(value)) + services = _split_values(environ.get("CLOUD_RUN_SERVICES")) + services.extend(_split_values(environ.get("CLOUD_RUN_SERVICE"))) + return list(dict.fromkeys(services)) + + +def _normalize_target( + item: Mapping[str, Any], + runtime_target: Mapping[str, Any], + defaults: Mapping[str, Any], + service: str, + environ: Mapping[str, str], + *, + use_global_market_fallback: bool, + profile_resolver: ProfileResolver | None, +) -> dict[str, Any]: + sources = [ + runtime_target, + item, + _mapping(item.get("env")), + defaults, + _mapping(defaults.get("env")), + ] + scheduler = next( + ( + value + for source in sources + if isinstance((value := source.get("scheduler")), dict) + ), + {}, + ) + scheduler = dict(scheduler) + if not str(scheduler.get("main_time") or "").strip(): + main_time = _target_field( + item, + defaults, + "CLOUD_SCHEDULER_MAIN_TIME", + "cloud_scheduler_main_time", + ) + if main_time is None: + main_time = environ.get("CLOUD_SCHEDULER_MAIN_TIME") + if main_time is not None and str(main_time).strip(): + scheduler["main_time"] = str(main_time).strip() + account_scope = _first_value( + sources, + ( + "account_scope", + "account_group", + "account_region", + "ACCOUNT_GROUP", + "ACCOUNT_REGION", + ), + ) + target_market_timezone = ( + _first_value(sources, ("market_timezone", "MARKET_TIMEZONE")) + or _suffix_value(sources, "_MARKET_TIMEZONE") + ) + global_market_timezone = ( + str(environ.get("RUNTIME_HEARTBEAT_MARKET_TIMEZONE") or "").strip() + or _suffix_value([environ], "_MARKET_TIMEZONE") + ) + market_timezone = target_market_timezone or ( + global_market_timezone if use_global_market_fallback else "" + ) + market = ( + _first_value(sources, ("market", "MARKET")) + or _suffix_value(sources, "_MARKET") + ).upper() + if not market: + market = _TIMEZONE_MARKETS.get(target_market_timezone, "") + if market not in _MARKET_DEFAULTS: + market = _TIMEZONE_MARKETS.get(str(scheduler.get("timezone") or "").strip(), "") + if not market: + platform_id = _first_value(sources, ("platform_id",)).lower() + market = _PLATFORM_MARKETS.get(platform_id, "") + if not market and account_scope.upper() in {"US", "HK", "CN"}: + market = account_scope.upper() + if market not in _MARKET_DEFAULTS and use_global_market_fallback: + market = str(environ.get("RUNTIME_HEARTBEAT_MARKET") or "").strip().upper() + target_market_calendar = ( + _first_value(sources, ("market_calendar", "MARKET_CALENDAR")) + or _suffix_value(sources, "_MARKET_CALENDAR") + ) + global_market_calendar = ( + str(environ.get("RUNTIME_HEARTBEAT_MARKET_CALENDAR") or "").strip() + or _suffix_value([environ], "_MARKET_CALENDAR") + ) + default_calendar, default_timezone = _MARKET_DEFAULTS.get(market, ("", "")) + market_calendar = ( + target_market_calendar + or (global_market_calendar if use_global_market_fallback else "") + or default_calendar + ) + market_timezone = ( + market_timezone + or default_timezone + or str(scheduler.get("timezone") or "").strip() + ) + if scheduler and not str(scheduler.get("timezone") or "").strip() and market_timezone: + scheduler["timezone"] = market_timezone + strategy_profile = _first_value( + sources, + ("strategy_profile", "strategy", "profile"), + ) + if strategy_profile and profile_resolver is not None: + strategy_profile = profile_resolver(strategy_profile) + return { + "service": str(service).strip(), + "strategy_profile": strategy_profile, + "account_scope": account_scope, + "scheduler": scheduler, + "market": market, + "market_calendar": market_calendar, + "market_timezone": market_timezone, + } + + +def _load_runtime_target_items( + environ: Mapping[str, str], +) -> tuple[list[Mapping[str, Any]], Mapping[str, Any]]: + raw_targets = str(environ.get("CLOUD_RUN_SERVICE_TARGETS_JSON") or "").strip() + items: list[Mapping[str, Any]] = [] + defaults: Mapping[str, Any] = {} + if raw_targets: + try: + payload = json.loads(raw_targets) + except json.JSONDecodeError as exc: + raise ValueError(f"CLOUD_RUN_SERVICE_TARGETS_JSON is invalid: {exc}") from exc + if isinstance(payload, Mapping): + targets = payload.get("targets") + defaults = _mapping(payload.get("defaults")) + else: + targets = payload + if isinstance(targets, list): + items = [target for target in targets if isinstance(target, Mapping)] + else: + raise ValueError( + "CLOUD_RUN_SERVICE_TARGETS_JSON must be an array or object with targets" + ) + + if not items: + raw_runtime_target = str(environ.get("RUNTIME_TARGET_JSON") or "").strip() + if raw_runtime_target: + try: + runtime_target = json.loads(raw_runtime_target) + except json.JSONDecodeError as exc: + raise ValueError(f"RUNTIME_TARGET_JSON is invalid: {exc}") from exc + if isinstance(runtime_target, Mapping): + items = [runtime_target] + else: + raise ValueError("RUNTIME_TARGET_JSON must decode to an object") + return items, defaults + + +def runtime_target_configuration_present(environ: Mapping[str, str]) -> bool: + return bool( + str(environ.get("CLOUD_RUN_SERVICE_TARGETS_JSON") or "").strip() + or str(environ.get("RUNTIME_TARGET_JSON") or "").strip() + ) + + +def runtime_target_configuration_has_enabled_targets( + environ: Mapping[str, str], +) -> bool: + items, defaults = _load_runtime_target_items(environ) + if not items: + return False + expected_scope = str(environ.get("RUNTIME_HEARTBEAT_ACCOUNT_SCOPE") or "").strip().lower() + for item in items: + runtime_target = _runtime_target(item, defaults) + target_scope = _first_value( + [ + runtime_target, + item, + _mapping(item.get("env")), + defaults, + _mapping(defaults.get("env")), + ], + ( + "account_scope", + "account_group", + "account_region", + "ACCOUNT_GROUP", + "ACCOUNT_REGION", + ), + ) + if expected_scope and target_scope and target_scope.lower() != expected_scope: + continue + enabled_value = _target_field( + item, + defaults, + "runtime_target_enabled", + "RUNTIME_TARGET_ENABLED", + ) + if enabled_value is None: + enabled_value = runtime_target.get("runtime_target_enabled") + if _enabled(enabled_value) and runtime_target_permits_standard_execution(runtime_target): + return True + return False + + +def load_runtime_targets( + environ: Mapping[str, str], + *, + include_disabled: bool = False, + profile_resolver: ProfileResolver | None = None, +) -> list[dict[str, Any]]: + items, defaults = _load_runtime_target_items(environ) + + expected_scope = str(environ.get("RUNTIME_HEARTBEAT_ACCOUNT_SCOPE") or "").strip().lower() + eligible: list[tuple[Mapping[str, Any], dict[str, Any], bool]] = [] + for item in items: + runtime_target = _runtime_target(item, defaults) + enabled_value = _target_field( + item, + defaults, + "runtime_target_enabled", + "RUNTIME_TARGET_ENABLED", + ) + if enabled_value is None: + enabled_value = runtime_target.get("runtime_target_enabled") + enabled = ( + _enabled(enabled_value) + and runtime_target_permits_standard_execution(runtime_target) + ) + if not enabled and not include_disabled: + continue + target_scope = _first_value( + [ + runtime_target, + item, + _mapping(item.get("env")), + defaults, + _mapping(defaults.get("env")), + ], + ( + "account_scope", + "account_group", + "account_region", + "ACCOUNT_GROUP", + "ACCOUNT_REGION", + ), + ) + if expected_scope and target_scope and target_scope.lower() != expected_scope: + continue + eligible.append((item, runtime_target, enabled)) + + normalized: list[dict[str, Any]] = [] + enabled_count = sum(1 for _item, _runtime_target_value, enabled in eligible if enabled) + for item, runtime_target, enabled in eligible: + for service in _service_values(item, runtime_target, defaults, environ): + target = _normalize_target( + item, + runtime_target, + defaults, + service, + environ, + use_global_market_fallback=enabled_count == 1, + profile_resolver=profile_resolver, + ) + if include_disabled: + target["enabled"] = enabled + key = target_key(target) + if key and all(target_key(existing) != key for existing in normalized): + normalized.append(target) + return normalized + + +def target_key(target: Mapping[str, Any]) -> str: + service = str(target.get("service") or "").strip().lower() + strategy = str(target.get("strategy_profile") or "").strip().lower() + scope = str(target.get("account_scope") or "").strip().lower() + return f"{service}|{strategy or '*'}|{scope or '*'}" if service else "" + + +def target_label(target: Mapping[str, Any]) -> str: + service = str(target.get("service") or "").strip() or "" + strategy = str(target.get("strategy_profile") or "").strip() + scope = str(target.get("account_scope") or "").strip() + qualifiers = "/".join(value for value in (strategy, scope) if value) + return f"{service}[{qualifiers}]" if qualifiers else service + + +def target_latest_due_at(target: Mapping[str, Any]) -> dt.datetime | None: + value = target.get(_LATEST_DUE_AT_KEY) + return value if isinstance(value, dt.datetime) else None + + +def _payload_value(payload: Mapping[str, Any], keys: tuple[str, ...]) -> str: + runtime_target = payload.get("runtime_target") + sources = [payload] + if isinstance(runtime_target, Mapping): + sources.append(runtime_target) + return _first_value(sources, keys) + + +def match_payload_target( + payload: Mapping[str, Any], + targets: list[dict[str, Any]], +) -> tuple[str | None, str]: + service = _payload_value( + payload, + ("service_name", "service", "cloud_run_service"), + ).lower() + strategy = _payload_value( + payload, + ("strategy_profile", "strategy", "profile"), + ).lower() + scope = _payload_value( + payload, + ("account_scope", "account_group", "account_region"), + ).lower() + for target in targets: + expected_service = str(target.get("service") or "").strip().lower() + expected_strategy = str(target.get("strategy_profile") or "").strip().lower() + expected_scope = str(target.get("account_scope") or "").strip().lower() + if service != expected_service: + continue + if expected_strategy and strategy != expected_strategy: + continue + if expected_scope and scope != expected_scope: + continue + return target_key(target), "matched runtime target" + return None, ( + f"runtime_target={service or '-'}/{strategy or '-'}/{scope or '-'}" + ) + + +def _cron_token_value(token: str, *, names: dict[str, int] | None = None) -> int: + normalized = token.strip().lower() + if names and normalized in names: + return names[normalized] + return int(normalized) + + +def _cron_field_values( + field: str, + *, + minimum: int, + maximum: int, + names: dict[str, int] | None = None, +) -> set[int] | None: + text = str(field or "").strip().lower() + if text in {"", "*"}: + return None + values: set[int] = set() + for raw_part in text.split(","): + part = raw_part.strip() + if not part: + continue + base, raw_step = part, "1" + if "/" in part: + base, raw_step = part.split("/", 1) + step = max(1, int(raw_step)) + if base == "*": + start, end = minimum, maximum + elif "-" in base: + raw_start, raw_end = base.split("-", 1) + start = _cron_token_value(raw_start, names=names) + end = _cron_token_value(raw_end, names=names) + else: + start = end = _cron_token_value(base, names=names) + for value in range(start, end + 1, step): + if minimum <= value <= maximum: + values.add(value) + elif maximum == 6 and value == 7: + values.add(0) + return values + + +def cron_matches(schedule: str, value: dt.datetime) -> bool: + fields = str(schedule or "").split() + if len(fields) == 2: + fields.extend(("*", "*", "*")) + if len(fields) != 5: + return False + minute, hour, day_of_month, month, day_of_week = fields + dow_names = { + "sun": 0, + "mon": 1, + "tue": 2, + "wed": 3, + "thu": 4, + "fri": 5, + "sat": 6, + } + minute_values = _cron_field_values(minute, minimum=0, maximum=59) + hour_values = _cron_field_values(hour, minimum=0, maximum=23) + dom_values = _cron_field_values(day_of_month, minimum=1, maximum=31) + month_values = _cron_field_values(month, minimum=1, maximum=12) + dow_values = _cron_field_values(day_of_week, minimum=0, maximum=6, names=dow_names) + if minute_values is not None and value.minute not in minute_values: + return False + if hour_values is not None and value.hour not in hour_values: + return False + if month_values is not None and value.month not in month_values: + return False + dom_matches = dom_values is None or value.day in dom_values + dow_matches = dow_values is None or value.isoweekday() % 7 in dow_values + if dom_values is not None and dow_values is not None: + return dom_matches or dow_matches + return dom_matches and dow_matches + + +def _market_session_dates( + calendar: str, + *, + start_date: dt.date, + end_date: dt.date, +) -> set[dt.date]: + import pandas_market_calendars as mcal + + schedule = mcal.get_calendar(calendar).schedule( + start_date=start_date, + end_date=end_date, + ) + return {value.date() for value in schedule.index} + + +def _target_due_status( + target: Mapping[str, Any], + *, + since: dt.datetime, + now: dt.datetime, + market_aware: bool, + session_dates_loader: SessionDatesLoader, + warning_logger: WarningLogger, + publication_grace: dt.timedelta, +) -> tuple[bool | None, dt.datetime | None]: + scheduler = target.get("scheduler") + if not isinstance(scheduler, Mapping): + return None, None + schedule = str(scheduler.get("main_time") or "").strip() + fields = schedule.split() + if len(fields) != 5: + return None, None + timezone_name = str(scheduler.get("timezone") or "UTC").strip() or "UTC" + try: + scheduler_timezone = ZoneInfo(timezone_name) + except Exception as exc: # noqa: BLE001 + warning_logger( + f"Unable to evaluate heartbeat scheduler timezone {timezone_name}: " + f"{type(exc).__name__}; keeping target required" + ) + return None, None + + since_utc = since.astimezone(dt.timezone.utc) + now_utc = now.astimezone(dt.timezone.utc) + cursor = since_utc.replace(second=0, microsecond=0) + if cursor < since_utc: + cursor += dt.timedelta(minutes=1) + matured_at = now_utc - max(publication_grace, dt.timedelta()) + cron_due_at: list[dt.datetime] = [] + while cursor <= now_utc: + local_time = cursor.astimezone(scheduler_timezone) + try: + matches = cron_matches(schedule, local_time) + except (TypeError, ValueError) as exc: + warning_logger( + f"Unable to evaluate heartbeat cron for {target_label(target)}: " + f"{type(exc).__name__}; keeping target required" + ) + return None, None + if matches and cursor <= matured_at: + cron_due_at.append(cursor) + cursor += dt.timedelta(minutes=1) + latest_cron_due_at = cron_due_at[-1] if cron_due_at else None + + session_dates: set[dt.date] | None = None + market_calendar = str(target.get("market_calendar") or "").strip() + market_timezone_name = ( + str(target.get("market_timezone") or "").strip() or timezone_name + ) + if market_aware and market_calendar: + try: + market_timezone = ZoneInfo(market_timezone_name) + session_dates = session_dates_loader( + market_calendar, + start_date=since_utc.astimezone(market_timezone).date(), + end_date=now_utc.astimezone(market_timezone).date(), + ) + except Exception as exc: # noqa: BLE001 + warning_logger( + f"Unable to evaluate heartbeat market calendar {market_calendar}: " + f"{type(exc).__name__}; keeping target required" + ) + return None, latest_cron_due_at + + latest_due_at = latest_cron_due_at + if session_dates is not None: + market_timezone = ZoneInfo(market_timezone_name) + latest_due_at = next( + ( + due_at + for due_at in reversed(cron_due_at) + if due_at.astimezone(market_timezone).date() in session_dates + ), + None, + ) + return latest_due_at is not None, latest_due_at + + +def filter_due_targets( + targets: list[dict[str, Any]], + *, + since: dt.datetime, + now: dt.datetime, + market_aware: bool = True, + publication_grace: dt.timedelta = dt.timedelta(minutes=30), + session_dates_loader: SessionDatesLoader = _market_session_dates, + warning_logger: WarningLogger = lambda message: print(message, file=sys.stderr), +) -> tuple[list[dict[str, Any]], bool]: + due: list[dict[str, Any]] = [] + evaluated = False + for target in targets: + status, latest_due_at = _target_due_status( + target, + since=since, + now=now, + market_aware=market_aware, + session_dates_loader=session_dates_loader, + warning_logger=warning_logger, + publication_grace=publication_grace, + ) + if status is not None: + evaluated = True + if status is not False: + due_target = dict(target) + if latest_due_at is not None: + due_target[_LATEST_DUE_AT_KEY] = latest_due_at + due.append(due_target) + return due, evaluated + + +def filter_services_for_targets( + services: list[str], + targets: list[dict[str, Any]], + *, + all_targets: list[dict[str, Any]] | None = None, +) -> list[str]: + if not targets: + return services + target_services = { + str(target.get("service") or "").strip() + for target in targets + if str(target.get("service") or "").strip() + } + configured_services = { + str(target.get("service") or "").strip() + for target in (all_targets or targets) + if str(target.get("service") or "").strip() + } + return [ + service + for service in services + if service not in configured_services or service in target_services + ] diff --git a/tests/test_runtime_heartbeat_policy.py b/tests/test_runtime_heartbeat_policy.py new file mode 100644 index 0000000..90b0956 --- /dev/null +++ b/tests/test_runtime_heartbeat_policy.py @@ -0,0 +1,436 @@ +from __future__ import annotations + +import datetime as dt +import json + +from quant_platform_kit.common.runtime_heartbeat_policy import filter_due_targets, load_runtime_targets, match_payload_target, runtime_target_configuration_present, target_key, target_latest_due_at + +def _target( + *, + service: str, + strategy: str, + scope: str, + timezone: str, + calendar: str, +) -> dict[str, object]: + return { + "service": service, + "runtime_target": { + "service_name": service, + "strategy_profile": strategy, + "account_scope": scope, + "scheduler": { + "timezone": timezone, + "main_time": "45 15 * * *", + }, + "market_calendar": calendar, + "market_timezone": timezone, + }, + } + +def test_due_targets_use_each_strategy_market_calendar() -> None: + targets = load_runtime_targets( + { + "CLOUD_RUN_SERVICE_TARGETS_JSON": json.dumps( + { + "targets": [ + _target( + service="svc-us", + strategy="us-strategy", + scope="US", + timezone="America/New_York", + calendar="NYSE", + ), + _target( + service="svc-hk", + strategy="hk-strategy", + scope="HK", + timezone="Asia/Hong_Kong", + calendar="XHKG", + ), + ] + } + ) + } + ) + + due, evaluated = filter_due_targets( + targets, + since=dt.datetime(2026, 7, 3, 0, 0, tzinfo=dt.timezone.utc), + now=dt.datetime(2026, 7, 3, 22, 0, tzinfo=dt.timezone.utc), + session_dates_loader=lambda calendar, **_kwargs: ( + {dt.date(2026, 7, 3)} if calendar == "XHKG" else set() + ), + ) + + assert evaluated is True + assert [target["strategy_profile"] for target in due] == ["hk-strategy"] + +def test_real_exchange_calendars_distinguish_us_holiday_from_hk_session() -> None: + targets = load_runtime_targets( + { + "CLOUD_RUN_SERVICE_TARGETS_JSON": json.dumps( + { + "targets": [ + _target( + service="svc-us", + strategy="us-strategy", + scope="US", + timezone="America/New_York", + calendar="NYSE", + ), + _target( + service="svc-hk", + strategy="hk-strategy", + scope="HK", + timezone="Asia/Hong_Kong", + calendar="XHKG", + ), + ] + } + ) + } + ) + + due, evaluated = filter_due_targets( + targets, + since=dt.datetime(2026, 7, 3, 0, 0, tzinfo=dt.timezone.utc), + now=dt.datetime(2026, 7, 3, 22, 0, tzinfo=dt.timezone.utc), + ) + + assert evaluated is True + assert [target["strategy_profile"] for target in due] == ["hk-strategy"] + +def test_july_29_us_month_end_target_is_due_at_1545_eastern() -> None: + raw_target = _target( + service="svc-us-monthly", + strategy="us-monthly", + scope="US", + timezone="America/New_York", + calendar="NYSE", + ) + raw_target["runtime_target"]["scheduler"]["main_time"] = "45 15 25-29 * *" + targets = load_runtime_targets( + { + "CLOUD_RUN_SERVICE_TARGETS_JSON": json.dumps( + {"targets": [raw_target]} + ) + } + ) + + due, evaluated = filter_due_targets( + targets, + since=dt.datetime(2026, 7, 29, 19, 40, tzinfo=dt.timezone.utc), + now=dt.datetime(2026, 7, 29, 20, 20, tzinfo=dt.timezone.utc), + ) + + assert evaluated is True + assert [target["strategy_profile"] for target in due] == ["us-monthly"] + assert target_latest_due_at(due[0]) == dt.datetime( + 2026, + 7, + 29, + 19, + 45, + tzinfo=dt.timezone.utc, + ) + +def test_neutral_daily_heartbeat_tracks_latest_due_time_per_market() -> None: + targets = load_runtime_targets( + { + "CLOUD_RUN_SERVICE_TARGETS_JSON": json.dumps( + { + "targets": [ + _target( + service="svc-us", + strategy="us-strategy", + scope="US", + timezone="America/New_York", + calendar="NYSE", + ), + _target( + service="svc-hk", + strategy="hk-strategy", + scope="HK", + timezone="Asia/Hong_Kong", + calendar="XHKG", + ), + ] + } + ) + } + ) + + due, evaluated = filter_due_targets( + targets, + since=dt.datetime(2026, 7, 28, 10, 20, tzinfo=dt.timezone.utc), + now=dt.datetime(2026, 7, 29, 22, 20, tzinfo=dt.timezone.utc), + session_dates_loader=lambda _calendar, **_kwargs: { + dt.date(2026, 7, 28), + dt.date(2026, 7, 29), + }, + ) + + assert evaluated is True + assert { + target["strategy_profile"]: target_latest_due_at(target) + for target in due + } == { + "us-strategy": dt.datetime( + 2026, + 7, + 29, + 19, + 45, + tzinfo=dt.timezone.utc, + ), + "hk-strategy": dt.datetime( + 2026, + 7, + 29, + 7, + 45, + tzinfo=dt.timezone.utc, + ), + } + +def test_scheduler_timezone_beats_account_region_when_market_is_not_explicit() -> None: + targets = load_runtime_targets( + { + "CLOUD_RUN_SERVICE_TARGETS_JSON": json.dumps( + { + "targets": [ + { + "service": "longbridge-sg-us-service", + "account_scope": "SG", + "runtime_target": { + "service_name": "longbridge-sg-us-service", + "strategy_profile": "us-strategy", + "account_scope": "SG", + "scheduler": { + "timezone": "America/New_York", + "main_time": "45 15 * * *", + }, + }, + } + ] + } + ) + } + ) + + assert targets[0]["market"] == "US" + assert targets[0]["market_calendar"] == "NYSE" + assert targets[0]["market_timezone"] == "America/New_York" + +def test_ambiguous_sg_account_does_not_guess_a_stock_exchange() -> None: + targets = load_runtime_targets( + { + "CLOUD_RUN_SERVICE_TARGETS_JSON": json.dumps( + { + "targets": [ + { + "service": "longbridge-sg-service", + "account_scope": "SG", + "runtime_target": { + "service_name": "longbridge-sg-service", + "strategy_profile": "unknown-strategy", + "account_scope": "SG", + "scheduler": { + "timezone": "UTC", + "main_time": "0 12 * * *", + }, + }, + } + ] + } + ) + } + ) + + assert targets[0]["market"] == "" + assert targets[0]["market_calendar"] == "" + +def test_calendar_failure_keeps_target_due_fail_closed() -> None: + targets = load_runtime_targets( + { + "RUNTIME_TARGET_JSON": json.dumps( + _target( + service="svc-us", + strategy="us-strategy", + scope="US", + timezone="America/New_York", + calendar="INVALID", + )["runtime_target"] + ) + } + ) + + def fail_calendar(_calendar: str, **_kwargs: object) -> set[dt.date]: + raise RuntimeError("calendar unavailable") + + due, evaluated = filter_due_targets( + targets, + since=dt.datetime(2026, 7, 3, 0, 0, tzinfo=dt.timezone.utc), + now=dt.datetime(2026, 7, 3, 22, 0, tzinfo=dt.timezone.utc), + session_dates_loader=fail_calendar, + warning_logger=lambda _message: None, + ) + + assert len(due) == 1 + assert target_latest_due_at(due[0]) == dt.datetime( + 2026, + 7, + 3, + 19, + 45, + tzinfo=dt.timezone.utc, + ) + assert evaluated is False + +def test_target_defaults_and_scheduler_aliases_are_normalized() -> None: + targets = load_runtime_targets( + { + "CLOUD_RUN_SERVICE_TARGETS_JSON": json.dumps( + { + "defaults": { + "env": { + "RUNTIME_TARGET_ENABLED": "false", + "CLOUD_SCHEDULER_MAIN_TIME": "45 15 25-29 * *", + }, + "market": "US", + }, + "targets": [ + { + "service": "disabled-service", + "runtime_target": { + "strategy_profile": "disabled-strategy", + }, + }, + { + "service": "enabled-service", + "RUNTIME_TARGET_ENABLED": "true", + "runtime_target": { + "strategy_profile": "enabled-strategy", + }, + }, + ], + } + ) + } + ) + + assert len(targets) == 1 + assert targets[0]["service"] == "enabled-service" + assert targets[0]["market"] == "US" + assert targets[0]["scheduler"] == { + "main_time": "45 15 25-29 * *", + "timezone": "America/New_York", + } + +def test_reconcile_only_target_is_not_an_execution_heartbeat_target() -> None: + environ = { + "CLOUD_RUN_SERVICE_TARGETS_JSON": json.dumps( + { + "targets": [ + { + "service": "reconcile-only-service", + "runtime_target": { + "service_name": "reconcile-only-service", + "strategy_profile": "strategy-a", + "live_continuity": {"state": "RECONCILE_ONLY"}, + }, + } + ] + } + ) + } + + assert load_runtime_targets(environ) == [] + assert runtime_target_configuration_present(environ) is True + +def test_strategy_profile_resolver_canonicalizes_aliases() -> None: + targets = load_runtime_targets( + { + "RUNTIME_TARGET_JSON": json.dumps( + { + "service_name": "alias-service", + "strategy_profile": "supported-alias", + "scheduler": { + "main_time": "45 15 * * *", + "timezone": "UTC", + }, + } + ) + }, + profile_resolver=lambda value: ( + "canonical-strategy" if value == "supported-alias" else value + ), + ) + + assert targets[0]["strategy_profile"] == "canonical-strategy" + +def test_publication_grace_uses_previous_matured_schedule_cutoff() -> None: + targets = load_runtime_targets( + { + "RUNTIME_TARGET_JSON": json.dumps( + { + "service_name": "grace-service", + "strategy_profile": "grace-strategy", + "scheduler": { + "main_time": "0 12 * * *", + "timezone": "UTC", + }, + } + ) + } + ) + since = dt.datetime(2026, 7, 28, 11, 30, tzinfo=dt.timezone.utc) + + within_grace, evaluated = filter_due_targets( + targets, + since=since, + now=dt.datetime(2026, 7, 29, 12, 5, tzinfo=dt.timezone.utc), + market_aware=False, + publication_grace=dt.timedelta(minutes=30), + ) + after_grace, _ = filter_due_targets( + targets, + since=since, + now=dt.datetime(2026, 7, 29, 12, 31, tzinfo=dt.timezone.utc), + market_aware=False, + publication_grace=dt.timedelta(minutes=30), + ) + + assert evaluated is True + assert target_latest_due_at(within_grace[0]) == dt.datetime( + 2026, + 7, + 28, + 12, + 0, + tzinfo=dt.timezone.utc, + ) + assert target_latest_due_at(after_grace[0]) == dt.datetime( + 2026, + 7, + 29, + 12, + 0, + tzinfo=dt.timezone.utc, + ) + +def test_runtime_target_configuration_presence_is_preserved_when_all_disabled() -> None: + environ = { + "CLOUD_RUN_SERVICE_TARGETS_JSON": json.dumps( + { + "defaults": {"runtime_target_enabled": False}, + "targets": [{"service": "disabled-service"}], + } + ) + } + + assert runtime_target_configuration_present(environ) is True + assert load_runtime_targets(environ) == [] + assert load_runtime_targets(environ, include_disabled=True)[0]["enabled"] is False + From f65d71633d9b1ae131456d33faec42e83d56ba9b Mon Sep 17 00:00:00 2001 From: Pigbibi <20649888+Pigbibi@users.noreply.github.com> Date: Wed, 7 Oct 2026 22:21:53 +0800 Subject: [PATCH 2/3] test: keep the optional market-calendar backend optional Co-Authored-By: Codex --- tests/test_runtime_heartbeat_policy.py | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/tests/test_runtime_heartbeat_policy.py b/tests/test_runtime_heartbeat_policy.py index 90b0956..fa2631d 100644 --- a/tests/test_runtime_heartbeat_policy.py +++ b/tests/test_runtime_heartbeat_policy.py @@ -2,6 +2,7 @@ import datetime as dt import json +import pytest from quant_platform_kit.common.runtime_heartbeat_policy import filter_due_targets, load_runtime_targets, match_payload_target, runtime_target_configuration_present, target_key, target_latest_due_at @@ -67,6 +68,7 @@ def test_due_targets_use_each_strategy_market_calendar() -> None: assert [target["strategy_profile"] for target in due] == ["hk-strategy"] def test_real_exchange_calendars_distinguish_us_holiday_from_hk_session() -> None: + pytest.importorskip("pandas_market_calendars") targets = load_runtime_targets( { "CLOUD_RUN_SERVICE_TARGETS_JSON": json.dumps( @@ -433,4 +435,3 @@ def test_runtime_target_configuration_presence_is_preserved_when_all_disabled() assert runtime_target_configuration_present(environ) is True assert load_runtime_targets(environ) == [] assert load_runtime_targets(environ, include_disabled=True)[0]["enabled"] is False - From 9aa0fa8fd96dc769624423728c9d8e9848076653 Mon Sep 17 00:00:00 2001 From: Pigbibi <20649888+Pigbibi@users.noreply.github.com> Date: Wed, 7 Oct 2026 22:31:48 +0800 Subject: [PATCH 3/3] test: inject session dates for the shared schedule unit test Co-Authored-By: Codex --- tests/test_runtime_heartbeat_policy.py | 1 + 1 file changed, 1 insertion(+) diff --git a/tests/test_runtime_heartbeat_policy.py b/tests/test_runtime_heartbeat_policy.py index fa2631d..384bb20 100644 --- a/tests/test_runtime_heartbeat_policy.py +++ b/tests/test_runtime_heartbeat_policy.py @@ -124,6 +124,7 @@ def test_july_29_us_month_end_target_is_due_at_1545_eastern() -> None: targets, since=dt.datetime(2026, 7, 29, 19, 40, tzinfo=dt.timezone.utc), now=dt.datetime(2026, 7, 29, 20, 20, tzinfo=dt.timezone.utc), + session_dates_loader=lambda _calendar, **_kwargs: {dt.date(2026, 7, 29)}, ) assert evaluated is True