Skip to content

Commit 3bb4505

Browse files
committed
Skip heartbeat outside Firstrade scheduler window
1 parent 8e0a40c commit 3bb4505

3 files changed

Lines changed: 115 additions & 2 deletions

File tree

.github/workflows/execution-report-heartbeat.yml

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -42,6 +42,7 @@ jobs:
4242
RUNTIME_HEARTBEAT_FAIL_WORKFLOW_ON_ALERT: ${{ inputs.fail_workflow_on_alert || vars.RUNTIME_HEARTBEAT_FAIL_WORKFLOW_ON_ALERT || 'true' }}
4343
RUNTIME_HEARTBEAT_ACCEPT_STAGES: ${{ vars.RUNTIME_HEARTBEAT_ACCEPT_STAGES }}
4444
RUNTIME_HEARTBEAT_REJECT_STAGES: ${{ vars.RUNTIME_HEARTBEAT_REJECT_STAGES }}
45+
RUNTIME_TARGET_JSON: ${{ vars.RUNTIME_TARGET_JSON }}
4546
FIRSTRADE_GCS_STATE_BUCKET: ${{ vars.FIRSTRADE_GCS_STATE_BUCKET }}
4647
FIRSTRADE_STATE_PREFIX: ${{ vars.FIRSTRADE_STATE_PREFIX }}
4748
GLOBAL_TELEGRAM_CHAT_ID: ${{ vars.GLOBAL_TELEGRAM_CHAT_ID }}

scripts/execution_report_heartbeat.py

Lines changed: 81 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -12,6 +12,7 @@
1212
import urllib.parse
1313
import urllib.request
1414
from typing import Any
15+
from zoneinfo import ZoneInfo, ZoneInfoNotFoundError
1516

1617

1718
DEFAULT_ACCEPT_STATUSES = {"ok", "skipped", "success", "completed", "no_action"}
@@ -42,6 +43,79 @@ def _env_bool(name: str, default: bool = False) -> bool:
4243
return value in {"1", "true", "yes", "y", "on"}
4344

4445

46+
def _parse_day_of_month_field(raw: str) -> set[int] | None:
47+
text = str(raw or "").strip()
48+
if not text or text in {"*", "?"}:
49+
return None
50+
days: set[int] = set()
51+
for part in text.split(","):
52+
part = part.strip()
53+
if not part:
54+
continue
55+
step = 1
56+
if "/" in part:
57+
part, step_text = part.split("/", 1)
58+
try:
59+
step = max(1, int(step_text))
60+
except ValueError:
61+
return None
62+
if "-" in part:
63+
start_text, end_text = part.split("-", 1)
64+
try:
65+
start = int(start_text)
66+
end = int(end_text)
67+
except ValueError:
68+
return None
69+
if start > end:
70+
return None
71+
days.update(day for day in range(start, end + 1, step) if 1 <= day <= 31)
72+
continue
73+
try:
74+
day = int(part)
75+
except ValueError:
76+
return None
77+
if 1 <= day <= 31:
78+
days.add(day)
79+
return days or None
80+
81+
82+
def _runtime_target_scheduler() -> dict[str, Any]:
83+
raw = (os.environ.get("RUNTIME_TARGET_JSON") or "").strip()
84+
if not raw:
85+
return {}
86+
try:
87+
payload = json.loads(raw)
88+
except json.JSONDecodeError:
89+
return {}
90+
scheduler = payload.get("scheduler") if isinstance(payload, dict) else None
91+
return scheduler if isinstance(scheduler, dict) else {}
92+
93+
94+
def _heartbeat_skip_reason_for_schedule(now: dt.datetime) -> str | None:
95+
scheduler = _runtime_target_scheduler()
96+
cron = str(scheduler.get("main_time") or "").strip()
97+
fields = cron.split()
98+
if len(fields) != 5:
99+
return None
100+
expected_days = _parse_day_of_month_field(fields[2])
101+
if not expected_days:
102+
return None
103+
timezone_name = str(scheduler.get("timezone") or "UTC").strip() or "UTC"
104+
try:
105+
timezone = ZoneInfo(timezone_name)
106+
except ZoneInfoNotFoundError:
107+
timezone = dt.timezone.utc
108+
timezone_name = "UTC"
109+
local_now = now.astimezone(timezone)
110+
if local_now.day in expected_days:
111+
return None
112+
day_text = ",".join(str(day) for day in sorted(expected_days))
113+
return (
114+
f"runtime scheduler main_time is not due today "
115+
f"({timezone_name} day={local_now.day}; expected day(s)={day_text})"
116+
)
117+
118+
45119
def _parse_timestamp(value: Any) -> dt.datetime | None:
46120
if not value:
47121
return None
@@ -353,7 +427,7 @@ def _send_telegram(message: str) -> bool:
353427
return ok
354428

355429

356-
def main() -> int:
430+
def main(now: dt.datetime | None = None) -> int:
357431
project = (
358432
os.environ.get("RUNTIME_HEARTBEAT_GCP_PROJECT_ID")
359433
or os.environ.get("GCP_PROJECT_ID")
@@ -365,7 +439,12 @@ def main() -> int:
365439
fail_workflow = _env_bool("RUNTIME_HEARTBEAT_FAIL_WORKFLOW_ON_ALERT", True)
366440
required_services = _load_required_services()
367441

368-
now = dt.datetime.now(dt.timezone.utc)
442+
now = now or dt.datetime.now(dt.timezone.utc)
443+
schedule_skip_reason = _heartbeat_skip_reason_for_schedule(now)
444+
if schedule_skip_reason:
445+
print(f"Execution report heartbeat skipped for {name}: {schedule_skip_reason}")
446+
return 0
447+
369448
since = now - dt.timedelta(hours=lookback_hours)
370449
globs = _report_globs(since, now)
371450
if not globs:

tests/test_execution_report_heartbeat.py

Lines changed: 33 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -52,3 +52,36 @@ def fake_run_gcloud(command):
5252
"firstradequant",
5353
]
5454

55+
56+
57+
def test_heartbeat_skips_outside_runtime_target_scheduler_day(monkeypatch, capsys):
58+
monkeypatch.setenv("RUNTIME_HEARTBEAT_NAME", "Firstrade monthly runtime")
59+
monkeypatch.setenv(
60+
"RUNTIME_TARGET_JSON",
61+
'{"scheduler":{"timezone":"America/New_York","main_time":"45 15 25-28 * *"}}',
62+
)
63+
monkeypatch.setattr(
64+
heartbeat,
65+
"_list_gcs_objects",
66+
lambda *_args, **_kwargs: (_ for _ in ()).throw(AssertionError("GCS should not be queried")),
67+
)
68+
69+
result = heartbeat.main(now=dt.datetime(2026, 6, 20, 23, 10, tzinfo=dt.timezone.utc))
70+
71+
assert result == 0
72+
output = capsys.readouterr().out
73+
assert "Execution report heartbeat skipped for Firstrade monthly runtime" in output
74+
assert "expected day(s)=25,26,27,28" in output
75+
76+
77+
def test_heartbeat_does_not_skip_inside_runtime_target_scheduler_day(monkeypatch):
78+
monkeypatch.setenv(
79+
"RUNTIME_TARGET_JSON",
80+
'{"scheduler":{"timezone":"America/New_York","main_time":"45 15 25-28 * *"}}',
81+
)
82+
83+
reason = heartbeat._heartbeat_skip_reason_for_schedule(
84+
dt.datetime(2026, 6, 25, 23, 10, tzinfo=dt.timezone.utc)
85+
)
86+
87+
assert reason is None

0 commit comments

Comments
 (0)