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
29 changes: 23 additions & 6 deletions .github/workflows/execution-report-heartbeat.yml
Original file line number Diff line number Diff line change
Expand Up @@ -22,34 +22,34 @@ on:
type: boolean
default: false
account_facts_target:
description: "Publish one fresh archived account-facts report for the selected target."
description: "Optional account-facts action; disabled by default."
required: false
type: choice
default: disabled
options:
- disabled
- live-u16608560
- ingress-diagnostic
schedule:
- cron: "20 22 * * *"

env:
GCP_PROJECT_ID: interactivebrokersquant
GCP_WORKLOAD_IDENTITY_PROVIDER: projects/303168642265/locations/global/workloadIdentityPools/github-actions/providers/github-main
GCP_WORKLOAD_IDENTITY_SERVICE_ACCOUNT: ibkr-platform-deploy@interactivebrokersquant.iam.gserviceaccount.com

concurrency:
group: ${{ github.workflow }}-${{ github.ref_name }}
cancel-in-progress: false

jobs:
heartbeat:
name: Check execution report heartbeat
if: ${{ github.event_name != 'workflow_dispatch' || inputs.account_facts_target != 'ingress-diagnostic' }}
runs-on: ubuntu-latest
timeout-minutes: 15
permissions:
contents: read
id-token: write
env:
GCP_PROJECT_ID: interactivebrokersquant
GCP_WORKLOAD_IDENTITY_PROVIDER: projects/303168642265/locations/global/workloadIdentityPools/github-actions/providers/github-main
GCP_WORKLOAD_IDENTITY_SERVICE_ACCOUNT: ibkr-platform-deploy@interactivebrokersquant.iam.gserviceaccount.com
RUNTIME_HEARTBEAT_NAME: InteractiveBrokersPlatform
RUNTIME_HEARTBEAT_REPORT_PLATFORM: interactive_brokers
RUNTIME_HEARTBEAT_REQUIRED_SERVICES: ${{ secrets.RUNTIME_HEARTBEAT_REQUIRED_SERVICES }}
Expand Down Expand Up @@ -135,3 +135,20 @@ jobs:
- name: Send configured daily dry-run digest
if: ${{ always() && (github.event_name == 'schedule' || inputs.send_daily_dry_run_digest) }}
run: uv run --no-sync python scripts/daily_dry_run_digest.py

account-facts-ingress-diagnostic:
name: Verify account-facts ingress authentication only
if: ${{ github.event_name == 'workflow_dispatch' && inputs.account_facts_target == 'ingress-diagnostic' }}
runs-on: ubuntu-latest
timeout-minutes: 2
permissions:
contents: read
env:
IBKR_ACCOUNT_FACTS_TARGET: ${{ inputs.account_facts_target }}
IBKR_ACCOUNT_FACTS_SYNC_TOKEN: ${{ secrets.IBKR_ACCOUNT_FACTS_SYNC_TOKEN }}
steps:
- name: Checkout repository
uses: actions/checkout@v6

