Skip to content

Commit 522d6f7

Browse files
Pigbibicodex
andcommitted
fix: close runtime heartbeat review gaps
Co-Authored-By: Codex <noreply@openai.com>
1 parent 7645a60 commit 522d6f7

10 files changed

Lines changed: 725 additions & 154 deletions

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

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -60,8 +60,9 @@ jobs:
6060
RUNTIME_HEARTBEAT_MARKET_AWARE: ${{ vars.RUNTIME_HEARTBEAT_MARKET_AWARE || 'true' }}
6161
RUNTIME_HEARTBEAT_MARKET_CALENDAR: ${{ vars.LONGBRIDGE_MARKET_CALENDAR }}
6262
RUNTIME_HEARTBEAT_MARKET_TIMEZONE: ${{ vars.LONGBRIDGE_MARKET_TIMEZONE }}
63+
RUNTIME_HEARTBEAT_PUBLICATION_GRACE_MINUTES: ${{ vars.RUNTIME_HEARTBEAT_PUBLICATION_GRACE_MINUTES || '30' }}
6364
RUNTIME_HEARTBEAT_SCHEDULER_AWARE: ${{ vars.RUNTIME_HEARTBEAT_SCHEDULER_AWARE || 'true' }}
64-
RUNTIME_HEARTBEAT_SCHEDULER_LOCATION: ${{ vars.RUNTIME_HEARTBEAT_SCHEDULER_LOCATION || vars.CLOUD_RUN_REGION }}
65+
RUNTIME_HEARTBEAT_SCHEDULER_LOCATION: ${{ vars.RUNTIME_HEARTBEAT_SCHEDULER_LOCATION || vars.CLOUD_RUN_REGION || 'us-central1' }}
6566
RUNTIME_HEARTBEAT_EXPECTED_DAY_OF_MONTH: ${{ vars.RUNTIME_HEARTBEAT_EXPECTED_DAY_OF_MONTH }}
6667
RUNTIME_HEARTBEAT_EXPECTED_TIMEZONE: ${{ vars.RUNTIME_HEARTBEAT_EXPECTED_TIMEZONE }}
6768
RUNTIME_TARGET_ENABLED: ${{ vars.RUNTIME_TARGET_ENABLED }}
@@ -70,6 +71,7 @@ jobs:
7071
CLOUD_RUN_SERVICE: ${{ vars.CLOUD_RUN_SERVICE }}
7172
CLOUD_RUN_SERVICES: ${{ vars.CLOUD_RUN_SERVICES }}
7273
CLOUD_RUN_SERVICE_TARGETS_JSON: ${{ vars.CLOUD_RUN_SERVICE_TARGETS_JSON }}
74+
CLOUD_SCHEDULER_MAIN_TIME: ${{ vars.CLOUD_SCHEDULER_MAIN_TIME }}
7375
GLOBAL_TELEGRAM_CHAT_ID: ${{ vars.GLOBAL_TELEGRAM_CHAT_ID }}
7476
TELEGRAM_TOKEN: ${{ secrets.TELEGRAM_TOKEN }}
7577
TELEGRAM_TOKEN_SECRET_NAME: ${{ vars.TELEGRAM_TOKEN_SECRET_NAME }}

main.py

