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
5 changes: 5 additions & 0 deletions .github/workflows/execution-report-heartbeat.yml
Original file line number Diff line number Diff line change
Expand Up @@ -48,11 +48,16 @@ jobs:
RUNTIME_HEARTBEAT_MARKET_AWARE: ${{ vars.RUNTIME_HEARTBEAT_MARKET_AWARE || 'true' }}
RUNTIME_HEARTBEAT_MARKET_CALENDAR: ${{ vars.SCHWAB_MARKET_CALENDAR }}
RUNTIME_HEARTBEAT_MARKET_TIMEZONE: ${{ vars.SCHWAB_MARKET_TIMEZONE }}
RUNTIME_HEARTBEAT_PUBLICATION_GRACE_MINUTES: ${{ vars.RUNTIME_HEARTBEAT_PUBLICATION_GRACE_MINUTES || '30' }}
RUNTIME_HEARTBEAT_SCHEDULER_AWARE: ${{ vars.RUNTIME_HEARTBEAT_SCHEDULER_AWARE || 'true' }}
RUNTIME_HEARTBEAT_SCHEDULER_LOCATION: ${{ vars.RUNTIME_HEARTBEAT_SCHEDULER_LOCATION || vars.CLOUD_RUN_REGION || 'us-central1' }}

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1 Badge Use the configured Cloud Scheduler location

When CLOUD_SCHEDULER_LOCATION differs from CLOUD_RUN_REGION, the deployment workflow creates or updates jobs in the former, but this heartbeat passes only the latter as the scheduler-location fallback. Any target with a two-field schedule then looks for its deployed job in the wrong region, reports an unable to resolve effective five-field scheduler cron policy error, and fails every scheduled heartbeat. Include vars.CLOUD_SCHEDULER_LOCATION in this fallback, consistent with .github/workflows/sync-cloud-run-env.yml.

Useful? React with 👍 / 👎.

RUNTIME_TARGET_ENABLED: ${{ vars.RUNTIME_TARGET_ENABLED }}
RUNTIME_TARGET_JSON: ${{ vars.RUNTIME_TARGET_JSON }}
CLOUD_RUN_REGION: ${{ vars.CLOUD_RUN_REGION }}
CLOUD_RUN_SERVICE: ${{ vars.CLOUD_RUN_SERVICE }}
CLOUD_RUN_SERVICES: ${{ vars.CLOUD_RUN_SERVICES }}
CLOUD_RUN_SERVICE_TARGETS_JSON: ${{ vars.CLOUD_RUN_SERVICE_TARGETS_JSON }}
CLOUD_SCHEDULER_MAIN_TIME: ${{ vars.CLOUD_SCHEDULER_MAIN_TIME }}
GLOBAL_TELEGRAM_CHAT_ID: ${{ vars.GLOBAL_TELEGRAM_CHAT_ID }}
TELEGRAM_TOKEN: ${{ secrets.TELEGRAM_TOKEN }}
steps:
Expand Down
158 changes: 112 additions & 46 deletions scripts/cloud_run_runtime_guard.py
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,8 @@ def _env_bool(name: str, default: bool = False) -> bool:

