diff --git a/.github/workflows/execution-report-heartbeat.yml b/.github/workflows/execution-report-heartbeat.yml index bd07f4f..8c571be 100644 --- a/.github/workflows/execution-report-heartbeat.yml +++ b/.github/workflows/execution-report-heartbeat.yml @@ -34,6 +34,7 @@ on: - additional-2 - additional-3 - paused-refresh-preflight + - paused-refresh-state-backup account_facts_report_name: description: "Optional exact UTC runtime report filename for account-facts publishing." required: false @@ -49,7 +50,7 @@ concurrency: jobs: heartbeat: name: Check execution report heartbeat - if: ${{ github.event_name != 'workflow_dispatch' || (inputs.account_facts_target != 'ingress-diagnostic' && inputs.account_facts_target != 'paused-refresh-preflight' && ((inputs.account_facts_target != 'primary-live' && inputs.account_facts_target != 'additional-1' && inputs.account_facts_target != 'additional-2' && inputs.account_facts_target != 'additional-3') || inputs.account_facts_report_name == '')) }} + if: ${{ github.event_name != 'workflow_dispatch' || (inputs.account_facts_target != 'ingress-diagnostic' && inputs.account_facts_target != 'paused-refresh-preflight' && inputs.account_facts_target != 'paused-refresh-state-backup' && ((inputs.account_facts_target != 'primary-live' && inputs.account_facts_target != 'additional-1' && inputs.account_facts_target != 'additional-2' && inputs.account_facts_target != 'additional-3') || inputs.account_facts_report_name == '')) }} runs-on: ubuntu-latest timeout-minutes: 15 permissions: @@ -192,8 +193,8 @@ jobs: IBKR_ACCOUNT_FACTS_ADDITIONAL_TARGETS_JSON: ${{ (inputs.account_facts_target == 'additional-1' || inputs.account_facts_target == 'additional-2' || inputs.account_facts_target == 'additional-3') && secrets.IBKR_ACCOUNT_FACTS_ADDITIONAL_TARGETS_JSON || '' }} paused-refresh-preflight: - name: Inspect archived paused IBKR account-facts state - if: ${{ github.event_name == 'workflow_dispatch' && github.ref == 'refs/heads/main' && inputs.account_facts_target == 'paused-refresh-preflight' }} + name: Inspect or back up archived paused IBKR account-facts state + if: ${{ github.event_name == 'workflow_dispatch' && github.ref == 'refs/heads/main' && (inputs.account_facts_target == 'paused-refresh-preflight' || inputs.account_facts_target == 'paused-refresh-state-backup') }} runs-on: ubuntu-latest timeout-minutes: 5 permissions: @@ -213,6 +214,7 @@ jobs: IBKR_ACCOUNT_FACTS_ACCOUNT_SELECTOR_JSON: ${{ secrets.IBKR_ACCOUNT_FACTS_ACCOUNT_SELECTOR_JSON }} IBKR_ACCOUNT_FACTS_DEPLOYMENT_SELECTOR: ${{ secrets.IBKR_ACCOUNT_FACTS_DEPLOYMENT_SELECTOR }} IBKR_PAUSED_FACTS_STATE_TRANSFER_URI: ${{ secrets.IBKR_PAUSED_FACTS_STATE_TRANSFER_URI }} + IBKR_PAUSED_FACTS_STATE_ACTION: ${{ inputs.account_facts_target == 'paused-refresh-state-backup' && 'backup' || 'inspect' }} steps: - name: Checkout repository uses: actions/checkout@v6 @@ -236,9 +238,10 @@ jobs: workload_identity_provider: ${{ env.GCP_WORKLOAD_IDENTITY_PROVIDER }} service_account: ${{ env.GCP_WORKLOAD_IDENTITY_SERVICE_ACCOUNT }} - - name: Inspect transferred state and effective permissions + - name: Inspect transferred state or back up the handover object env: IBKR_ACCOUNT_FACTS_ADDITIONAL_TARGETS_JSON: ${{ secrets.IBKR_ACCOUNT_FACTS_ADDITIONAL_TARGETS_JSON }} + IBKR_PAUSED_FACTS_STATE_PREFIX: ${{ inputs.account_facts_target == 'paused-refresh-state-backup' && secrets.IBKR_PAUSED_FACTS_STATE_PREFIX || '' }} run: uv run --no-sync python scripts/inspect_paused_account_facts_state.py account-facts-ingress-diagnostic: diff --git a/scripts/inspect_paused_account_facts_state.py b/scripts/inspect_paused_account_facts_state.py index d990948..2567d22 100644 --- a/scripts/inspect_paused_account_facts_state.py +++ b/scripts/inspect_paused_account_facts_state.py @@ -109,7 +109,7 @@ def _bounded_content(response: object, limit: int) -> bytes: return b"".join(chunks) -def _read_transfer(raw_uri: str) -> dict[str, object]: +def _read_transfer(raw_uri: str) -> tuple[bytes, dict[str, object]]: bucket, obj = _gs_uri(raw_uri) response = _request( "GET", f"https://storage.googleapis.com/storage/v1/b/{quote(bucket, safe='')}/o/{quote(obj, safe='')}?alt=media", @@ -122,19 +122,36 @@ def _read_transfer(raw_uri: str) -> dict[str, object]: response.close() if not isinstance(parsed, dict): raise PreflightError("invalid_object") - return parsed - - -def _request(method: str, url: str, *, expected: int, json_body: object | None = None): + return body, parsed + + +def _request( + method: str, + url: str, + *, + expected: int | tuple[int, ...], + json_body: object | None = None, + data: bytes | None = None, + headers: dict[str, str] | None = None, +): session = _SESSION + request_kwargs: dict[str, object] = { + "timeout": TIMEOUT_SECONDS, + "allow_redirects": False, + "stream": True, + } + if json_body is not None: + request_kwargs["json"] = json_body + if data is not None: + request_kwargs["data"] = data + if headers is not None: + request_kwargs["headers"] = headers try: - response = session.request( - method, url, json=json_body, timeout=TIMEOUT_SECONDS, - allow_redirects=False, stream=True, - ) + response = session.request(method, url, **request_kwargs) except Exception: raise PreflightError("request_failed") from None - if response.status_code != expected or 300 <= response.status_code < 400: + accepted = (expected,) if isinstance(expected, int) else expected + if response.status_code not in accepted or 300 <= response.status_code < 400: response.close() raise PreflightError("http_failed") return response @@ -294,7 +311,7 @@ def _bucket_permissions(bucket: str, permissions: list[str]) -> dict[str, bool]: def inspect(uri: str, protected: dict[str, str]) -> dict[str, object]: transfer_bucket, _ = _gs_uri(uri) - payload = _read_transfer(uri) + _, payload = _read_transfer(uri) target, decoded = _verify_payload(payload, protected) del decoded project = target["project"] @@ -320,6 +337,62 @@ def inspect(uri: str, protected: dict[str, str]) -> dict[str, object]: return results +def _state_backup_object(prefix_uri: str, transfer_bucket: str) -> tuple[str, str]: + try: + bucket, path = _gs_uri(prefix_uri) + except PreflightError: + raise PreflightError("invalid_backup_prefix") from None + parsed = urlsplit(prefix_uri) + if bucket != transfer_bucket or not parsed.path.endswith("/") or not path.rstrip("/"): + raise PreflightError("invalid_backup_prefix") + object_name = f"{path}archived-handover.json" + return bucket, object_name + + +def backup_handover_state(uri: str, prefix_uri: str, protected: dict[str, str]) -> dict[str, str]: + transfer_bucket, _ = _gs_uri(uri) + bucket, object_name = _state_backup_object(prefix_uri, transfer_bucket) + object_bytes, payload = _read_transfer(uri) + _verify_payload(payload, protected) + encoded_bucket = quote(bucket, safe="") + encoded_object = quote(object_name, safe="") + metadata_url = f"https://storage.googleapis.com/storage/v1/b/{encoded_bucket}/o/{encoded_object}" + try: + existing = _request("GET", metadata_url, expected=(200, 404)) + existed = existing.status_code == 200 + existing.close() + except PreflightError: + return {"state": "blocked", "phase": "backup_check", "reason": "backup_check_failed"} + except Exception: + return {"state": "blocked", "phase": "backup_check", "reason": "backup_check_failed"} + if existed: + return {"state": "blocked", "phase": "backup_check", "reason": "backup_already_exists"} + + upload_query = urlencode({"uploadType": "media", "name": object_name, "ifGenerationMatch": "0"}) + upload_url = f"https://storage.googleapis.com/upload/storage/v1/b/{encoded_bucket}/o?{upload_query}" + try: + uploaded = _request( + "POST", upload_url, expected=(200, 201), data=object_bytes, + headers={"Content-Type": "application/json"}, + ) + uploaded.close() + except Exception: + return {"state": "unknown", "phase": "upload", "reason": "upload_outcome_unknown"} + + readback_url = f"{metadata_url}?alt=media" + try: + readback = _request("GET", readback_url, expected=200) + try: + readback_bytes = _bounded_content(readback, MAX_OBJECT_BYTES) + finally: + readback.close() + except Exception: + return {"state": "unknown", "phase": "readback", "reason": "readback_unknown"} + if hashlib.sha256(readback_bytes).digest() != hashlib.sha256(object_bytes).digest(): + return {"state": "unknown", "phase": "readback", "reason": "readback_mismatch"} + return {"state": "backup_verified", "phase": "readback_verified", "reason": "none"} + + def main() -> int: global _SESSION try: @@ -331,13 +404,22 @@ def main() -> int: if protected.get("IBKR_ACCOUNT_FACTS_TARGET") != "additional-1": raise PreflightError("target_mismatch") uri = os.environ.get("IBKR_PAUSED_FACTS_STATE_TRANSFER_URI", "") - _gs_uri(uri) + action = os.environ.get("IBKR_PAUSED_FACTS_STATE_ACTION", "inspect") + if action not in {"inspect", "backup"}: + raise PreflightError("invalid_action") + if action == "backup" and os.environ.get("IBKR_ACCOUNT_FACTS_TARGET") != "additional-1": + raise PreflightError("target_mismatch") credentials, _ = google.auth.default(scopes=["https://www.googleapis.com/auth/cloud-platform"]) _SESSION = AuthorizedSession(credentials, max_refresh_attempts=0) _SESSION.mount("https://", HTTPAdapter(max_retries=0)) - result = inspect(uri, protected) + if action == "backup": + result = backup_handover_state( + uri, os.environ.get("IBKR_PAUSED_FACTS_STATE_PREFIX", ""), protected + ) + else: + result = inspect(uri, protected) print(json.dumps(result, sort_keys=True)) - return 0 + return 0 if result.get("state") in {"archived_state_verified", "backup_verified"} else 1 except PreflightError as exc: print(json.dumps({"state": "blocked", "reason": str(exc)}, sort_keys=True)) return 1 diff --git a/tests/test_paused_account_facts_state.py b/tests/test_paused_account_facts_state.py index 85f9225..387da1c 100644 --- a/tests/test_paused_account_facts_state.py +++ b/tests/test_paused_account_facts_state.py @@ -109,6 +109,37 @@ def request(self, method, url, **kwargs): return _Response({"permissions": kwargs["json"]["permissions"]}) +class _BackupSession: + def __init__(self, object_bytes, *, existing=False, upload_error=False, readback_error=False): + self.object_bytes = object_bytes + self.existing = existing + self.upload_error = upload_error + self.readback_error = readback_error + self.uploaded = None + self.calls = [] + + def request(self, method, url, **kwargs): + self.calls.append((method, url, kwargs)) + if method == "GET" and url.endswith("/o/state.json?alt=media"): + return _Response(self.object_bytes) + if method == "GET" and url.endswith("/o/private%2Fstate%2Farchived-handover.json"): + return _Response(b"{}", status=200 if self.existing else 404) + if method == "POST" and "/upload/storage/v1/b/" in url: + if self.upload_error: + raise TimeoutError("synthetic timeout") + self.uploaded = kwargs["data"] + return _Response({"kind": "storage#object"}) + if method == "GET" and url.endswith("/o/private%2Fstate%2Farchived-handover.json?alt=media"): + if self.readback_error: + return _Response(b"", status=503) + return _Response(self.uploaded) + raise AssertionError("unexpected request") + + +def _transfer_bytes(payload): + return json.dumps(payload).encode() + + def test_inspection_uses_only_expected_read_and_permission_calls(monkeypatch): payload, protected = _fixture() session = _Session(payload) @@ -127,7 +158,8 @@ def test_inspection_uses_only_expected_read_and_permission_calls(monkeypatch): assert all("fixture-account" not in url for _, url, _ in session.calls) assert session.calls[0][1] == "https://storage.googleapis.com/storage/v1/b/fixture-transfer/o/state.json?alt=media" assert session.calls[1][1] == "https://storage.googleapis.com/storage/v1/b/fixture-transfer/iam/testPermissions?permissions=storage.objects.get&permissions=storage.objects.create" - assert session.calls[1][2]["json"] is None + assert "json" not in session.calls[1][2] + assert "data" not in session.calls[1][2] assert session.calls[2][1] == "https://run.googleapis.com/v2/projects/fixture-project/locations/us-central1/services/fixture-service:testIamPermissions" assert session.calls[3][1] == "https://cloudresourcemanager.googleapis.com/v1/projects/fixture-project:testIamPermissions" assert all("run.app" not in url and "scheduler.googleapis.com" not in url for _, url, _ in session.calls) @@ -135,6 +167,10 @@ def test_inspection_uses_only_expected_read_and_permission_calls(monkeypatch): "state", "preserved_file_count", "absent_execution_state_count", "bucket_permissions", "service_permissions", "scheduler_permissions", } + assert all( + method == "GET" or url.endswith(":testIamPermissions") + for method, url, _ in session.calls + ) @pytest.mark.parametrize("uri", [ @@ -204,18 +240,128 @@ def request(self, method, url, **kwargs): inspector.inspect("gs://fixture-transfer/state.json", protected) +def test_backup_stops_when_exact_marker_already_exists(monkeypatch): + payload, protected = _fixture() + session = _BackupSession(_transfer_bytes(payload), existing=True) + monkeypatch.setattr(inspector, "_SESSION", session) + + result = inspector.backup_handover_state( + "gs://fixture-transfer/state.json", "gs://fixture-transfer/private/state/", protected + ) + + assert result == {"state": "blocked", "phase": "backup_check", "reason": "backup_already_exists"} + assert [method for method, _, _ in session.calls] == ["GET", "GET"] + assert not any(method == "POST" for method, _, _ in session.calls) + + +def test_backup_first_attempt_is_create_only_and_reads_back_original_bytes(monkeypatch): + payload, protected = _fixture() + original = _transfer_bytes(payload) + session = _BackupSession(original) + monkeypatch.setattr(inspector, "_SESSION", session) + + result = inspector.backup_handover_state( + "gs://fixture-transfer/state.json", "gs://fixture-transfer/private/state/", protected + ) + + assert result == {"state": "backup_verified", "phase": "readback_verified", "reason": "none"} + assert session.uploaded == original + assert [method for method, _, _ in session.calls] == ["GET", "GET", "POST", "GET"] + assert session.calls[1][1] == "https://storage.googleapis.com/storage/v1/b/fixture-transfer/o/private%2Fstate%2Farchived-handover.json" + upload = session.calls[2] + assert upload[1] == "https://storage.googleapis.com/upload/storage/v1/b/fixture-transfer/o?uploadType=media&name=private%2Fstate%2Farchived-handover.json&ifGenerationMatch=0" + assert upload[2]["data"] == original + assert upload[2]["headers"] == {"Content-Type": "application/json"} + assert upload[2]["allow_redirects"] is False + assert session.calls[3][1] == "https://storage.googleapis.com/storage/v1/b/fixture-transfer/o/private%2Fstate%2Farchived-handover.json?alt=media" + + +def test_backup_upload_outcome_unknown_is_not_retried(monkeypatch): + payload, protected = _fixture() + session = _BackupSession(_transfer_bytes(payload), upload_error=True) + monkeypatch.setattr(inspector, "_SESSION", session) + + result = inspector.backup_handover_state( + "gs://fixture-transfer/state.json", "gs://fixture-transfer/private/state/", protected + ) + + assert result == {"state": "unknown", "phase": "upload", "reason": "upload_outcome_unknown"} + assert [method for method, _, _ in session.calls] == ["GET", "GET", "POST"] + + +def test_backup_readback_error_remains_unknown_without_reupload(monkeypatch): + payload, protected = _fixture() + session = _BackupSession(_transfer_bytes(payload), readback_error=True) + monkeypatch.setattr(inspector, "_SESSION", session) + + result = inspector.backup_handover_state( + "gs://fixture-transfer/state.json", "gs://fixture-transfer/private/state/", protected + ) + + assert result == {"state": "unknown", "phase": "readback", "reason": "readback_unknown"} + assert [method for method, _, _ in session.calls] == ["GET", "GET", "POST", "GET"] + + +def test_backup_readback_mismatch_remains_unknown_without_reupload(monkeypatch): + payload, protected = _fixture() + session = _BackupSession(_transfer_bytes(payload)) + session.readback_error = False + session.readback_bytes = b"different bytes" + original_request = session.request + + def request(method, url, **kwargs): + response = original_request(method, url, **kwargs) + if method == "GET" and url.endswith("/o/private%2Fstate%2Farchived-handover.json?alt=media"): + return _Response(session.readback_bytes) + return response + + session.request = request + monkeypatch.setattr(inspector, "_SESSION", session) + + result = inspector.backup_handover_state( + "gs://fixture-transfer/state.json", "gs://fixture-transfer/private/state/", protected + ) + + assert result == {"state": "unknown", "phase": "readback", "reason": "readback_mismatch"} + assert [method for method, _, _ in session.calls] == ["GET", "GET", "POST", "GET"] + + +@pytest.mark.parametrize("prefix", [ + "gs://other-transfer/private/state/", "gs://fixture-transfer/private/state", + "gs://fixture-transfer/private/../state/", "gs://fixture-transfer/private/state/?generation=1", + "gs://fixture-transfer/", +]) +def test_backup_rejects_invalid_prefix_before_any_cloud_request(monkeypatch, prefix): + payload, protected = _fixture() + session = _BackupSession(_transfer_bytes(payload)) + monkeypatch.setattr(inspector, "_SESSION", session) + + with pytest.raises(inspector.PreflightError, match="invalid_backup_prefix"): + inspector.backup_handover_state("gs://fixture-transfer/state.json", prefix, protected) + assert session.calls == [] + + def test_action_workflow_is_manual_main_only_and_separate_from_existing_paths(): from pathlib import Path workflow = Path(".github/workflows/execution-report-heartbeat.yml").read_text() assert "default: disabled" in workflow assert "- paused-refresh-preflight" in workflow + assert "- paused-refresh-state-backup" in workflow assert "github.ref == 'refs/heads/main'" in workflow assert "inputs.account_facts_target == 'paused-refresh-preflight'" in workflow + assert "inputs.account_facts_target == 'paused-refresh-state-backup'" in workflow assert "IBKR_PAUSED_FACTS_STATE_TRANSFER_URI: ${{ secrets.IBKR_PAUSED_FACTS_STATE_TRANSFER_URI }}" in workflow + assert "IBKR_PAUSED_FACTS_STATE_PREFIX" in workflow + assert "secrets.IBKR_PAUSED_FACTS_STATE_PREFIX" in workflow assert "script='" not in workflow assert "inputs.account_facts_target != 'paused-refresh-preflight'" in workflow + assert "inputs.account_facts_target != 'paused-refresh-state-backup'" in workflow assert "Publish one validated report" in workflow assert "name: Check execution report heartbeat" in workflow publisher_condition = workflow.split(" account-facts-publisher:", 1)[1].split(" paused-refresh-preflight:", 1)[0] assert "paused-refresh-preflight" not in publisher_condition + assert "paused-refresh-state-backup" not in publisher_condition + ingress_condition = workflow.split(" account-facts-ingress-diagnostic:", 1)[1] + assert "inputs.account_facts_target == 'ingress-diagnostic'" in ingress_condition + assert "inputs.account_facts_target == 'paused-refresh-state-backup'" not in ingress_condition