Lines changed: 42 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -623,7 +623,15 @@ def run_strategy(*, force_run: bool = False, validation_only: bool = False, vali
623623
composer = build_composer(dry_run_only_override=True if validation_only else None)
624624
reporting_adapters = composer.build_reporting_adapters()
625625
log_context, report = reporting_adapters.start_run()
626-
notification_adapters = composer.build_notification_adapters()
626+
notification_delivery_events: list[dict] = []
627+
try:
628+
notification_adapters = composer.build_notification_adapters(
629+
delivery_events=notification_delivery_events,
630+
)
631+
except TypeError as exc:
632+
if "delivery_events" not in str(exc):
633+
raise
634+
notification_adapters = composer.build_notification_adapters()
627635
strategy_plugin_signals, strategy_plugin_error = composer.load_strategy_plugin_signals(
628636
getattr(RUNTIME_SETTINGS, "strategy_plugin_mounts_json", None)
629637
)
@@ -717,7 +725,6 @@ def run_strategy(*, force_run: bool = False, validation_only: bool = False, vali
717725
return True
718726
if not validation_only:
719727
publish_strategy_plugin_alerts(strategy_plugin_signals, report=report)
720-
notification_delivery_events: list[dict] = []
721728
try:
722729
rebalance_runtime = composer.build_rebalance_runtime(
723730
silent_cycle_notifications=validation_only,
@@ -792,7 +799,6 @@ def run_strategy(*, force_run: bool = False, validation_only: bool = False, vali
792799
message=str(exc),
793800
error_type=type(exc).__name__,
794801
)
795-
finalize_runtime_report(report, status="error")
796802
reporting_adapters.log_event(
797803
log_context,
798804
"strategy_cycle_failed",
@@ -802,9 +808,39 @@ def run_strategy(*, force_run: bool = False, validation_only: bool = False, vali
802808
error_message=str(exc),
803809
)
804810
err = traceback.format_exc()
805-
notification_adapters.publish_cycle_notification(
806-
detailed_text=f"Strategy error:\n{err}",
807-
compact_text=_compact_error_notification(exc),
811+
try:
812+
notification_adapters.publish_cycle_notification(
813+
detailed_text=f"Strategy error:\n{err}",
814+
compact_text=_compact_error_notification(exc),
815+
)
816+
except Exception as notification_exc:
817+
notification_delivery_events.append(
818+
{
819+
"sink": "telegram",
820+
"delivery_status": "failed",
821+
"transport_acknowledged": False,
822+
"error_type": type(notification_exc).__name__,
823+
}
824+
)
825+
reporting_adapters.log_event(
826+
log_context,
827+
"strategy_error_notification_failed",
828+
message="Strategy error notification failed",
829+
severity="ERROR",
830+
error_type=type(notification_exc).__name__,
831+
)
832+
error_summary = {}
833+
notification_delivery_summary = _build_notification_delivery_summary(
834+
notification_delivery_events
835+
)
836+
if notification_delivery_summary:
837+
error_summary["notification_delivery_summary"] = (
838+
notification_delivery_summary
839+
)
840+
finalize_runtime_report(
841+
report,
842+
status="error",
843+
summary=error_summary or None,
808844
)
809845
return False
810846
finally:

scripts/cloud_run_runtime_guard.py

Lines changed: 114 additions & 48 deletions
Original file line numberDiff line numberDiff line change
@@ -41,50 +41,42 @@ def _env_bool(name: str, default: bool = False) -> bool:
4141

4242
def _load_services() -> list[str]:
4343
services = []
44+
enabled_target_services = []
45+
disabled_target_services = []
4446
for name in (
4547
"RUNTIME_GUARD_CLOUD_RUN_SERVICES",
4648
"CLOUD_RUN_SERVICES",
4749
"CLOUD_RUN_SERVICE",
4850
):
4951
services.extend(_split_values(os.environ.get(name)))
50-
if services:
51-
return list(dict.fromkeys(services))
52+
explicit_services = bool(services)
5253

5354
raw_targets = (os.environ.get("CLOUD_RUN_SERVICE_TARGETS_JSON") or "").strip()
5455
if raw_targets:
5556
try:
5657
payload = json.loads(raw_targets)
58+
defaults = payload.get("defaults") if isinstance(payload, dict) else {}
59+
defaults = defaults if isinstance(defaults, dict) else {}
5760
targets = payload.get("targets") if isinstance(payload, dict) else payload
5861
if isinstance(targets, list):
5962
for target in targets:
6063
if not isinstance(target, dict):
6164
continue
62-
if not _target_enabled(target):
63-
continue
64-
runtime_target = target.get("runtime_target") or target.get(
65-
"runtime_target_json"
66-
)
67-
if isinstance(runtime_target, str):
68-
try:
69-
runtime_target = json.loads(runtime_target)
70-
except json.JSONDecodeError:
71-
runtime_target = {}
72-
for key in ("service", "service_name", "cloud_run_service"):
73-
value = target.get(key) or (
74-
runtime_target.get(key)
75-
if isinstance(runtime_target, dict)
76-
else None
77-
)
78-
if value:
79-
services.extend(_split_values(str(value)))
80-
break
65+
target_services = _target_service_names(target, defaults)
66+
if _target_enabled(target, defaults):
67+
enabled_target_services.extend(target_services)
68+
else:
69+
disabled_target_services.extend(target_services)
8170
except json.JSONDecodeError as exc:
8271
raise RuntimeError(f"CLOUD_RUN_SERVICE_TARGETS_JSON is invalid: {exc}") from exc
8372

73+
if not explicit_services:
74+
services.extend(enabled_target_services)
75+
disabled = set(disabled_target_services) - set(enabled_target_services)
8476
seen = set()
8577
unique = []
8678
for service in services:
87-
if service not in seen:
79+
if service not in seen and service not in disabled:
8880
seen.add(service)
8981
unique.append(service)
9082
return unique
@@ -114,9 +106,29 @@ def _service_job_aliases(service: str) -> list[str]:
114106
def _scheduler_job_pattern_for_services(services: list[str]) -> str:
115107
candidates: list[str] = []
116108
for service in services:
117-
candidates.extend(_service_job_aliases(service))
109+
candidates.extend(_scheduler_job_names(service))
118110
unique = list(dict.fromkeys(candidates))
119-
return "|".join(re.escape(candidate) for candidate in unique)
111+
if not unique:
112+
return ""
113+
return r"^(?:" + "|".join(re.escape(candidate) for candidate in unique) + r")\Z"
114+
115+
116+
def _scheduler_job_names(service: str) -> list[str]:
117+
names = []
118+
for alias in _service_job_aliases(service):
119+
names.extend(
120+
(
121+
f"{alias}-scheduler",
122+
f"{alias}-probe-scheduler",
123+
f"{alias}-precheck-scheduler",
124+
)
125+
)
126+
return list(dict.fromkeys(names))
127+
128+
129+
def _job_matches_service(job_name: str, service: str) -> bool:
130+
normalized = str(job_name or "").strip().rsplit("/", 1)[-1]
131+
return normalized in _scheduler_job_names(service)
120132

121133

122134
def _entry_job_name(entry: dict[str, Any]) -> str:
@@ -133,7 +145,7 @@ def _scheduler_entry_since(
133145
matches = [
134146
service_since
135147
for service, service_since in service_since_by_name.items()
136-
if any(alias and alias in job_name for alias in _service_job_aliases(service))
148+
if _job_matches_service(job_name, service)
137149
]
138150
return max(matches) if matches else fallback
139151

@@ -149,7 +161,7 @@ def _is_duplicate_scheduler_failure(
149161

150162
tolerance = dt.timedelta(seconds=SCHEDULER_CLOUD_RUN_DEDUP_SECONDS)
151163
for service, failures in cloud_run_failures_by_service.items():
152-
if not any(alias and alias in job_name for alias in _service_job_aliases(service)):
164+
if not _job_matches_service(job_name, service):
153165
continue
154166
for failure in failures:
155167
cloud_run_timestamp = _parse_timestamp(failure.get("timestamp"))
@@ -235,22 +247,51 @@ def _format_timestamp(value: dt.datetime) -> str:
235247
return value.astimezone(dt.timezone.utc).isoformat().replace("+00:00", "Z")
236248

237249

238-
def _target_payloads() -> list[dict[str, Any]]:
250+
def _target_configuration() -> tuple[list[dict[str, Any]], dict[str, Any]]:
239251
raw_targets = (os.environ.get("CLOUD_RUN_SERVICE_TARGETS_JSON") or "").strip()
240252
if not raw_targets:
241-
return []
253+
return [], {}
242254
try:
243255
payload = json.loads(raw_targets)
244256
except json.JSONDecodeError:
245-
return []
257+
return [], {}
258+
defaults = payload.get("defaults") if isinstance(payload, dict) else {}
259+
defaults = defaults if isinstance(defaults, dict) else {}
246260
targets = payload.get("targets") if isinstance(payload, dict) else payload
247261
if not isinstance(targets, list):
248-
return []
249-
return [target for target in targets if isinstance(target, dict)]
262+
return [], defaults
263+
return [target for target in targets if isinstance(target, dict)], defaults
250264

251265

252-
def _runtime_target(target: dict[str, Any]) -> dict[str, Any]:
253-
runtime_target = target.get("runtime_target") or target.get("runtime_target_json")
266+
def _target_payloads() -> list[dict[str, Any]]:
267+
targets, _defaults = _target_configuration()
268+
return targets
269+
270+
271+
def _target_field(
272+
target: dict[str, Any],
273+
defaults: dict[str, Any],
274+
*names: str,
275+
) -> Any:
276+
target_env = target.get("env") if isinstance(target.get("env"), dict) else {}
277+
defaults_env = defaults.get("env") if isinstance(defaults.get("env"), dict) else {}
278+
for source in (target, target_env, defaults, defaults_env):
279+
for name in names:
280+
if name in source:
281+
return source[name]
282+
return None
283+
284+
285+
def _runtime_target(
286+
target: dict[str, Any],
287+
defaults: dict[str, Any] | None = None,
288+
) -> dict[str, Any]:
289+
runtime_target = _target_field(
290+
target,
291+
defaults or {},
292+
"runtime_target",
293+
"runtime_target_json",
294+
)
254295
if isinstance(runtime_target, str):
255296
try:
256297
runtime_target = json.loads(runtime_target)
@@ -270,32 +311,57 @@ def _coerce_bool(value: Any, default: bool) -> bool:
270311
return text in {"1", "true", "yes", "y", "on"}
271312

272313

273-
def _target_enabled(target: dict[str, Any]) -> bool:
274-
runtime_target = _runtime_target(target)
314+
def _target_enabled(
315+
target: dict[str, Any],
316+
defaults: dict[str, Any] | None = None,
317+
) -> bool:
318+
defaults = defaults or {}
319+
runtime_target = _runtime_target(target, defaults)
320+
value = _target_field(
321+
target,
322+
defaults,
323+
"runtime_target_enabled",
324+
"RUNTIME_TARGET_ENABLED",
325+
)
326+
if value is not None:
327+
return _coerce_bool(value, True)
275328
for key in ("runtime_target_enabled", "RUNTIME_TARGET_ENABLED"):
276-
if key in target:
277-
return _coerce_bool(target.get(key), True)
278329
if key in runtime_target:
279330
return _coerce_bool(runtime_target.get(key), True)
280331
return True
281332

282333

283-
def _target_service_names(target: dict[str, Any]) -> list[str]:
284-
runtime_target = _runtime_target(target)
285-
for key in ("service", "service_name", "cloud_run_service"):
286-
value = target.get(key) or runtime_target.get(key)
287-
if value:
288-
return _split_values(str(value))
334+
def _target_service_names(
335+
target: dict[str, Any],
336+
defaults: dict[str, Any] | None = None,
337+
) -> list[str]:
338+
defaults = defaults or {}
339+
runtime_target = _runtime_target(target, defaults)
340+
value = _target_field(
341+
target,
342+
defaults,
343+
"service",
344+
"service_name",
345+
"cloud_run_service",
346+
)
347+
if value is None:
348+
for key in ("service", "service_name", "cloud_run_service"):
349+
if runtime_target.get(key):
350+
value = runtime_target[key]
351+
break
352+
if value:
353+
return _split_values(str(value))
289354
return []
290355

291356

292357
def _region_for_service(service: str) -> str:
293-
for target in _target_payloads():
294-
if service not in _target_service_names(target):
358+
targets, defaults = _target_configuration()
359+
for target in targets:
360+
if service not in _target_service_names(target, defaults):
295361
continue
296-
runtime_target = _runtime_target(target)
362+
runtime_target = _runtime_target(target, defaults)
297363
for key in ("region", "cloud_run_region", "location"):
298-
value = target.get(key) or runtime_target.get(key)
364+
value = _target_field(target, defaults, key) or runtime_target.get(key)
299365
if value:
300366
return str(value).strip()
301367
return (
@@ -598,7 +664,7 @@ def main() -> int:
598664
entries = [
599665
entry
600666
for entry in entries
601-
if regex.search(str(_labels(entry).get("job_id") or _labels(entry).get("job_name") or ""))
667+
if regex.search(_entry_job_name(entry).rsplit("/", 1)[-1])
602668
]
603669
failures = []
604670
for entry in entries:

0 commit comments

Comments
 (0)