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
35 changes: 28 additions & 7 deletions scripts/publish_daily_runtime_projection.py
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,7 @@
from typing import Any
from urllib.parse import urlsplit

from google.api_core.exceptions import GoogleAPICallError, PreconditionFailed
from google.api_core.exceptions import Forbidden, GoogleAPICallError, PreconditionFailed

from scripts import execution_report_heartbeat as heartbeat
from scripts.daily_runtime_projection import project_listed_reports
Expand All @@ -46,6 +46,20 @@ class _Rejected(Exception):
"""The projection was not published. Nothing was written."""


class _KnownFailure(RuntimeError):
"""A safe, fixed diagnostic category for a known failure phase."""

_CODES = frozenset(
{"schedule_unavailable", "storage_write_permission_denied", "write_unknown"}
)

def __init__(self, code: str) -> None:
if code not in self._CODES:
raise ValueError("invalid diagnostic code")
self.code = code
super().__init__(code)


def _prefix(value: str) -> str:
try:
parsed = urlsplit(value.strip())
Expand Down Expand Up @@ -417,8 +431,12 @@ def _upload(client: Any, prefix: str, business_date: str, body: str) -> str:
)
except PreconditionFailed:
return "already_recorded"
except GoogleAPICallError as exc:
raise RuntimeError("write_unknown") from exc
except Forbidden:
raise _KnownFailure("storage_write_permission_denied") from None
except GoogleAPICallError:
raise _KnownFailure("write_unknown") from None
except Exception:
raise _KnownFailure("write_unknown") from None
return uri


