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
11 changes: 7 additions & 4 deletions .github/workflows/execution-report-heartbeat.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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:
Expand Down Expand Up @@ -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:
Expand All @@ -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
Expand All @@ -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:
Expand Down
110 changes: 96 additions & 14 deletions scripts/inspect_paused_account_facts_state.py
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand All @@ -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
Expand Down Expand Up @@ -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"]
Expand All @@ -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:
Expand All @@ -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
Expand Down
148 changes: 147 additions & 1 deletion tests/test_paused_account_facts_state.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand All @@ -127,14 +158,19 @@ 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)
assert set(result) == {
"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", [
Expand Down Expand Up @@ -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
Loading