diff --git a/.github/workflows/runtime-target-lifecycle.yml b/.github/workflows/runtime-target-lifecycle.yml index b6d479f..3991c26 100644 --- a/.github/workflows/runtime-target-lifecycle.yml +++ b/.github/workflows/runtime-target-lifecycle.yml @@ -2,6 +2,12 @@ name: Runtime Target Lifecycle on: workflow_dispatch: + inputs: + metadata_only: + description: Inspect the configured serving revision without invoking the service + type: boolean + required: true + default: false workflow_run: workflows: ["Deploy Cloud Run"] types: [completed] @@ -20,6 +26,7 @@ concurrency: jobs: lifecycle: name: Publish Firstrade target lifecycle + if: ${{ !(github.event_name == 'workflow_dispatch' && inputs.metadata_only) }} runs-on: ubuntu-latest timeout-minutes: 15 permissions: @@ -205,3 +212,32 @@ jobs: echo echo "This workflow is read-only: it cannot enable a target, alter its lane, or submit an order." } >> "$GITHUB_STEP_SUMMARY" + + account_data_readiness: + name: Inspect Firstrade account-data configuration + if: ${{ github.event_name == 'workflow_dispatch' && inputs.metadata_only }} + runs-on: ubuntu-latest + timeout-minutes: 5 + permissions: + contents: read + id-token: write + env: + CLOUD_RUN_REGION: ${{ vars.CLOUD_RUN_REGION }} + CLOUD_RUN_SERVICE: ${{ secrets.CLOUD_RUN_SERVICE }} + steps: + - name: Checkout repository + uses: actions/checkout@v6 + + - name: Authenticate to Google Cloud + uses: google-github-actions/auth@v3 + with: + workload_identity_provider: ${{ env.GCP_WORKLOAD_IDENTITY_PROVIDER }} + service_account: ${{ env.GCP_WORKLOAD_IDENTITY_SERVICE_ACCOUNT }} + + - name: Set up gcloud + uses: google-github-actions/setup-gcloud@v3 + with: + project_id: ${{ env.GCP_PROJECT_ID }} + + - name: Inspect serving account-data configuration + run: python3 scripts/inspect_account_data_readiness.py diff --git a/scripts/inspect_account_data_readiness.py b/scripts/inspect_account_data_readiness.py new file mode 100644 index 0000000..5391422 --- /dev/null +++ b/scripts/inspect_account_data_readiness.py @@ -0,0 +1,283 @@ +"""Read Cloud Run metadata for the configured Firstrade serving revision.""" + +from __future__ import annotations + +import hashlib +import json +import os +import re +import subprocess +import sys +import urllib.error +import urllib.parse +import urllib.request +from collections.abc import Callable, Mapping +from typing import Any + +PROJECT_ID = "firstradequant" +API_ROOT = "https://run.googleapis.com/v2" +SERVICE_NAME_RE = re.compile(r"[a-z][a-z0-9-]{0,62}\Z") +REGION_PATH_PATTERN = r"[a-z]+(?:-[a-z0-9]+)+[0-9]" +REGION_RE = re.compile(REGION_PATH_PATTERN + r"\Z") +COMMIT_RE = re.compile(r"[0-9a-fA-F]{40}\Z") +MAX_RESPONSE_BYTES = 2_000_000 + + +class DiagnosticFailure(Exception): + def __init__(self, reason: str, http_status: int | None = None) -> None: + self.reason = reason + self.http_status = http_status + + +def _safe_resource_part(value: str, pattern: re.Pattern[str]) -> str: + if not isinstance(value, str) or not pattern.fullmatch(value): + raise DiagnosticFailure("invalid_target_configuration") + return value + + +def _service_path(region: str, service: str) -> str: + return ( + f"projects/{PROJECT_ID}/locations/{region}/services/{service}" + ) + + +class _NoRedirect(urllib.request.HTTPRedirectHandler): + def redirect_request(self, req: Any, fp: Any, code: int, msg: str, headers: Any, newurl: str) -> None: + raise DiagnosticFailure("redirect_rejected", code) + + +def _request_json(url: str, token: str) -> Mapping[str, Any]: + parsed = urllib.parse.urlsplit(url) + if ( + parsed.scheme != "https" + or parsed.netloc != "run.googleapis.com" + or parsed.query + or parsed.fragment + or not re.fullmatch( + rf"/v2/projects/{PROJECT_ID}/locations/{REGION_PATH_PATTERN}" + r"/services/[a-z][a-z0-9-]{0,62}(?:/revisions/[a-z][a-z0-9-]{0,62})?", + parsed.path, + ) + ): + raise DiagnosticFailure("request_target_rejected") + + request = urllib.request.Request( + url, + headers={"Authorization": f"Bearer {token}", "Accept": "application/json"}, + method="GET", + ) + opener = urllib.request.build_opener(_NoRedirect()) + try: + with opener.open(request, timeout=15) as response: + raw = response.read(MAX_RESPONSE_BYTES + 1) + except urllib.error.HTTPError as exc: + raise DiagnosticFailure("metadata_http_error", int(exc.code)) from None + except DiagnosticFailure: + raise + except (urllib.error.URLError, TimeoutError, OSError): + raise DiagnosticFailure("metadata_transport_error") from None + if len(raw) > MAX_RESPONSE_BYTES: + raise DiagnosticFailure("metadata_response_too_large") + + try: + payload = json.loads(raw) + except (UnicodeDecodeError, json.JSONDecodeError): + raise DiagnosticFailure("metadata_response_invalid") from None + if not isinstance(payload, Mapping): + raise DiagnosticFailure("metadata_response_invalid") + return payload + + +def _ready(document: Mapping[str, Any]) -> bool: + if document.get("reconciling", False) is not False: + return False + generation = document.get("generation") + observed_generation = document.get("observedGeneration") + if not isinstance(generation, str) or not generation or generation != observed_generation: + return False + conditions = document.get("conditions") + if not isinstance(conditions, list): + return False + ready = [ + condition + for condition in conditions + if isinstance(condition, Mapping) and condition.get("type") == "Ready" + ] + return len(ready) == 1 and ready[0].get("state") == "CONDITION_SUCCEEDED" + + +def _service_ready(service: Mapping[str, Any]) -> bool: + if service.get("reconciling", False) is not False: + return False + generation = service.get("generation") + observed_generation = service.get("observedGeneration") + if not isinstance(generation, str) or not generation or generation != observed_generation: + return False + terminal = service.get("terminalCondition") + return isinstance(terminal, Mapping) and terminal.get("state") == "CONDITION_SUCCEEDED" + + +def _service_fingerprint(service: Mapping[str, Any]) -> str: + selected = { + "etag": service.get("etag"), + "template": service.get("template"), + "trafficStatuses": service.get("trafficStatuses"), + "reconciling": service.get("reconciling"), + "generation": service.get("generation"), + "observedGeneration": service.get("observedGeneration"), + "terminalCondition": service.get("terminalCondition"), + "conditions": service.get("conditions"), + } + try: + encoded = json.dumps(selected, sort_keys=True, separators=(",", ":")) + except (TypeError, ValueError): + raise DiagnosticFailure("metadata_response_invalid") from None + return hashlib.sha256(encoded.encode("utf-8")).hexdigest() + + +def _serving_revision(service: Mapping[str, Any], service_path: str) -> str: + traffic = service.get("trafficStatuses") + if not isinstance(traffic, list) or not traffic: + raise DiagnosticFailure("serving_traffic_unavailable") + total = 0 + positive: list[str] = [] + for entry in traffic: + if not isinstance(entry, Mapping): + raise DiagnosticFailure("serving_traffic_invalid") + percent = entry.get("percent", 0) + revision = entry.get("revision") + if isinstance(percent, bool) or not isinstance(percent, int) or not 0 <= percent <= 100: + raise DiagnosticFailure("serving_traffic_invalid") + if not isinstance(revision, str) or not revision: + raise DiagnosticFailure("serving_traffic_invalid") + prefix = f"{service_path}/revisions/" + if revision.startswith(prefix): + revision = revision.removeprefix(prefix) + elif "/" in revision: + raise DiagnosticFailure("serving_traffic_invalid") + if not SERVICE_NAME_RE.fullmatch(revision): + raise DiagnosticFailure("serving_traffic_invalid") + total += percent + if percent > 0: + positive.append(revision) + if total != 100 or len(positive) != 1: + raise DiagnosticFailure("serving_traffic_ambiguous") + return positive[0] + + +def _selector_configuration(revision: Mapping[str, Any]) -> tuple[bool, str]: + containers = revision.get("containers") + if not isinstance(containers, list) or len(containers) != 1: + raise DiagnosticFailure("revision_container_invalid") + env = containers[0].get("env") if isinstance(containers[0], Mapping) else None + if env is None: + env = [] + if not isinstance(env, list): + raise DiagnosticFailure("revision_environment_invalid") + matches = [ + entry + for entry in env + if isinstance(entry, Mapping) and entry.get("name") == "FIRSTRADE_ACCOUNT" + ] + if len(matches) > 1: + raise DiagnosticFailure("selector_configuration_ambiguous") + if not matches: + return False, "absent" + entry = matches[0] + value_source = entry.get("valueSource") + if isinstance(value_source, Mapping) and isinstance( + value_source.get("secretKeyRef"), Mapping + ): + return True, "secret_reference" + if isinstance(entry.get("value"), str) and entry["value"].strip(): + return True, "literal" + return False, "absent" + + +def inspect_readiness( + service_name: str, + region: str, + token: str, + *, + request_json: Callable[[str, str], Mapping[str, Any]] = _request_json, +) -> dict[str, Any]: + service_name = _safe_resource_part(service_name, SERVICE_NAME_RE) + region = _safe_resource_part(region, REGION_RE) + if not token: + raise DiagnosticFailure("authentication_unavailable") + service_path = _service_path(region, service_name) + service_url = f"{API_ROOT}/{service_path}" + before = request_json(service_url, token) + if before.get("name") != service_path: + raise DiagnosticFailure("service_identity_mismatch") + if not _service_ready(before): + raise DiagnosticFailure("service_not_ready") + if not before.get("etag"): + raise DiagnosticFailure("service_version_unavailable") + revision_name = _serving_revision(before, service_path) + revision_path = f"{service_path}/revisions/{urllib.parse.quote(revision_name, safe='-') }" + revision = request_json(f"{API_ROOT}/{revision_path}", token) + expected_revision_path = f"{service_path}/revisions/{revision_name}" + if revision.get("name") != expected_revision_path: + raise DiagnosticFailure("revision_identity_mismatch") + if not _ready(revision): + raise DiagnosticFailure("revision_not_ready") + source_commit = (revision.get("labels") or {}).get("commit-sha") + if not isinstance(source_commit, str) or not COMMIT_RE.fullmatch(source_commit): + raise DiagnosticFailure("source_commit_unavailable") + selector_configured, selector_source = _selector_configuration(revision) + after = request_json(service_url, token) + if after.get("name") != service_path: + raise DiagnosticFailure("service_identity_mismatch") + if _service_fingerprint(before) != _service_fingerprint(after): + raise DiagnosticFailure("service_changed_during_read") + return { + "status": "verified", + "reason": "metadata_only_verified", + "http_status": None, + "source_commit": source_commit.lower(), + "selector_configured": selector_configured, + "selector_source": selector_source, + } + + +def _access_token() -> str: + try: + result = subprocess.run( + ["gcloud", "auth", "print-access-token"], + check=False, + capture_output=True, + text=True, + timeout=15, + ) + except (OSError, subprocess.TimeoutExpired): + raise DiagnosticFailure("authentication_unavailable") from None + if result.returncode != 0 or not result.stdout.strip(): + raise DiagnosticFailure("authentication_unavailable") + return result.stdout.strip() + + +def main() -> int: + try: + result = inspect_readiness( + os.environ.get("CLOUD_RUN_SERVICE", ""), + os.environ.get("CLOUD_RUN_REGION", ""), + _access_token(), + ) + exit_code = 0 + except DiagnosticFailure as exc: + result = { + "status": "blocked", + "reason": exc.reason, + "http_status": exc.http_status, + } + exit_code = 1 + except Exception: # noqa: BLE001 - keep unexpected provider errors out of logs + result = {"status": "blocked", "reason": "diagnostic_failed", "http_status": None} + exit_code = 1 + print(json.dumps(result, sort_keys=True, separators=(",", ":"))) + return exit_code + + +if __name__ == "__main__": + sys.exit(main()) diff --git a/tests/test_account_data_readiness.py b/tests/test_account_data_readiness.py new file mode 100644 index 0000000..cd0a388 --- /dev/null +++ b/tests/test_account_data_readiness.py @@ -0,0 +1,257 @@ +from __future__ import annotations + +import json +import urllib.error +from collections.abc import Mapping +from pathlib import Path +from typing import Any, Self + +import pytest + +from scripts import inspect_account_data_readiness as readiness + +PROJECT = "projects/firstradequant/locations/us-central1/services/firstrade-platform" +REVISION = f"{PROJECT}/revisions/firstrade-platform-abc123" +COMMIT = "a" * 40 +SENTINEL = "private-selector-value-do-not-print" + + +def _service() -> dict[str, Any]: + return { + "name": PROJECT, + "etag": "etag-one", + "reconciling": False, + "generation": "7", + "observedGeneration": "7", + "terminalCondition": {"state": "CONDITION_SUCCEEDED"}, + "conditions": [], + "template": { + "containers": [ + {"env": [{"name": "FIRSTRADE_ACCOUNT", "value": SENTINEL}]} + ] + }, + "trafficStatuses": [{"revision": REVISION, "percent": 100}], + } + + +def _revision(*, env: list[dict[str, Any]] | None = None) -> dict[str, Any]: + return { + "name": REVISION, + "reconciling": False, + "generation": "3", + "observedGeneration": "3", + "conditions": [{"type": "Ready", "state": "CONDITION_SUCCEEDED"}], + "labels": {"commit-sha": COMMIT}, + "containers": [{"env": env or []}], + } + + +def _runner( + service_responses: list[Mapping[str, Any]] | None = None, + revision: Mapping[str, Any] | None = None, +) -> tuple[Any, list[str]]: + queue = list(service_responses or [_service(), _service()]) + urls: list[str] = [] + + def request(url: str, token: str) -> Mapping[str, Any]: + assert token == "token" + urls.append(url) + if url.endswith("/revisions/firstrade-platform-abc123"): + return revision or _revision() + return queue.pop(0) + + return request, urls + + +def test_inspects_only_exact_project_region_service_and_serving_revision() -> None: + request, urls = _runner(revision=_revision(env=[{"name": "FIRSTRADE_ACCOUNT", "value": SENTINEL}])) + + result = readiness.inspect_readiness( + "firstrade-platform", "us-central1", "token", request_json=request + ) + + assert result == { + "status": "verified", + "reason": "metadata_only_verified", + "http_status": None, + "source_commit": COMMIT, + "selector_configured": True, + "selector_source": "literal", + } + assert urls == [ + f"{readiness.API_ROOT}/{PROJECT}", + f"{readiness.API_ROOT}/{REVISION}", + f"{readiness.API_ROOT}/{PROJECT}", + ] + assert all(url.startswith("https://run.googleapis.com/v2/projects/firstradequant/") for url in urls) + + +@pytest.mark.parametrize( + ("env", "configured", "source"), + [ + ([{"name": "FIRSTRADE_ACCOUNT", "valueSource": {"secretKeyRef": {"secret": SENTINEL}}}], True, "secret_reference"), + ([{"name": "OTHER_SETTING", "value": SENTINEL}], False, "absent"), + ([{"name": "FIRSTRADE_ACCOUNT", "value": " "}], False, "absent"), + ], +) +def test_selector_reports_only_presence_and_source( + env: list[dict[str, Any]], configured: bool, source: str +) -> None: + request, _ = _runner(revision=_revision(env=env)) + + result = readiness.inspect_readiness( + "firstrade-platform", "us-central1", "token", request_json=request + ) + + assert result["selector_configured"] is configured + assert result["selector_source"] == source + assert SENTINEL not in json.dumps(result) + + +@pytest.mark.parametrize( + "traffic", + [ + [], + [{"revision": "rev-a", "percent": 0}], + [{"revision": "rev-a", "percent": 50}, {"revision": "rev-b", "percent": 50}], + [{"revision": "rev-a", "percent": 99}], + [{"revision": "rev-a", "percent": True}], + ], +) +def test_rejects_missing_ambiguous_or_invalid_actual_traffic(traffic: list[dict[str, Any]]) -> None: + service = _service() + service["trafficStatuses"] = traffic + request, _ = _runner([service, service]) + + with pytest.raises(readiness.DiagnosticFailure): + readiness.inspect_readiness( + "firstrade-platform", "us-central1", "token", request_json=request + ) + + +def test_accepts_protojson_zero_percent_default_without_guessing_revision() -> None: + service = _service() + service.pop("reconciling") + service["trafficStatuses"] = [ + {"revision": "firstrade-platform-zero", "tag": "preview"}, + {"revision": REVISION, "percent": 100}, + ] + revision = _revision() + revision.pop("reconciling") + request, _ = _runner([service, service], revision=revision) + + result = readiness.inspect_readiness( + "firstrade-platform", "us-central1", "token", request_json=request + ) + + assert result["status"] == "verified" + + +def test_rejects_service_change_during_metadata_read() -> None: + before = _service() + after = _service() + after["etag"] = "etag-two" + request, _ = _runner([before, after]) + + with pytest.raises(readiness.DiagnosticFailure, match="service_changed_during_read"): + readiness.inspect_readiness( + "firstrade-platform", "us-central1", "token", request_json=request + ) + + +def test_rejects_invalid_resource_parts_without_request() -> None: + request, urls = _runner() + + with pytest.raises(readiness.DiagnosticFailure): + readiness.inspect_readiness("service/other", "us-central1", "token", request_json=request) + + assert urls == [] + + +def test_http_reader_uses_only_get_to_fixed_google_api_and_preserves_status(monkeypatch: pytest.MonkeyPatch) -> None: + class FakeResponse: + def __enter__(self) -> Self: + return self + + def __exit__(self, *args: object) -> None: + return None + + def read(self, size: int) -> bytes: + assert size == readiness.MAX_RESPONSE_BYTES + 1 + return b'{"name":"ok"}' + + class FakeOpener: + def open(self, request: Any, timeout: int) -> FakeResponse: + assert request.method == "GET" + assert request.full_url.startswith("https://run.googleapis.com/v2/projects/firstradequant/") + assert "SENTINEL" not in request.full_url + assert timeout == 15 + return FakeResponse() + + monkeypatch.setattr(readiness.urllib.request, "build_opener", lambda *handlers: FakeOpener()) + assert readiness._request_json(f"{readiness.API_ROOT}/{PROJECT}", "token") == {"name": "ok"} + + def fail_open(self, request: Any, timeout: int) -> None: + raise urllib.error.HTTPError(request.full_url, 403, SENTINEL, {}, None) + + monkeypatch.setattr(FakeOpener, "open", fail_open) + with pytest.raises(readiness.DiagnosticFailure) as caught: + readiness._request_json(f"{readiness.API_ROOT}/{PROJECT}", "token") + assert caught.value.http_status == 403 + assert SENTINEL not in str(caught.value) + + +@pytest.mark.parametrize( + "url", + [ + "https://run.googleapis.com/v2/projects/other/locations/us-central1/services/firstrade-platform", + "https://example.test/v2/projects/firstradequant/locations/us-central1/services/firstrade-platform", + f"{readiness.API_ROOT}/{PROJECT}?alt=json", + f"{readiness.API_ROOT}/projects/firstradequant/locations/us-central1/services/firstrade-platform/revisions/", + ], +) +def test_http_reader_rejects_non_target_uris_before_network( + monkeypatch: pytest.MonkeyPatch, url: str +) -> None: + def no_request(*args: Any, **kwargs: Any) -> None: + raise AssertionError("unexpected network request") + + monkeypatch.setattr(readiness.urllib.request, "build_opener", no_request) + with pytest.raises(readiness.DiagnosticFailure, match="request_target_rejected"): + readiness._request_json(url, "token") + + +def test_http_reader_rejects_oversized_metadata_response(monkeypatch: pytest.MonkeyPatch) -> None: + class LargeResponse: + def __enter__(self) -> Self: + return self + + def __exit__(self, *args: object) -> None: + return None + + def read(self, size: int) -> bytes: + assert size == readiness.MAX_RESPONSE_BYTES + 1 + return b"x" * size + + class FakeOpener: + def open(self, request: Any, timeout: int) -> LargeResponse: + return LargeResponse() + + monkeypatch.setattr(readiness.urllib.request, "build_opener", lambda *handlers: FakeOpener()) + with pytest.raises(readiness.DiagnosticFailure, match="metadata_response_too_large"): + readiness._request_json(f"{readiness.API_ROOT}/{PROJECT}", "token") + + +def test_workflow_metadata_path_is_opt_in_and_isolated() -> None: + workflow = Path(readiness.__file__).resolve().parents[1] / ".github/workflows/runtime-target-lifecycle.yml" + content = workflow.read_text(encoding="utf-8") + + assert "metadata_only:" in content + assert "default: false" in content + assert "github.event_name == 'workflow_dispatch' && inputs.metadata_only" in content + assert "account_data_readiness:" in content + assert "scripts/inspect_account_data_readiness.py" in content + assert "workflow_run:" in content and 'cron: "37 * * * *"' in content + assert "uv sync --frozen --no-dev" in content + assert content.count("uv sync --frozen --no-dev") == 1 + assert "probe" not in content.split("account_data_readiness:", 1)[1].lower() diff --git a/tests/test_runtime_monitor_workflows.py b/tests/test_runtime_monitor_workflows.py index 75b29ee..4a7b6a8 100644 --- a/tests/test_runtime_monitor_workflows.py +++ b/tests/test_runtime_monitor_workflows.py @@ -1,7 +1,6 @@ import re from pathlib import Path - ROOT = Path(__file__).resolve().parents[1] @@ -107,6 +106,22 @@ def test_lifecycle_observes_completed_sync_regardless_of_conclusion() -> None: assert "github.event.workflow_run.head_sha" not in workflow +def test_metadata_only_dispatch_is_opt_in_and_skips_lifecycle_job() -> None: + workflow = (ROOT / ".github/workflows/runtime-target-lifecycle.yml").read_text() + lifecycle = workflow.split(" lifecycle:", 1)[1].split(" account_data_readiness:", 1)[0] + metadata_job = workflow.split(" account_data_readiness:", 1)[1] + + assert "metadata_only:" in workflow + assert "type: boolean" in workflow + assert "default: false" in workflow + assert "if: ${{ !(github.event_name == 'workflow_dispatch' && inputs.metadata_only) }}" in lifecycle + assert "if: ${{ github.event_name == 'workflow_dispatch' && inputs.metadata_only }}" in metadata_job + assert "uv sync --frozen --no-dev" not in metadata_job + assert "scripts/inspect_account_data_readiness.py" in metadata_job + assert "CLOUD_RUN_SERVICE: ${{ secrets.CLOUD_RUN_SERVICE }}" in metadata_job + assert "CLOUD_RUN_REGION: ${{ vars.CLOUD_RUN_REGION }}" in metadata_job + + def test_lifecycle_publishes_read_only_observation_for_exact_service() -> None: workflow = (ROOT / ".github/workflows/runtime-target-lifecycle.yml").read_text() publisher = workflow.split("- name: Publish lifecycle to the unified control plane", 1)[1]