Expand Down Expand Up @@ -574,20 +592,20 @@ def publish(
try:
payload = heartbeat._describe_cloud_run_service(service, project=project)
if not isinstance(payload, dict):
raise RuntimeError("schedule_unavailable")
raise _KnownFailure("schedule_unavailable")
deployed = heartbeat._deployed_runtime_target(payload)
hydrated = _profiles_from_same_readback(target, payload, project=project)
hydrated = heartbeat._hydrate_runtime_target_schedules(hydrated, project=project)
except RuntimeError as exc:
raise RuntimeError("schedule_unavailable") from exc
raise _KnownFailure("schedule_unavailable") from None
if not _deployment_matches_paper(target, deployed) or not _deployment_enabled(payload, deployed):
raise _Rejected("deployed_target_rejected")
if len(hydrated) != 1:
raise _Rejected("target_not_unique")
try:
hydrated[0], scheduler_error = _confirmed_scheduler(hydrated[0], project=project)
except RuntimeError as exc:
raise RuntimeError("schedule_unavailable") from exc
raise _KnownFailure("schedule_unavailable") from None
since = observed - dt.timedelta(hours=lookback)
globs = (report_globs or heartbeat._report_globs)(since, observed)
if not globs:
Expand Down Expand Up @@ -643,7 +661,7 @@ def read_payload(
raise _Rejected("target_not_unique")
business_date = records[0].get("business_date")
if not isinstance(business_date, str) or not business_date:
raise RuntimeError("schedule_unavailable")
raise _KnownFailure("schedule_unavailable")
# QRS stores the fixed PAPER target with its canonical lower-case scope.
# Keep the actual scope check above case-insensitive and serialize this
# normalized representation once so GCS and the optional POST share bytes.
Expand Down Expand Up @@ -673,6 +691,9 @@ def main(argv: Sequence[str] | None = None) -> int:
del argv
try:
status, business_date, sync_status = publish(os.environ, now=dt.datetime.now(dt.timezone.utc))
except _KnownFailure as exc:
print(f"daily runtime projection failed; category={exc.code}")
return 1
except _Rejected:
print("daily runtime projection rejected")
return 2
Expand Down
54 changes: 50 additions & 4 deletions tests/test_publish_daily_runtime_projection.py
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,7 @@
from zoneinfo import ZoneInfo

import pytest
from google.api_core.exceptions import GoogleAPICallError, PreconditionFailed
from google.api_core.exceptions import Forbidden, GoogleAPICallError, PreconditionFailed

from scripts import execution_report_heartbeat as heartbeat
from scripts import publish_daily_runtime_projection as publisher
Expand Down Expand Up @@ -401,8 +401,10 @@ def test_existing_object_is_not_overwritten_and_unknown_write_is_not_retried(mon
assert blob.uploads[0]["kwargs"]["if_generation_match"] == 0

unknown = _Blob(fail=RuntimeError("token=secret"))
with pytest.raises(RuntimeError, match="token=secret"):
with pytest.raises(publisher._KnownFailure) as failure:
_publish(monkeypatch, _env(), [], blob=unknown)
assert failure.value.code == "write_unknown"
assert "token=secret" not in str(failure.value)
assert len(unknown.uploads) == 1


Expand All @@ -417,7 +419,7 @@ def cat(uri: str, *, project: str | None):
assert "secret payload" not in capsys.readouterr().err


def test_main_unknown_write_prints_no_secret_and_does_not_retry(monkeypatch, capsys) -> None:
def test_main_unknown_write_prints_safe_category_and_does_not_retry(monkeypatch, capsys) -> None:
blob = _Blob(fail=RuntimeError("token=secret"))
monkeypatch.setattr(publisher, "_storage_client", lambda: _Client(blob))
monkeypatch.setattr(heartbeat, "_describe_cloud_run_service", _describe)
Expand All @@ -428,11 +430,55 @@ def test_main_unknown_write_prints_no_secret_and_does_not_retry(monkeypatch, cap
monkeypatch.setenv(key, value)
assert publisher.main([]) == 1
captured = capsys.readouterr()
assert captured.out.strip() == "daily runtime projection failed"
assert captured.out.strip() == "daily runtime projection failed; category=write_unknown"
assert "token=secret" not in captured.out + captured.err
assert len(blob.uploads) == 1


def test_forbidden_write_is_classified_once_without_qrs_post(monkeypatch) -> None:
blob = _Blob(fail=Forbidden("token=secret permission denied"))
posts = []
with pytest.raises(publisher._KnownFailure) as failure:
_publish(
monkeypatch,
_sync_env(),
[],
blob=blob,
describe=_sync_describe,
http_post=lambda *args, **kwargs: posts.append((args, kwargs)),
)
assert failure.value.code == "storage_write_permission_denied"
assert "token=secret" not in str(failure.value)
assert len(blob.uploads) == 1
assert posts == []


def test_main_keeps_unclassified_exceptions_generic_and_redacted(monkeypatch, capsys) -> None:
monkeypatch.setattr(
publisher,
"publish",
lambda *_args, **_kwargs: (_ for _ in ()).throw(RuntimeError("token=secret detail")),
)
assert publisher.main([]) == 1
captured = capsys.readouterr()
assert captured.out.strip() == "daily runtime projection failed"
assert "token=secret" not in captured.out + captured.err


def test_main_displays_only_safe_schedule_failure_category(monkeypatch, capsys) -> None:
for key, value in _env().items():
monkeypatch.setenv(key, value)

def fail_describe(*_args, **_kwargs):
raise RuntimeError("token=secret detail")

monkeypatch.setattr(heartbeat, "_describe_cloud_run_service", fail_describe)
assert publisher.main([]) == 1
captured = capsys.readouterr()
assert captured.out.strip() == "daily runtime projection failed; category=schedule_unavailable"
assert "token=secret" not in captured.out + captured.err


def test_main_hides_errors_and_keeps_heartbeat_uninvoked(monkeypatch, capsys) -> None:
assert "heartbeat.main" not in (ROOT / "scripts/publish_daily_runtime_projection.py").read_text(encoding="utf-8")
monkeypatch.delenv("RUNTIME_DAILY_PROJECTION_ENABLED", raising=False)
Expand Down
Loading