- name: Send one fixed account-facts ingress diagnostic
run: python3 scripts/publish_account_facts_from_report.py
249 changes: 242 additions & 7 deletions scripts/publish_account_facts_from_report.py
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,9 @@
SOURCE_BINDING_KIND = "deployment_runtime_account"
IBKR_ACCOUNT_FACTS_SYNC_TOKEN_ENV = "IBKR_ACCOUNT_FACTS_SYNC_TOKEN"
IBKR_ACCOUNT_FACTS_SYNC_URL = "https://qsl-strategy-switch-console.pigbibi.workers.dev/api/account-facts/sync"
IBKR_ACCOUNT_FACTS_INGRESS_DIAGNOSTIC_TARGET = "ingress-diagnostic"
_INGRESS_DIAGNOSTIC_BODY = b"{}"
_INGRESS_DIAGNOSTIC_ERROR = "invalid_account_facts_history"
_REPORT_RUN_ID = re.compile(r"^\d{8}T\d{6}Z\.json$")
_ACCOUNT_ID = re.compile(r"^(?:U|DU)\d+$")
_ALLOWED_CASH_TAGS = frozenset(
Expand All @@ -35,6 +38,41 @@
_CURRENCY = re.compile(r"^[A-Z]{3}$")
_DECIMAL_TEXT = re.compile(r"^-?(?:0|[1-9]\d*)(?:\.\d+)?$")
_TARGET_ID = re.compile(r"^[a-z0-9](?:[a-z0-9-]{0,62}[a-z0-9])?$")
_QRS_ACCOUNT_FACTS_ERROR_CODES = frozenset(
{
"account_facts_account_unattributed",
"account_facts_bindings_missing",
"account_facts_caller_identity_forbidden",
"account_facts_disabled",
"account_facts_identity_mismatch",
"account_facts_store_unavailable",
"account_facts_sync_token_ambiguous",
"account_facts_sync_token_invalid",
"account_facts_sync_token_not_configured",
"account_facts_sync_token_platform_mismatch",
"account_facts_target_mismatch",
"account_facts_binding_mismatch",
"account_facts_observation_window",
"account_facts_future_observation",
"invalid_account_facts_account",
"invalid_account_facts_balances",
"invalid_account_facts_bindings",
"invalid_account_facts_cash",
"invalid_account_facts_date",
"invalid_account_facts_history",
"invalid_account_facts_money_magnitude",
"invalid_account_facts_money_scale",
"invalid_account_facts_scope",
"invalid_account_facts_source_binding",
"invalid_account_facts_target",
"invalid_account_facts_time",
"account_facts_stored_invalid",
"duplicate_account_facts_binding",
"invalid_account_facts_stored",
"unsupported_account_facts_action",
}
)
_HTTP_ERROR_BODY_LIMIT = 16 * 1024


class _ProjectionError(ValueError):
Expand Down Expand Up @@ -253,14 +291,201 @@ def publish_ibkr_account_facts_history(
method="POST",
)
opener = build_opener(_NoRedirect())
with opener.open(request, timeout=15) as response:
if not 200 <= response.status < 300:
raise _ProjectionError("publish_failed")
try:
with opener.open(request, timeout=15) as response:
if not 200 <= response.status < 300:
return _publish_failed(
stage="http_response",
category="http_status",
http_status=_numeric_http_status(response.status),
qrs_error_code="unknown",
)
except HTTPError as exc:
return _publish_failed(
stage="http_response",
category="http_error",
http_status=_numeric_http_status(exc.code),
qrs_error_code=_safe_qrs_error_code(exc),
)
except TimeoutError:
return _publish_failed(stage="request", category="timeout")
except URLError:
return _publish_failed(stage="request", category="url_error")
except OSError:
return _publish_failed(stage="request", category="transport_error")
except Exception:
return _publish_failed(stage="unknown", category="unknown")
return {"status": "published"}
except _ProjectionError as exc:
return {"status": "skipped", "reason": exc.reason}
except (HTTPError, URLError, TimeoutError, OSError):
return {"status": "skipped", "reason": "publish_failed"}


def _numeric_http_status(value: object) -> int | None:
if isinstance(value, int) and not isinstance(value, bool) and 100 <= value <= 599:
return value
return None


def _safe_qrs_error_code(error: HTTPError) -> str:
try:
body = error.read(_HTTP_ERROR_BODY_LIMIT + 1)
except Exception:
return "unknown"
return _safe_qrs_error_code_from_body(body)


def _safe_qrs_error_code_from_body(body: object) -> str:
if not isinstance(body, bytes) or len(body) > _HTTP_ERROR_BODY_LIMIT:
return "unknown"
try:
payload = json.loads(body)
except Exception:
return "unknown"
code = payload.get("error") if isinstance(payload, Mapping) else None
return code if isinstance(code, str) and code in _QRS_ACCOUNT_FACTS_ERROR_CODES else "unknown"


def _publish_failed(
*,
stage: str,
category: str,
http_status: int | None = None,
qrs_error_code: str = "unknown",
) -> dict[str, Any]:
return {
"status": "skipped",
"reason": "publish_failed",
"diagnostics": {
"stage": stage,
"category": category,
"http_status": http_status,
"qrs_error_code": qrs_error_code,
"outcome": "unknown",
},
}


def _ingress_diagnostic_result(
*,
http_status: int | None,
qrs_error_code: str,
category: str = "http_error",
stage: str = "http_response",
) -> dict[str, Any]:
verified = http_status == 400 and qrs_error_code == _INGRESS_DIAGNOSTIC_ERROR
return {
"status": "verified" if verified else "skipped",
"reason": (
"ingress_authentication_and_schema_rejection_verified"
if verified
else "ingress_diagnostic_unverified"
),
"diagnostics": {
"stage": stage,
"category": category,
"http_status": http_status,
"qrs_error_code": qrs_error_code,
"outcome": "rejected_before_storage" if verified else "unknown",
},
}


def diagnose_account_facts_ingress(*, sync_token: str) -> dict[str, Any]:
"""Send one fixed empty schema probe to verify protected QRS ingress only."""
if not _text(sync_token):
return {"status": "skipped", "reason": "publish_auth_unavailable"}
request = Request(
IBKR_ACCOUNT_FACTS_SYNC_URL,
data=_INGRESS_DIAGNOSTIC_BODY,
headers={
"Authorization": f"Bearer {_text(sync_token)}",
"Content-Type": "application/json",
},
method="POST",
)
try:
opener = build_opener(_NoRedirect())
with opener.open(request, timeout=15) as response:
status = _numeric_http_status(getattr(response, "status", None))
if status != 400:
return _ingress_diagnostic_result(
http_status=status,
qrs_error_code="unknown",
category="http_status",
)
try:
body = response.read(_HTTP_ERROR_BODY_LIMIT + 1)
except Exception:
body = None
return _ingress_diagnostic_result(
http_status=status,
qrs_error_code=_safe_qrs_error_code_from_body(body),
category="http_status",
)
except HTTPError as exc:
status = _numeric_http_status(exc.code)
code = _safe_qrs_error_code(exc)
return _ingress_diagnostic_result(http_status=status, qrs_error_code=code)
except TimeoutError:
return _ingress_diagnostic_result(
http_status=None,
qrs_error_code="unknown",
category="timeout",
stage="request",
)
except URLError:
return _ingress_diagnostic_result(
http_status=None,
qrs_error_code="unknown",
category="url_error",
stage="request",
)
except OSError:
return _ingress_diagnostic_result(
http_status=None,
qrs_error_code="unknown",
category="transport_error",
stage="request",
)
except Exception:
return _ingress_diagnostic_result(
http_status=None,
qrs_error_code="unknown",
category="unknown",
stage="unknown",
)


def _format_cli_result(result: Mapping[str, Any]) -> str:
status = result.get("status") if result.get("status") in {"published", "verified", "skipped"} else "skipped"
reason = result.get("reason") if isinstance(result.get("reason"), str) else "unknown"
diagnostic = result.get("diagnostics")
if not isinstance(diagnostic, Mapping):
return f"{status}:{reason}"
stage = (
diagnostic.get("stage")
if diagnostic.get("stage") in {"http_response", "request", "unknown"}
else "unknown"
)
category = (
diagnostic.get("category")
if diagnostic.get("category") in {
"http_status", "http_error", "timeout", "url_error", "transport_error", "unknown"
}
else "unknown"
)
http_status = _numeric_http_status(diagnostic.get("http_status"))
qrs_error_code = diagnostic.get("qrs_error_code")
if not isinstance(qrs_error_code, str) or qrs_error_code not in _QRS_ACCOUNT_FACTS_ERROR_CODES:
qrs_error_code = "unknown"
outcome = diagnostic.get("outcome")
if outcome not in {"rejected_before_storage", "unknown"}:
outcome = "unknown"
return (
f"{status}:{reason}:stage={stage}:category={category}"
f":http_status={http_status if http_status is not None else 'none'}"
f":qrs_error_code={qrs_error_code}:outcome={outcome}"
)


class _NoRedirect(HTTPRedirectHandler):
Expand Down Expand Up @@ -322,6 +547,12 @@ def main() -> int:
"""Workflow-only, explicitly enabled single-report publisher."""
try:
target = _text(os.environ.get("IBKR_ACCOUNT_FACTS_TARGET"))
if target == IBKR_ACCOUNT_FACTS_INGRESS_DIAGNOSTIC_TARGET:
result = diagnose_account_facts_ingress(
sync_token=os.environ.get(IBKR_ACCOUNT_FACTS_SYNC_TOKEN_ENV, "")
)
print(_format_cli_result(result))
return 0 if result.get("status") == "verified" else 1
if target != "live-u16608560":
print("skipped:target_disabled")
return 0
Expand Down Expand Up @@ -357,7 +588,7 @@ def main() -> int:
sync_token=os.environ.get(IBKR_ACCOUNT_FACTS_SYNC_TOKEN_ENV, ""),
**expected,
)
print(f"{result['status']}:{result.get('reason', 'ok')}")
print(_format_cli_result(result))
return 0 if result["status"] == "published" else 1
except _ProjectionError as exc:
print(f"skipped:{exc.reason}")
Expand Down Expand Up @@ -470,7 +701,11 @@ def _project_ibkr_account_facts_history(
}


__all__ = ["project_ibkr_account_facts_history", "publish_ibkr_account_facts_history"]
__all__ = [
"diagnose_account_facts_ingress",
"project_ibkr_account_facts_history",
"publish_ibkr_account_facts_history",
]


if __name__ == "__main__":
Expand Down
Loading
Loading