def _load_services() -> list[str]:
services = []
enabled_target_services = []
disabled_target_services = []
for name in (
"RUNTIME_GUARD_CLOUD_RUN_SERVICES",
"CLOUD_RUN_SERVICES",
Expand All @@ -52,37 +54,27 @@ def _load_services() -> list[str]:
if raw_targets:
try:
payload = json.loads(raw_targets)
defaults = payload.get("defaults") if isinstance(payload, dict) else {}
defaults = defaults if isinstance(defaults, dict) else {}
targets = payload.get("targets") if isinstance(payload, dict) else payload
if isinstance(targets, list):
for target in targets:
if not isinstance(target, dict):
continue
if not _target_enabled(target):
continue
runtime_target = target.get("runtime_target") or target.get(
"runtime_target_json"
)
if isinstance(runtime_target, str):
try:
runtime_target = json.loads(runtime_target)
except json.JSONDecodeError:
runtime_target = {}
for key in ("service", "service_name", "cloud_run_service"):
value = target.get(key) or (
runtime_target.get(key)
if isinstance(runtime_target, dict)
else None
)
if value:
services.extend(_split_values(str(value)))
break
target_services = _target_service_names(target, defaults)
if _target_enabled(target, defaults):
enabled_target_services.extend(target_services)
else:
disabled_target_services.extend(target_services)
except json.JSONDecodeError as exc:
raise RuntimeError(f"CLOUD_RUN_SERVICE_TARGETS_JSON is invalid: {exc}") from exc

services.extend(enabled_target_services)
disabled = set(disabled_target_services) - set(enabled_target_services)
seen = set()
unique = []
for service in services:
if service not in seen:
if service not in seen and service not in disabled:
seen.add(service)
unique.append(service)
return unique
Expand Down Expand Up @@ -112,9 +104,29 @@ def _service_job_aliases(service: str) -> list[str]:
def _scheduler_job_pattern_for_services(services: list[str]) -> str:
candidates: list[str] = []
for service in services:
candidates.extend(_service_job_aliases(service))
candidates.extend(_scheduler_job_names(service))
unique = list(dict.fromkeys(candidates))
return "|".join(re.escape(candidate) for candidate in unique)
if not unique:
return ""
return r"^(?:" + "|".join(re.escape(candidate) for candidate in unique) + r")\Z"


def _scheduler_job_names(service: str) -> list[str]:
names = []
for alias in _service_job_aliases(service):
names.extend(
(
f"{alias}-scheduler",
f"{alias}-probe-scheduler",
f"{alias}-precheck-scheduler",
)
)
return list(dict.fromkeys(names))


def _job_matches_service(job_name: str, service: str) -> bool:
normalized = str(job_name or "").strip().rsplit("/", 1)[-1]
return normalized in _scheduler_job_names(service)


def _entry_job_name(entry: dict[str, Any]) -> str:
Expand All @@ -131,7 +143,7 @@ def _scheduler_entry_since(
matches = [
service_since
for service, service_since in service_since_by_name.items()
if any(alias and alias in job_name for alias in _service_job_aliases(service))
if _job_matches_service(job_name, service)
]
return max(matches) if matches else fallback

Expand All @@ -147,7 +159,7 @@ def _is_duplicate_scheduler_failure(

tolerance = dt.timedelta(seconds=SCHEDULER_CLOUD_RUN_DEDUP_SECONDS)
for service, failures in cloud_run_failures_by_service.items():
if not any(alias and alias in job_name for alias in _service_job_aliases(service)):
if not _job_matches_service(job_name, service):
continue
for failure in failures:
cloud_run_timestamp = _parse_timestamp(failure.get("timestamp"))
Expand Down Expand Up @@ -233,22 +245,51 @@ def _format_timestamp(value: dt.datetime) -> str:
return value.astimezone(dt.timezone.utc).isoformat().replace("+00:00", "Z")


def _target_payloads() -> list[dict[str, Any]]:
def _target_configuration() -> tuple[list[dict[str, Any]], dict[str, Any]]:
raw_targets = (os.environ.get("CLOUD_RUN_SERVICE_TARGETS_JSON") or "").strip()
if not raw_targets:
return []
return [], {}
try:
payload = json.loads(raw_targets)
except json.JSONDecodeError:
return []
return [], {}
defaults = payload.get("defaults") if isinstance(payload, dict) else {}
defaults = defaults if isinstance(defaults, dict) else {}
targets = payload.get("targets") if isinstance(payload, dict) else payload
if not isinstance(targets, list):
return []
return [target for target in targets if isinstance(target, dict)]
return [], defaults
return [target for target in targets if isinstance(target, dict)], defaults


def _runtime_target(target: dict[str, Any]) -> dict[str, Any]:
runtime_target = target.get("runtime_target") or target.get("runtime_target_json")
def _target_payloads() -> list[dict[str, Any]]:
targets, _defaults = _target_configuration()
return targets


def _target_field(
target: dict[str, Any],
defaults: dict[str, Any],
*names: str,
) -> Any:
target_env = target.get("env") if isinstance(target.get("env"), dict) else {}
defaults_env = defaults.get("env") if isinstance(defaults.get("env"), dict) else {}
for source in (target, target_env, defaults, defaults_env):
for name in names:
if name in source:
return source[name]
return None


def _runtime_target(
target: dict[str, Any],
defaults: dict[str, Any] | None = None,
) -> dict[str, Any]:
runtime_target = _target_field(
target,
defaults or {},
"runtime_target",
"runtime_target_json",
)
if isinstance(runtime_target, str):
try:
runtime_target = json.loads(runtime_target)
Expand All @@ -268,32 +309,57 @@ def _coerce_bool(value: Any, default: bool) -> bool:
return text in {"1", "true", "yes", "y", "on"}


def _target_enabled(target: dict[str, Any]) -> bool:
runtime_target = _runtime_target(target)
def _target_enabled(
target: dict[str, Any],
defaults: dict[str, Any] | None = None,
) -> bool:
defaults = defaults or {}
runtime_target = _runtime_target(target, defaults)
value = _target_field(
target,
defaults,
"runtime_target_enabled",
"RUNTIME_TARGET_ENABLED",
)
if value is not None:
return _coerce_bool(value, True)
for key in ("runtime_target_enabled", "RUNTIME_TARGET_ENABLED"):
if key in target:
return _coerce_bool(target.get(key), True)
if key in runtime_target:
return _coerce_bool(runtime_target.get(key), True)
return True


def _target_service_names(target: dict[str, Any]) -> list[str]:
runtime_target = _runtime_target(target)
for key in ("service", "service_name", "cloud_run_service"):
value = target.get(key) or runtime_target.get(key)
if value:
return _split_values(str(value))
def _target_service_names(
target: dict[str, Any],
defaults: dict[str, Any] | None = None,
) -> list[str]:
defaults = defaults or {}
runtime_target = _runtime_target(target, defaults)
value = _target_field(
target,
defaults,
"service",
"service_name",
"cloud_run_service",
)
if value is None:
for key in ("service", "service_name", "cloud_run_service"):
if runtime_target.get(key):
value = runtime_target[key]
break
if value:
return _split_values(str(value))
return []


def _region_for_service(service: str) -> str:
for target in _target_payloads():
if service not in _target_service_names(target):
targets, defaults = _target_configuration()
for target in targets:
if service not in _target_service_names(target, defaults):
continue
runtime_target = _runtime_target(target)
runtime_target = _runtime_target(target, defaults)
for key in ("region", "cloud_run_region", "location"):
value = target.get(key) or runtime_target.get(key)
value = _target_field(target, defaults, key) or runtime_target.get(key)
if value:
return str(value).strip()
return (
Expand Down Expand Up @@ -595,7 +661,7 @@ def main() -> int:
entries = [
entry
for entry in entries
if regex.search(str(_labels(entry).get("job_id") or _labels(entry).get("job_name") or ""))
if regex.search(_entry_job_name(entry).rsplit("/", 1)[-1])
]
failures = []
for entry in entries:
Expand Down
Loading