diff --git a/.github/workflows/execution-report-heartbeat.yml b/.github/workflows/execution-report-heartbeat.yml index ff9d40b..621256f 100644 --- a/.github/workflows/execution-report-heartbeat.yml +++ b/.github/workflows/execution-report-heartbeat.yml @@ -122,6 +122,23 @@ jobs: - name: Set up gcloud uses: google-github-actions/setup-gcloud@v3 + - name: Record and sync daily paper account snapshot + id: account_history + if: ${{ !cancelled() && matrix.target.label == 'PAPER' && vars.ACCOUNT_HISTORY_RECORDING_ENABLED == 'true' }} + continue-on-error: true + env: + ACCOUNT_HISTORY_RECORDING_ENABLED: ${{ vars.ACCOUNT_HISTORY_RECORDING_ENABLED }} + ACCOUNT_HISTORY_SERVICE_URL: ${{ vars.ACCOUNT_HISTORY_SERVICE_URL }} + ACCOUNT_HISTORY_GCS_PREFIX: ${{ vars.ACCOUNT_HISTORY_GCS_PREFIX }} + ACCOUNT_HISTORY_TARGET_ID: ${{ matrix.target.id }} + ACCOUNT_HISTORY_EXPECTED_SCOPE: ${{ matrix.target.label }} + ACCOUNT_HISTORY_EXPECTED_SOURCE_BINDING_ID: ${{ vars.ACCOUNT_HISTORY_EXPECTED_SOURCE_BINDING_ID }} + GOOGLE_CLOUD_PROJECT: ${{ env.GCP_PROJECT_ID }} + ACCOUNT_FACTS_SYNC_ENABLED: ${{ vars.ACCOUNT_FACTS_SYNC_ENABLED }} + ACCOUNT_FACTS_SYNC_URL: ${{ vars.ACCOUNT_FACTS_SYNC_URL }} + ACCOUNT_FACTS_SYNC_TOKEN: ${{ vars.ACCOUNT_FACTS_SYNC_ENABLED == 'true' && secrets.ACCOUNT_FACTS_SYNC_TOKEN || '' }} + run: uv run --no-sync python scripts/record_daily_account_snapshot.py + - name: Check recent execution report run: uv run --no-sync python scripts/execution_report_heartbeat.py @@ -136,20 +153,6 @@ jobs: env: EXECUTION_EVIDENCE_SYNC_TOKEN: ${{ secrets.EXECUTION_EVIDENCE_SYNC_TOKEN }} - - name: Record daily paper account snapshot - if: ${{ success() && matrix.target.label == 'PAPER' && vars.ACCOUNT_HISTORY_RECORDING_ENABLED == 'true' }} - env: - ACCOUNT_HISTORY_RECORDING_ENABLED: ${{ vars.ACCOUNT_HISTORY_RECORDING_ENABLED }} - ACCOUNT_HISTORY_SERVICE_URL: ${{ vars.ACCOUNT_HISTORY_SERVICE_URL }} - ACCOUNT_HISTORY_GCS_PREFIX: ${{ vars.ACCOUNT_HISTORY_GCS_PREFIX }} - ACCOUNT_HISTORY_TARGET_ID: ${{ matrix.target.id }} - ACCOUNT_HISTORY_EXPECTED_SCOPE: ${{ matrix.target.label }} - GOOGLE_CLOUD_PROJECT: ${{ env.GCP_PROJECT_ID }} - ACCOUNT_FACTS_SYNC_ENABLED: ${{ vars.ACCOUNT_FACTS_SYNC_ENABLED }} - ACCOUNT_FACTS_SYNC_URL: ${{ vars.ACCOUNT_FACTS_SYNC_URL }} - ACCOUNT_FACTS_SYNC_TOKEN: ${{ vars.ACCOUNT_FACTS_SYNC_ENABLED == 'true' && secrets.ACCOUNT_FACTS_SYNC_TOKEN || '' }} - run: uv run --no-sync python scripts/record_daily_account_snapshot.py - - name: Publish daily runtime projection if: ${{ !cancelled() && matrix.target.label == 'PAPER' && vars.RUNTIME_DAILY_PROJECTION_ENABLED == 'true' && (steps.gcp_auth_primary.outcome == 'success' || steps.gcp_auth_retry.outcome == 'success') }} env: @@ -161,3 +164,7 @@ jobs: EXECUTION_EVIDENCE_SYNC_TOKEN: ${{ vars.RUNTIME_DAILY_SYNC_ENABLED == 'true' && secrets.EXECUTION_EVIDENCE_SYNC_TOKEN || '' }} GOOGLE_CLOUD_PROJECT: ${{ env.GCP_PROJECT_ID }} run: uv run --no-sync python scripts/publish_daily_runtime_projection.py + + - name: Fail after completing heartbeat checks if account snapshot sync failed + if: ${{ !cancelled() && steps.account_history.outcome == 'failure' }} + run: exit 1 diff --git a/.github/workflows/sync-cloud-run-env.yml b/.github/workflows/sync-cloud-run-env.yml index a375c90..c117b35 100644 --- a/.github/workflows/sync-cloud-run-env.yml +++ b/.github/workflows/sync-cloud-run-env.yml @@ -204,6 +204,7 @@ jobs: # Trusted control code stays on this main checkout. The application archive is chosen below. approved_candidate=0b939723c1db3ef59175535998b470cbcd4b8824 approved_http_snapshot_candidate=d8314a61df697cae1dd03a78ddc5c2fc4179ec67 + approved_probe_snapshot_candidate=a2921d157efb887e9210fad6734ca040ea6e5293 history_count=0 for history_value in \ "${ACCOUNT_HISTORY_RECORDING_ENABLED:-}" \ @@ -252,6 +253,10 @@ jobs: elif [ "${WORKFLOW_TARGET}" = "PAPER" ] && [ "${history_count}" -eq 0 ] \ && [ "${SOURCE_COMMIT}" = "${approved_http_snapshot_candidate}" ]; then image_mode=snapshot + elif [ "${WORKFLOW_TARGET}" = "PAPER" ] && [ "${history_count}" -eq 0 ] \ + && [ -z "${snapshot_setting}" ] \ + && [ "${SOURCE_COMMIT}" = "${approved_probe_snapshot_candidate}" ]; then + image_mode=probe else echo "Image source is not approved." >&2 exit 1 @@ -305,6 +310,7 @@ jobs: set -euo pipefail approved_candidate=0b939723c1db3ef59175535998b470cbcd4b8824 approved_http_snapshot_candidate=d8314a61df697cae1dd03a78ddc5c2fc4179ec67 + approved_probe_snapshot_candidate=a2921d157efb887e9210fad6734ca040ea6e5293 history_count=0 for history_value in \ "${ACCOUNT_HISTORY_RECORDING_ENABLED:-}" \ @@ -316,6 +322,7 @@ jobs: history_count=$((history_count + 1)) fi done + snapshot_setting="${ACCOUNT_SNAPSHOT_ENABLED_INPUT:-}" if [ "${SOURCE_COMMIT}" = "${GITHUB_SHA}" ] && [ "${history_count}" -eq 0 ]; then image_mode=same elif [ "${WORKFLOW_TARGET}" = "PAPER" ] && [ "${history_count}" -eq 4 ] \ @@ -324,6 +331,10 @@ jobs: elif [ "${WORKFLOW_TARGET}" = "PAPER" ] && [ "${history_count}" -eq 0 ] \ && [ "${SOURCE_COMMIT}" = "${approved_http_snapshot_candidate}" ]; then image_mode=snapshot + elif [ "${WORKFLOW_TARGET}" = "PAPER" ] && [ "${history_count}" -eq 0 ] \ + && [ -z "${snapshot_setting}" ] \ + && [ "${SOURCE_COMMIT}" = "${approved_probe_snapshot_candidate}" ]; then + image_mode=probe else echo "Image source is not approved." >&2 exit 1 @@ -334,6 +345,9 @@ jobs: elif [ "${image_mode}" = "snapshot" ]; then archive_ref="${approved_http_snapshot_candidate}" image_tag="${approved_http_snapshot_candidate}" + elif [ "${image_mode}" = "probe" ]; then + archive_ref="${approved_probe_snapshot_candidate}" + image_tag="${approved_probe_snapshot_candidate}" else archive_ref="${approved_candidate}" image_tag="${approved_candidate}" diff --git a/docs/account_snapshot_history.md b/docs/account_snapshot_history.md index a9482db..f096e6b 100644 --- a/docs/account_snapshot_history.md +++ b/docs/account_snapshot_history.md @@ -1,6 +1,6 @@ # 日频账户资产记录 -这是关闭默认的工程记录,不是已启用的资产曲线。每日 `execution-report-heartbeat` 在原有检查之后可以多跑一步;前面的 heartbeat 步骤失败时,这一步按 `success()` 跳过。不新增 Cloud Scheduler,也不改 `/run`、`/probe`、`/dry-run`。 +这是关闭默认的工程记录,不是已启用的资产曲线。每日 `execution-report-heartbeat` 在现有 Google Cloud 认证后单独运行 PAPER 采样步骤,随后继续原 heartbeat 检查。采样失败会保留失败状态但不跳过原检查,最后再使整个 workflow 失败。不新增 Cloud Scheduler,也不改 `/run`、`/probe`、`/dry-run`。 ## 何时会写 @@ -8,14 +8,14 @@ - 当前 matrix 目标的 `label` 是 `PAPER` - GitHub variable `ACCOUNT_HISTORY_RECORDING_ENABLED` 精确等于 `true` -- 同一步提供 `ACCOUNT_HISTORY_SERVICE_URL`、`ACCOUNT_HISTORY_GCS_PREFIX` +- 同一步提供 `ACCOUNT_HISTORY_SERVICE_URL`、`ACCOUNT_HISTORY_GCS_PREFIX` 和非敏感 `ACCOUNT_HISTORY_EXPECTED_SOURCE_BINDING_ID` - 目标 ID 来自 `matrix.target.id`,期望 scope 来自 `matrix.target.label`,项目来自已有 `GCP_PROJECT_ID` -`PAPER` 只表示这份清单配置的范围,不能据此推断券商账户身份。脚本再要求期望 scope 精确为 `PAPER`,服务根必须是没有 userinfo、query、fragment 和额外 path 的 HTTPS `*.run.app`,并且只 GET 该源站的 `/account-snapshot`。GCS 前缀最后一段必须是 `account_snapshots`,不能落在 execution report 路径上。缺任何一项就失败,不补默认值。 +`PAPER` 只表示这份清单配置的范围,不能据此推断券商账户身份。脚本从已校验的 runtime target manifest 读取准确 PAPER service/region,并先读回 `{service}-probe-scheduler` 的完整 job。只有 job 完整资源名、`ENABLED`、`POST {service_url}/probe`、空 body、Scheduler OIDC service account/audience 及零重试配置全部匹配时,才调用一次 Cloud Scheduler `jobs:run`。不直连 internal Cloud Run,也不使用 `/account-snapshot`。Scheduler 请求结果未知时不重触发。缺任何配置就失败,不补默认值。 -可选 QRS 发布仍默认关闭;只有 `ACCOUNT_FACTS_SYNC_ENABLED` 精确为 `true` 时才使用 `ACCOUNT_FACTS_SYNC_URL` 与专用 `ACCOUNT_FACTS_SYNC_TOKEN`。启用时会先校验 URL 和非空、无首尾空白的 token;配置错误在读取快照或创建日对象前失败,不消耗该日的 create-only 名额。URL 必须是 HTTPS 的精确 `/api/account-facts/sync`,不接受 userinfo、端口、query、fragment 或其他路径;POST 不跟随重定向,设置超时并限制响应体大小。token 只作为该 workflow step 的环境变量和 Authorization header 使用,不打印。 +`ACCOUNT_HISTORY_EXPECTED_SOURCE_BINDING_ID` 必须由部署后的可信读回提供,并与 QRS 的预期绑定完全相同;不能从第一个 GCS 对象反推信任。GCS 前缀最后一段必须是 `account_snapshots`,不能落在 execution report 路径上。脚本只查本次触发 UTC 日期与当前 UTC 日期下准确的 target/source 前缀,限单页和 64KiB 对象。列举与固定 generation 读取都带超时、`retry=None`;若分页截断、对象变大或 generation 改变则失败,不把部分结果当完整结果。 -OIDC 使用 heartbeat 里已经配置的 gcloud:`gcloud auth print-identity-token --audiences=<服务根> --quiet`。stdout 只留在内存,不打印,也不放进参数;失败只报短类别。HTTP 有超时、不跟随重定向、不重试。观察时钟在响应收齐之后读取。 +接受 `jobs:run` 后,最多等待 180 秒,每 5 秒只读查询 GCS。候选必须是本次触发之后开始、15 分钟内的完整同源观察;取最新合法对象。若没有可信新鲜对象就失败,不追加 Scheduler 触发。QRS POST 前再次检查观察起始时间仍在 15 分钟窗口内。 ## 保存什么 @@ -23,9 +23,9 @@ OIDC 使用 heartbeat 里已经配置的 gcloud:`gcloud auth print-identity-to 对象只保留分币种 `broker_reported_balances`(`currency`、`net_assets`、`total_cash`)和 `cash`(`currency`、`available_cash`、`frozen_cash`、`settling_cash`)。金额必须是有限十进制字符串;负数保持原值,不改成零,也不把币种加总。不保存持仓、订单、token、secret 或完整响应。 -路径是 `{prefix}/{target_id}/{source_binding_id}/{YYYY-MM-DD}.json`。目标与来源绑定分成不同路径段,不把两段来源拼进同一个对象。`create_text` 只创建:已有对象时结果是 `already_recorded`,不覆盖、不重新拉取。存储结果不明则停止,不写占位点,也不重试。错误输出只有短类别。 +producer 为每次观察写入 `{prefix}/paper/{source_binding_id}/{observation_date}/{HHMMSSffffffZ.json}`,其中日期来自 `observed_started_at`,文件名使用 `observed_finished_at`。consumer 固定读取 listing 返回的 generation、保留原字节,并在 POST 前重新验证合同、路径、来源、UTC日期和时间,不改写时间或拼接来源。 -仅当本次 `create_text` 明确返回新建成功时,脚本才把完全相同的历史 JSON body POST 到 QRS。`already_recorded` 明确跳过发布,不能把这次新读取的内容冒充为已保存对象;`store_unknown` 不 POST。QRS 发布状态与历史记录状态分开输出:发布拒绝或结果未知不会撤销已写入历史;未知 POST 不自动重试。现有流程没有已存对象的补送入口,且接收端默认拒绝超过 15 分钟观察窗口的记录,因此不能靠下一次日常运行可靠补送;配置预检不解决请求结果未知的情况。QRS `ok=true` 且回读的目标、观察日、观察结束时间匹配,只表示接收端确认保存,不证明页面已经展示或数据完成物理账户身份核验。接收端按其可信配置绑定目标与来源,调用方不传账户 key 或身份结论。 +只有从该受限 GCS listing 选出的同一对象才会被 POST 到 QRS,发送 body 是读取到的原始字节,不重新序列化。QRS 发布状态与观察读取状态分开输出;发布拒绝或结果未知不会改变原对象,未知 POST 不自动重试。QRS `ok=true` 且回读的目标、观察日、观察结束时间匹配,只表示接收端确认保存,不证明页面已经展示或数据完成物理账户身份核验。接收端按其可信配置绑定目标与来源,调用方不传账户 key 或身份结论。 这份记录不是 TWR,不是收益率,也不授予 live 权限。 @@ -33,7 +33,7 @@ OIDC 使用 heartbeat 里已经配置的 gcloud:`gcloud auth print-identity-to `sync-cloud-run-env.yml` 的 image-only 入口仍由同仓 main workflow 控制。PAPER 可以额外把已审查候选 `0b939723c1db3ef59175535998b470cbcd4b8824` 做成无流量镜像,并只更新 image、commit、run 标签和四个 history 环境变量。HK/SG 不接受这个候选。准入读取候选的 `uv.lock`、`pyproject.toml` 和 `qsl.toml`,不执行候选脚本。本轮只改源码和离线测试,没有执行云端暂存,也没有切流。 -该固定候选使用自然运行周期中已经读取的余额生成记录,与上文 main 的 HTTP heartbeat 快照入口不同。无流量暂存不会启用 HTTP 快照链,不设置 `ACCOUNT_HISTORY_SERVICE_URL`,也不调用 `/run` 或产生首个样本。上文 HTTP 入口条件不能作为该候选的启用步骤;正式采用及自然周期首样本需要分别核验。 +已审固定候选 `a2921d157efb887e9210fad6734ca040ea6e5293` 在自然 PAPER `/probe` 周期中复用一次原始余额读数生成记录。无流量暂存不会切换流量或产生首个样本。日常 heartbeat 现在经已有内部 probe Scheduler 触发该路径,再读取可信 GCS 对象;源码接线本身不证明已产生对象或网站已显示。 ## 尚未启用 diff --git a/scripts/record_daily_account_snapshot.py b/scripts/record_daily_account_snapshot.py index 4f57e79..e30a840 100644 --- a/scripts/record_daily_account_snapshot.py +++ b/scripts/record_daily_account_snapshot.py @@ -1,9 +1,5 @@ #!/usr/bin/env python3 -"""Record one paper account snapshot under an account_snapshots object prefix. - -The daily heartbeat may call this script. It stays idle unless recording is -explicitly enabled, and it never follows redirects, retries, or overwrites. -""" +"""Trigger one internal PAPER probe and sync its bounded GCS observation.""" from __future__ import annotations @@ -11,6 +7,7 @@ import os import re import sys +import time from collections.abc import Callable, Mapping from dataclasses import dataclass from datetime import datetime, timedelta, timezone @@ -22,25 +19,30 @@ HISTORY_SCHEMA = "longbridge_account_snapshot_history.v1" SNAPSHOT_SCHEMA = "longbridge_account_snapshot.v1" SOURCE_KIND = "deployment_scope_token_version" +ACCOUNT_FACTS_SYNC_PATH = "/api/account-facts/sync" EXPECTED_SCOPE = "PAPER" +SCHEDULER_SERVICE_ACCOUNT = "longbridge-platform-scheduler@longbridgequant.iam.gserviceaccount.com" OBSERVATION_WINDOW = timedelta(minutes=15) -HTTP_TIMEOUT_SECONDS = 20 -MAX_RESPONSE_BYTES = 256 * 1024 -MAX_ACCOUNT_FACTS_SYNC_BODY_BYTES = 64 * 1024 -MAX_ACCOUNT_FACTS_SYNC_RESPONSE_BYTES = 64 * 1024 -ACCOUNT_FACTS_SYNC_PATH = "/api/account-facts/sync" -ACCOUNT_FACTS_SYNC_TIMEOUT_SECONDS = 20 -_RUN_APP_HOST = re.compile(r"^(?:[a-z0-9](?:[a-z0-9-]{0,61}[a-z0-9])?\.)+run\.app$") +WAIT_SECONDS = 180.0 +POLL_SECONDS = 5.0 +HTTP_TIMEOUT_SECONDS = 20.0 +GCS_TIMEOUT_SECONDS = 15.0 +MAX_OBJECTS_PER_DAY = 64 +MAX_OBJECT_BYTES = 64 * 1024 +MAX_QRS_BODY_BYTES = 64 * 1024 +MAX_QRS_RESPONSE_BYTES = 64 * 1024 +_RUN_APP_HOST = re.compile(r"^(?:[a-z0-9](?:[a-z0-9-]{0,61}[a-z0-9])?\.)+[a-z0-9](?:[a-z0-9-]{0,61}[a-z0-9])?\.run\.app$") _HTTPS_HOST = re.compile(r"^(?:[a-z0-9](?:[a-z0-9-]{0,61}[a-z0-9])?\.)+[a-z0-9](?:[a-z0-9-]{0,61}[a-z0-9])?$") _BUCKET = re.compile(r"^[a-z0-9][a-z0-9._-]{1,61}[a-z0-9]$") -_TARGET_ID = re.compile(r"^[a-z0-9](?:[a-z0-9-]{0,62}[a-z0-9])?$") _PROJECT_ID = re.compile(r"^[a-z][a-z0-9-]{4,28}[a-z0-9]$") +_TARGET_ID = re.compile(r"^[a-z0-9](?:[a-z0-9-]{0,62}[a-z0-9])?$") _BINDING_ID = re.compile(r"^[0-9a-f]{64}$") _CURRENCY = re.compile(r"^[A-Z]{3}$") _DECIMAL_TEXT = re.compile(r"^-?(?:0|[1-9]\d*)(?:\.\d+)?$") +_FILENAME = re.compile(r"^\d{12}Z\.json$") _FORBIDDEN_PREFIX_PARTS = ("execution-report", "execution_report", "runtime-report") -_CASH_FIELDS = ("currency", "available_cash", "frozen_cash", "settling_cash") _BALANCE_FIELDS = ("currency", "net_assets", "total_cash") +_CASH_FIELDS = ("currency", "available_cash", "frozen_cash", "settling_cash") @dataclass(frozen=True) @@ -58,239 +60,145 @@ def __init__(self, category: str) -> None: @dataclass(frozen=True) class _Config: - audience: str - snapshot_url: str + project_id: str + region: str + service: str + service_url: str prefix: str + bucket: str + prefix_path: str target_id: str - project_id: str + source_binding_id: str + scheduler_job: str + scheduler_resource: str def record_daily_account_snapshot( env: Mapping[str, str], *, - fetch_id_token: Callable[[str], str], - http_get: Callable[..., Any], open_store: Callable[[str], Any], + session_factory: Callable[[], Any], + http_post: Callable[..., Any], now_reader: Callable[[], datetime], - http_post: Callable[..., Any] | None = None, + monotonic: Callable[[], float] = time.monotonic, + sleep: Callable[[float], None] = time.sleep, ) -> DailyAccountRecordResult: - """Validate one snapshot response and create the first object for that day.""" + """Validate, trigger once, select a fresh same-source object, and POST once.""" if str(env.get("ACCOUNT_HISTORY_RECORDING_ENABLED") or "").strip() != "true": return DailyAccountRecordResult("disabled", publish_status="disabled") try: + config = _config(env) sync_config = _account_facts_sync_config(env) except _Rejected as rejected: - return DailyAccountRecordResult( - "error", rejected.category, "rejected", rejected.category - ) + return DailyAccountRecordResult("error", rejected.category, "rejected", rejected.category) + + session = None try: - config = _config(env) + session = session_factory() + job = _scheduler_get(session, config) + _validate_scheduler_job(job, config) except _Rejected as rejected: + _close(session) return DailyAccountRecordResult("error", rejected.category) - try: - token = fetch_id_token(config.audience) - except Exception: - return DailyAccountRecordResult("error", "token_unavailable") - if not isinstance(token, str) or not token.strip(): - return DailyAccountRecordResult("error", "token_unavailable") - try: - response = http_get( - config.snapshot_url, - headers={"Authorization": f"Bearer {token}", "Accept": "application/json"}, - timeout=HTTP_TIMEOUT_SECONDS, - ) except Exception: - return DailyAccountRecordResult("error", "http_failed") - if getattr(response, "status_code", None) != 200 or getattr(response, "is_redirect", False): - return DailyAccountRecordResult("error", "http_failed") + _close(session) + return DailyAccountRecordResult("error", "scheduler_read_failed") try: - observed_at = now_reader() - body, uri = _history_object(response, config=config, now=observed_at) + started_at = _utc_now(now_reader) + _scheduler_run(session, config) except _Rejected as rejected: return DailyAccountRecordResult("error", rejected.category) - try: - created = open_store(config.project_id).create_text(uri, body, "application/json") except Exception: - return DailyAccountRecordResult( - "error", "store_unknown", "skipped_store_unknown" - ) - if created is True: - publish_status, publish_category = _publish_account_facts( - sync_config, body, http_post=http_post or _http_post - ) - return DailyAccountRecordResult( - "recorded", "", publish_status, publish_category - ) - if created is False: - return DailyAccountRecordResult( - "already_recorded", "", "skipped_already_recorded" - ) - return DailyAccountRecordResult( - "error", "store_unknown", "skipped_store_unknown" - ) - - -def _publish_account_facts( - sync_config: tuple[str, str] | None, - body: str, - *, - http_post: Callable[..., Any], -) -> tuple[str, str]: - if sync_config is None: - return "disabled", "" - url, token = sync_config + # The request may have reached Scheduler. Never trigger it a second time. + return DailyAccountRecordResult("error", "scheduler_run_unknown") + finally: + _close(session) - encoded_body = body.encode("utf-8") - if len(encoded_body) > MAX_ACCOUNT_FACTS_SYNC_BODY_BYTES: - return "rejected", "qrs_payload_too_large" try: - response = http_post( - url, - headers={ - "Authorization": f"Bearer {token}", - "Accept": "application/json", - "Content-Type": "application/json", - }, - data=encoded_body, - timeout=ACCOUNT_FACTS_SYNC_TIMEOUT_SECONDS, - allow_redirects=False, - stream=True, + store = open_store(config.project_id) + client = _storage_client(store) + deadline = monotonic() + WAIT_SECONDS + selected = _wait_for_observation( + client, + config, + started_at=started_at, + deadline=deadline, + now_reader=now_reader, + monotonic=monotonic, + sleep=sleep, ) + if monotonic() >= deadline: + return DailyAccountRecordResult("error", "observation_timeout") + except _Rejected as rejected: + return DailyAccountRecordResult("error", rejected.category) except Exception: - return "unknown", "qrs_request_unknown" - - try: - status_code = getattr(response, "status_code", None) - if not isinstance(status_code, int): - return "unknown", "qrs_response_invalid" - if getattr(response, "is_redirect", False) or 300 <= status_code < 400: - return "rejected", "qrs_redirect_rejected" - if 400 <= status_code < 500: - return "rejected", "qrs_http_rejected" - if status_code >= 500 or not 200 <= status_code < 300: - return "unknown", "qrs_response_unknown" - try: - payload = _bounded_response_json( - response, MAX_ACCOUNT_FACTS_SYNC_RESPONSE_BYTES - ) - except _Rejected as rejected: - return "unknown", rejected.category - if isinstance(payload, dict) and payload.get("ok") is False: - return "rejected", "qrs_application_rejected" - try: - sent = json.loads(body) - except json.JSONDecodeError: - return "unknown", "qrs_request_body_invalid" - if ( - isinstance(payload, dict) - and payload.get("ok") is True - and payload.get("stored") is True - and payload.get("target_id") == sent.get("target_id") - and payload.get("observation_date") == sent.get("observation_date") - and payload.get("observed_finished_at") - == sent.get("observed_finished_at") - ): - return "published", "" - return "unknown", "qrs_response_unconfirmed" - finally: - close = getattr(response, "close", None) - if callable(close): - try: - close() - except Exception: - pass - - -def _account_facts_sync_url(value: str) -> str: - try: - parsed = urlsplit(value.strip()) - port = parsed.port - except ValueError: - raise _Rejected("qrs_config_invalid") from None - host = parsed.hostname or "" - if ( - parsed.scheme != "https" - or parsed.username - or parsed.password - or port is not None - or parsed.query - or parsed.fragment - or parsed.path != ACCOUNT_FACTS_SYNC_PATH - or _HTTPS_HOST.fullmatch(host) is None - ): - raise _Rejected("qrs_config_invalid") - return f"https://{host}{ACCOUNT_FACTS_SYNC_PATH}" - - -def _account_facts_sync_config(env: Mapping[str, str]) -> tuple[str, str] | None: - if str(env.get("ACCOUNT_FACTS_SYNC_ENABLED") or "").strip() != "true": - return None - url = _account_facts_sync_url(str(env.get("ACCOUNT_FACTS_SYNC_URL") or "")) - token = str(env.get("ACCOUNT_FACTS_SYNC_TOKEN") or "") - if not token or token != token.strip(): - raise _Rejected("qrs_config_invalid") - return url, token - + return DailyAccountRecordResult("error", "gcs_read_failed") -def _bounded_response_json(response: Any, max_bytes: int) -> Any: - chunks: list[bytes] = [] - total = 0 - iterator = getattr(response, "iter_content", None) - if callable(iterator): - try: - for chunk in iterator(chunk_size=8192): - if not chunk: - continue - if not isinstance(chunk, bytes): - raise _Rejected("qrs_response_invalid") - total += len(chunk) - if total > max_bytes: - raise _Rejected("qrs_response_too_large") - chunks.append(chunk) - except _Rejected: - raise - except Exception: - raise _Rejected("qrs_response_invalid") from None - content = b"".join(chunks) - else: - content = getattr(response, "content", b"") - if not isinstance(content, bytes) or len(content) > max_bytes: - raise _Rejected("qrs_response_too_large") - try: - return json.loads(content) - except (UnicodeDecodeError, json.JSONDecodeError): - raise _Rejected("qrs_response_invalid") from None + if selected is None: + return DailyAccountRecordResult("error", "observation_timeout") + body, record = selected + if sync_config is None: + return DailyAccountRecordResult("recorded", publish_status="disabled") + if not _is_fresh(record, _utc_now(now_reader)): + return DailyAccountRecordResult("recorded", publish_status="rejected", publish_category="observation_stale") + publish_status, publish_category = _publish_account_facts( + sync_config, body, record, http_post=http_post + ) + return DailyAccountRecordResult("recorded", "", publish_status, publish_category) def _config(env: Mapping[str, str]) -> _Config: - audience, snapshot_url = _service_url(str(env.get("ACCOUNT_HISTORY_SERVICE_URL") or "")) - prefix = _gcs_prefix(str(env.get("ACCOUNT_HISTORY_GCS_PREFIX") or "")) - target_id = str(env.get("ACCOUNT_HISTORY_TARGET_ID") or "").strip() - scope = str(env.get("ACCOUNT_HISTORY_EXPECTED_SCOPE") or "").strip() + prefix, bucket, prefix_path = _gcs_prefix(str(env.get("ACCOUNT_HISTORY_GCS_PREFIX") or "")) project_id = str(env.get("GOOGLE_CLOUD_PROJECT") or "").strip() + target_id = str(env.get("ACCOUNT_HISTORY_TARGET_ID") or "").strip() + expected_scope = str(env.get("ACCOUNT_HISTORY_EXPECTED_SCOPE") or "").strip() + source_binding_id = str(env.get("ACCOUNT_HISTORY_EXPECTED_SOURCE_BINDING_ID") or "").strip() if ( - _TARGET_ID.fullmatch(target_id) is None - or scope != EXPECTED_SCOPE - or _PROJECT_ID.fullmatch(project_id) is None + _PROJECT_ID.fullmatch(project_id) is None + or target_id != "paper" + or _TARGET_ID.fullmatch(target_id) is None + or expected_scope != EXPECTED_SCOPE + or _BINDING_ID.fullmatch(source_binding_id) is None ): raise _Rejected("config_invalid") + try: + from application.runtime_target_manifest import load_runtime_target_manifest + + matches = [target for target in load_runtime_target_manifest().targets if target.id == "paper"] + except Exception: + raise _Rejected("config_invalid") from None + if len(matches) != 1 or matches[0].mode != "paper": + raise _Rejected("config_invalid") + service = matches[0].service + region = matches[0].region + service_url = _service_url( + str(env.get("ACCOUNT_HISTORY_SERVICE_URL") or ""), service=service, region=region + ) + scheduler_job = f"{service}-probe-scheduler" + scheduler_resource = f"projects/{project_id}/locations/{region}/jobs/{scheduler_job}" return _Config( - audience=audience, - snapshot_url=snapshot_url, + project_id=project_id, + region=region, + service=service, + service_url=service_url, prefix=prefix, + bucket=bucket, + prefix_path=prefix_path, target_id=target_id, - project_id=project_id, + source_binding_id=source_binding_id, + scheduler_job=scheduler_job, + scheduler_resource=scheduler_resource, ) -def _service_url(value: str) -> tuple[str, str]: +def _service_url(value: str, *, service: str, region: str) -> str: try: parsed = urlsplit(value.strip()) - host = parsed.hostname or "" port = parsed.port except ValueError: raise _Rejected("config_invalid") from None + host = parsed.hostname or "" if ( parsed.scheme != "https" or parsed.username @@ -302,11 +210,10 @@ def _service_url(value: str) -> tuple[str, str]: or _RUN_APP_HOST.fullmatch(host) is None ): raise _Rejected("config_invalid") - origin = f"https://{host}" - return origin, f"{origin}/account-snapshot" + return f"https://{host}" -def _gcs_prefix(value: str) -> str: +def _gcs_prefix(value: str) -> tuple[str, str, str]: try: parsed = urlsplit(value.strip()) port = parsed.port @@ -324,113 +231,340 @@ def _gcs_prefix(value: str) -> str: or not segments or segments[-1] != "account_snapshots" or any(segment in {".", ".."} for segment in segments) - or any( - part in segment.lower() - for segment in segments - for part in _FORBIDDEN_PREFIX_PARTS - ) + or any(denied in segment.lower() for segment in segments for denied in _FORBIDDEN_PREFIX_PARTS) ): raise _Rejected("config_invalid") - return f"gs://{parsed.netloc}/{'/'.join(segments)}" + path = "/".join(segments) + return f"gs://{parsed.netloc}/{path}", parsed.netloc, path -def _history_object(response: Any, *, config: _Config, now: datetime) -> tuple[str, str]: - content = getattr(response, "content", b"") - if not isinstance(content, (bytes, bytearray)) or len(content) > MAX_RESPONSE_BYTES: - raise _Rejected("response_invalid") +def _authorized_session(): + import google.auth + from google.auth.transport.requests import AuthorizedSession + from requests.adapters import HTTPAdapter + + credentials, _ = google.auth.default(scopes=["https://www.googleapis.com/auth/cloud-platform"]) + session = AuthorizedSession(credentials, max_refresh_attempts=0) + adapter = HTTPAdapter(max_retries=0) + session.mount("https://", adapter) + session.mount("http://", adapter) + return session + + +def _scheduler_get(session: Any, config: _Config) -> Mapping[str, Any]: + response = session.get( + f"https://cloudscheduler.googleapis.com/v1/{config.scheduler_resource}", + timeout=HTTP_TIMEOUT_SECONDS, + allow_redirects=False, + ) try: - payload = json.loads(content) - except (UnicodeDecodeError, json.JSONDecodeError): - raise _Rejected("response_invalid") from None - if not isinstance(payload, dict): - raise _Rejected("response_invalid") - _require_snapshot_contract(payload) - binding_id = _binding_id(payload.get("source_binding")) - started = _aware(payload.get("observed_started_at")) - finished = _aware(payload.get("observed_finished_at")) - if now.tzinfo is None or now.utcoffset() is None: - raise _Rejected("response_invalid") - current = now.astimezone(timezone.utc) - if finished < started or finished > current or started < current - OBSERVATION_WINDOW: - raise _Rejected("response_invalid") - record = { - "schema_version": HISTORY_SCHEMA, - "snapshot_schema_version": SNAPSHOT_SCHEMA, - "account_scope": EXPECTED_SCOPE, - "target_id": config.target_id, - "source_binding": { - "kind": SOURCE_KIND, - "status": "bound", - "id": binding_id, - }, - "observed_started_at": started.isoformat(), - "observed_finished_at": finished.isoformat(), - "snapshot_atomic": False, - "observation_date": started.date().isoformat(), - "broker_reported_balances": _money_rows( - payload.get("broker_reported_balances"), _BALANCE_FIELDS - ), - "cash": _money_rows(payload.get("cash"), _CASH_FIELDS), - } - body = json.dumps(record, ensure_ascii=True, separators=(",", ":"), sort_keys=True) - uri = f"{config.prefix}/{config.target_id}/{binding_id}/{record['observation_date']}.json" - return body, uri + if getattr(response, "status_code", None) != 200 or getattr(response, "is_redirect", False): + raise _Rejected("scheduler_read_failed") + payload = response.json() + if not isinstance(payload, Mapping): + raise _Rejected("scheduler_read_invalid") + return payload + except _Rejected: + raise + except Exception: + raise _Rejected("scheduler_read_invalid") from None + finally: + _close(response) + + +def _validate_scheduler_job(job: Mapping[str, Any], config: _Config) -> None: + target = job.get("httpTarget") + target = target if isinstance(target, Mapping) else {} + auth = target.get("oidcToken") + auth = auth if isinstance(auth, Mapping) else {} + retry = job.get("retryConfig") + retry = retry if isinstance(retry, Mapping) else {} + body = target.get("body") + if isinstance(body, str): + body_empty = not body + elif isinstance(body, (bytes, bytearray)): + body_empty = not body + else: + body_empty = body is None + if ( + job.get("name") != config.scheduler_resource + or job.get("state") != "ENABLED" + or target.get("httpMethod") != "POST" + or target.get("uri") != f"{config.service_url}/probe" + or not body_empty + or auth.get("serviceAccountEmail") != SCHEDULER_SERVICE_ACCOUNT + or auth.get("audience") != config.service_url + or retry.get("retryCount", 0) != 0 + or isinstance(retry.get("retryCount", 0), bool) + or not _zero_duration(retry.get("maxRetryDuration")) + ): + raise _Rejected("scheduler_job_mismatch") + + +def _zero_duration(value: object) -> bool: + if value is None: + return True + if isinstance(value, str): + return value in {"0s", "0.0s"} + if isinstance(value, Mapping): + return value.get("seconds", 0) == 0 and value.get("nanos", 0) == 0 + return False -def _require_snapshot_contract(payload: Mapping[str, Any]) -> None: +def _scheduler_run(session: Any, config: _Config) -> None: + response = session.post( + f"https://cloudscheduler.googleapis.com/v1/{config.scheduler_resource}:run", + timeout=HTTP_TIMEOUT_SECONDS, + allow_redirects=False, + ) + try: + status = getattr(response, "status_code", None) + if getattr(response, "is_redirect", False) or not isinstance(status, int): + raise _Rejected("scheduler_run_rejected") + if 200 <= status < 300: + return + if status < 500: + raise _Rejected("scheduler_run_rejected") + raise _Rejected("scheduler_run_unknown") + finally: + _close(response) + + +def _storage_client(store: Any) -> Any: + client = getattr(store, "client", None) + if client is None: + raise _Rejected("gcs_config_invalid") + return client + + +def _wait_for_observation( + client: Any, + config: _Config, + *, + started_at: datetime, + deadline: float, + now_reader: Callable[[], datetime], + monotonic: Callable[[], float], + sleep: Callable[[float], None], +) -> tuple[bytes, dict[str, Any]] | None: + while True: + if monotonic() >= deadline: + return None + now = _utc_now(now_reader) + dates = {started_at.date(), now.date()} + candidates = _list_candidates( + client, config, dates, deadline=deadline, monotonic=monotonic + ) + if monotonic() >= deadline: + return None + valid: list[tuple[datetime, bytes, dict[str, Any]]] = [] + for candidate in candidates: + try: + payload, raw = _read_candidate( + client, config, candidate, deadline=deadline, monotonic=monotonic + ) + record = _validate_history_object(payload, config, candidate, started_at, now) + except _Rejected as rejected: + if rejected.category == "generation_changed": + raise + continue + if monotonic() >= deadline: + return None + valid.append((datetime.fromisoformat(record["observed_finished_at"]), raw, record)) + if valid: + _finished, raw, record = max(valid, key=lambda entry: entry[0]) + if monotonic() < deadline and _is_fresh(record, _utc_now(now_reader)): + return raw, record + remaining = deadline - monotonic() + if remaining <= 0: + return None + sleep(min(POLL_SECONDS, remaining)) + + +def _list_candidates( + client: Any, + config: _Config, + dates: set[Any], + *, + deadline: float, + monotonic: Callable[[], float], +) -> list[dict[str, Any]]: + candidates: list[dict[str, Any]] = [] + for day in sorted(dates): + remaining = deadline - monotonic() + if remaining <= 0: + raise _Rejected("observation_timeout") + prefix = f"{config.prefix_path}/paper/{config.source_binding_id}/{day.isoformat()}/" + try: + blobs = client.list_blobs( + config.bucket, + prefix=prefix, + max_results=MAX_OBJECTS_PER_DAY + 1, + page_size=MAX_OBJECTS_PER_DAY + 1, + timeout=min(GCS_TIMEOUT_SECONDS, remaining), + retry=None, + fields="items(name,generation,size),nextPageToken", + ) + pages = getattr(blobs, "pages", None) + if pages is None: + page = list(blobs) + has_next_page = False + else: + first_page = next(iter(pages), ()) + page = list(first_page) + has_next_page = bool(getattr(blobs, "next_page_token", None)) + except Exception: + raise _Rejected("gcs_list_failed") from None + if monotonic() >= deadline: + raise _Rejected("observation_timeout") + if has_next_page or len(page) > MAX_OBJECTS_PER_DAY: + raise _Rejected("gcs_listing_truncated") + for blob in page: + name = getattr(blob, "name", None) + generation = getattr(blob, "generation", None) + size = getattr(blob, "size", None) + if ( + not isinstance(name, str) + or not name.startswith(prefix) + or not _FILENAME.fullmatch(name[len(prefix):]) + or isinstance(generation, bool) + or not str(generation or "").isdigit() + or isinstance(size, bool) + or not isinstance(size, int) + or size < 0 + ): + continue + if size > MAX_OBJECT_BYTES: + raise _Rejected("gcs_object_too_large") + candidates.append({"name": name, "generation": int(generation), "size": size}) + if len(candidates) > MAX_OBJECTS_PER_DAY * 2: + raise _Rejected("gcs_listing_truncated") + return candidates + + +def _read_candidate( + client: Any, + config: _Config, + candidate: Mapping[str, Any], + *, + deadline: float, + monotonic: Callable[[], float], +) -> tuple[dict[str, Any], bytes]: + blob = client.bucket(config.bucket).blob(candidate["name"], generation=candidate["generation"]) + remaining = deadline - monotonic() + if remaining <= 0: + raise _Rejected("observation_timeout") + try: + raw = blob.download_as_bytes( + start=0, + end=MAX_OBJECT_BYTES - 1, + if_generation_match=candidate["generation"], + timeout=min(GCS_TIMEOUT_SECONDS, remaining), + retry=None, + ) + except Exception as exc: + text = str(exc).lower() + if "generation" in text or "precondition" in text or "conditionnotmet" in text: + raise _Rejected("generation_changed") from None + raise _Rejected("gcs_read_failed") from None + if monotonic() >= deadline: + raise _Rejected("observation_timeout") + if not isinstance(raw, (bytes, bytearray)) or len(raw) > MAX_OBJECT_BYTES or len(raw) != candidate["size"]: + raise _Rejected("gcs_object_truncated") + try: + payload = json.loads(raw) + except (UnicodeError, json.JSONDecodeError): + raise _Rejected("gcs_object_invalid") from None + if not isinstance(payload, dict): + raise _Rejected("gcs_object_invalid") + return payload, bytes(raw) + + +def _validate_history_object( + payload: Mapping[str, Any], + config: _Config, + candidate: Mapping[str, Any], + triggered_at: datetime, + now: datetime, +) -> dict[str, Any]: + expected_fields = { + "schema_version", "snapshot_schema_version", "account_scope", "target_id", + "source_binding", "observed_started_at", "observed_finished_at", "snapshot_atomic", + "observation_date", "broker_reported_balances", "cash", + } + binding = payload.get("source_binding") if ( - payload.get("schema_version") != SNAPSHOT_SCHEMA - or payload.get("status") != "partial" + set(payload) != expected_fields + or payload.get("schema_version") != HISTORY_SCHEMA + or payload.get("snapshot_schema_version") != SNAPSHOT_SCHEMA or payload.get("account_scope") != EXPECTED_SCOPE - or payload.get("positions_complete") is not True - or payload.get("cash_complete") is not True - or payload.get("no_order") is not True - or payload.get("live_authority_granted") is not False + or payload.get("target_id") != config.target_id or payload.get("snapshot_atomic") is not False + or not isinstance(binding, Mapping) + or binding.get("kind") != SOURCE_KIND + or binding.get("status") != "bound" + or binding.get("id") != config.source_binding_id ): - raise _Rejected("response_invalid") - - -def _binding_id(value: object) -> str: - if not isinstance(value, Mapping): - raise _Rejected("response_invalid") - binding_id = value.get("id") + raise _Rejected("gcs_object_invalid") + started = _aware(payload.get("observed_started_at")) + finished = _aware(payload.get("observed_finished_at")) + expected_name = ( + f"{config.prefix_path}/paper/{config.source_binding_id}/" + f"{started.date().isoformat()}/{_filename_for(finished)}" + ) if ( - value.get("kind") != SOURCE_KIND - or value.get("status") != "bound" - or not isinstance(binding_id, str) - or _BINDING_ID.fullmatch(binding_id) is None + started < triggered_at + or finished < started + or finished > now + or payload.get("observation_date") != started.date().isoformat() + or candidate["name"] != expected_name + or not _is_fresh(payload, now) ): - raise _Rejected("response_invalid") - return binding_id + raise _Rejected("gcs_object_invalid") + balances = _money_rows(payload.get("broker_reported_balances"), _BALANCE_FIELDS) + cash = _money_rows(payload.get("cash"), _CASH_FIELDS) + return dict(payload) + + +def _filename_for(finished: datetime) -> str: + return finished.astimezone(timezone.utc).strftime("%H%M%S%fZ.json") + + +def _is_fresh(record: Mapping[str, Any], now: datetime) -> bool: + try: + started = _aware(record.get("observed_started_at")) + except _Rejected: + return False + current = now.astimezone(timezone.utc) + return started <= current and started >= current - OBSERVATION_WINDOW def _aware(value: object) -> datetime: if not isinstance(value, str): - raise _Rejected("response_invalid") + raise _Rejected("gcs_object_invalid") try: parsed = datetime.fromisoformat(value.replace("Z", "+00:00")) except ValueError: - raise _Rejected("response_invalid") from None + raise _Rejected("gcs_object_invalid") from None if parsed.tzinfo is None or parsed.utcoffset() is None: - raise _Rejected("response_invalid") + raise _Rejected("gcs_object_invalid") return parsed.astimezone(timezone.utc) def _money_rows(value: object, fields: tuple[str, ...]) -> list[dict[str, str]]: if not isinstance(value, list) or not value: - raise _Rejected("response_invalid") + raise _Rejected("gcs_object_invalid") seen: set[str] = set() - rows: list[dict[str, str]] = [] + result: list[dict[str, str]] = [] for item in value: if not isinstance(item, Mapping): - raise _Rejected("response_invalid") + raise _Rejected("gcs_object_invalid") row: dict[str, str] = {} for field in fields: cell = item.get(field) if field == "currency": if not isinstance(cell, str) or _CURRENCY.fullmatch(cell) is None or cell in seen: - raise _Rejected("response_invalid") + raise _Rejected("gcs_object_invalid") seen.add(cell) elif ( isinstance(cell, bool) @@ -438,67 +572,155 @@ def _money_rows(value: object, fields: tuple[str, ...]) -> list[dict[str, str]]: or _DECIMAL_TEXT.fullmatch(cell) is None or not Decimal(cell).is_finite() ): - raise _Rejected("response_invalid") + raise _Rejected("gcs_object_invalid") row[field] = cell - rows.append(row) - return rows + result.append(row) + return result -def _fetch_id_token(audience: str) -> str: - """Read one WIF identity token from gcloud. Stdout stays in memory.""" +def _account_facts_sync_config(env: Mapping[str, str]) -> tuple[str, str] | None: + if str(env.get("ACCOUNT_FACTS_SYNC_ENABLED") or "").strip() != "true": + return None + url = _account_facts_sync_url(str(env.get("ACCOUNT_FACTS_SYNC_URL") or "")) + token = str(env.get("ACCOUNT_FACTS_SYNC_TOKEN") or "") + if not token or token != token.strip(): + raise _Rejected("qrs_config_invalid") + return url, token - import subprocess - command = [ - "gcloud", - "auth", - "print-identity-token", - f"--audiences={audience}", - "--quiet", - ] +def _account_facts_sync_url(value: str) -> str: try: - completed = subprocess.run( - command, - capture_output=True, + parsed = urlsplit(value.strip()) + port = parsed.port + except ValueError: + raise _Rejected("qrs_config_invalid") from None + host = parsed.hostname or "" + if ( + parsed.scheme != "https" + or parsed.username + or parsed.password + or port is not None + or parsed.query + or parsed.fragment + or parsed.path != ACCOUNT_FACTS_SYNC_PATH + or _HTTPS_HOST.fullmatch(host) is None + ): + raise _Rejected("qrs_config_invalid") + return f"https://{host}{ACCOUNT_FACTS_SYNC_PATH}" + + +def _publish_account_facts( + sync_config: tuple[str, str], + body: bytes, + record: Mapping[str, Any], + *, + http_post: Callable[..., Any], +) -> tuple[str, str]: + if len(body) > MAX_QRS_BODY_BYTES: + return "rejected", "qrs_payload_too_large" + url, token = sync_config + try: + response = http_post( + url, + headers={ + "Authorization": f"Bearer {token}", + "Accept": "application/json", + "Content-Type": "application/json", + }, + data=body, timeout=HTTP_TIMEOUT_SECONDS, - check=False, - shell=False, + allow_redirects=False, + stream=True, ) - except (OSError, subprocess.TimeoutExpired): - raise _Rejected("token_unavailable") from None - if completed.returncode != 0: - raise _Rejected("token_unavailable") - token = completed.stdout.decode("utf-8", "replace").strip() - if not token: - raise _Rejected("token_unavailable") - return token + except Exception: + return "unknown", "qrs_request_unknown" + try: + status = getattr(response, "status_code", None) + if not isinstance(status, int): + return "unknown", "qrs_response_invalid" + if getattr(response, "is_redirect", False) or 300 <= status < 400: + return "rejected", "qrs_redirect_rejected" + if 400 <= status < 500: + return "rejected", "qrs_http_rejected" + if status >= 500 or not 200 <= status < 300: + return "unknown", "qrs_response_unknown" + try: + payload = _bounded_response_json(response) + except _Rejected as rejected: + return "unknown", rejected.category + if isinstance(payload, dict) and payload.get("ok") is False: + return "rejected", "qrs_application_rejected" + if ( + isinstance(payload, dict) + and payload.get("ok") is True + and payload.get("stored") is True + and payload.get("target_id") == record.get("target_id") + and payload.get("observation_date") == record.get("observation_date") + and payload.get("observed_finished_at") == record.get("observed_finished_at") + ): + return "published", "" + return "unknown", "qrs_response_unconfirmed" + finally: + _close(response) + + +def _bounded_response_json(response: Any) -> Any: + chunks: list[bytes] = [] + total = 0 + iterator = getattr(response, "iter_content", None) + if callable(iterator): + try: + for chunk in iterator(chunk_size=8192): + if not chunk: + continue + if not isinstance(chunk, bytes): + raise _Rejected("qrs_response_invalid") + total += len(chunk) + if total > MAX_QRS_RESPONSE_BYTES: + raise _Rejected("qrs_response_too_large") + chunks.append(chunk) + except _Rejected: + raise + except Exception: + raise _Rejected("qrs_response_invalid") from None + content = b"".join(chunks) + else: + content = getattr(response, "content", b"") + if not isinstance(content, bytes) or len(content) > MAX_QRS_RESPONSE_BYTES: + raise _Rejected("qrs_response_too_large") + try: + return json.loads(content) + except (UnicodeDecodeError, json.JSONDecodeError): + raise _Rejected("qrs_response_invalid") from None -def _http_get(url: str, *, headers: Mapping[str, str], timeout: float): +def _default_http_post(url: str, **kwargs: Any): import requests + from requests.adapters import HTTPAdapter - return requests.get(url, headers=dict(headers), timeout=timeout, allow_redirects=False) + session = requests.Session() + session.mount("https://", HTTPAdapter(max_retries=0)) + try: + response = session.post(url, **kwargs) + except Exception: + session.close() + raise + return _OwnedResponse(response, session) -def _http_post( - url: str, - *, - headers: Mapping[str, str], - data: bytes, - timeout: float, - allow_redirects: bool, - stream: bool, -): - import requests +class _OwnedResponse: + def __init__(self, response: Any, session: Any) -> None: + self._response = response + self._session = session - return requests.post( - url, - headers=dict(headers), - data=data, - timeout=timeout, - allow_redirects=allow_redirects, - stream=stream, - ) + def __getattr__(self, name: str) -> Any: + return getattr(self._response, name) + + def close(self) -> None: + try: + self._response.close() + finally: + self._session.close() def _open_store(project_id: str): @@ -507,35 +729,43 @@ def _open_store(project_id: str): return get_object_store(project_id=project_id) +def _utc_now(reader: Callable[[], datetime]) -> datetime: + value = reader() + if not isinstance(value, datetime) or value.tzinfo is None or value.utcoffset() is None: + raise _Rejected("clock_invalid") + return value.astimezone(timezone.utc) + + +def _close(value: Any) -> None: + close = getattr(value, "close", None) + if callable(close): + try: + close() + except Exception: + pass + + def main( argv: list[str] | None = None, *, environ: Mapping[str, str] | None = None, - fetch_id_token: Callable[[str], str] | None = None, - http_get: Callable[..., Any] | None = None, - http_post: Callable[..., Any] | None = None, open_store: Callable[[str], Any] | None = None, + session_factory: Callable[[], Any] | None = None, + http_post: Callable[..., Any] | None = None, now_reader: Callable[[], datetime] | None = None, ) -> int: - """Record from the environment. Extra arguments are rejected without I/O.""" - args = sys.argv[1:] if argv is None else argv if args: print("error: config_invalid") return 1 result = record_daily_account_snapshot( os.environ if environ is None else environ, - fetch_id_token=fetch_id_token or _fetch_id_token, - http_get=http_get or _http_get, - http_post=http_post or _http_post, open_store=open_store or _open_store, + session_factory=session_factory or _authorized_session, + http_post=http_post or _default_http_post, now_reader=now_reader or (lambda: datetime.now(timezone.utc)), ) - record_part = ( - f"record=error:{result.category}" - if result.status == "error" - else f"record={result.status}" - ) + record_part = f"record=error:{result.category}" if result.status == "error" else f"record={result.status}" publish_part = f"account_facts_publish={result.publish_status}" print(f"{record_part} {publish_part}") if result.status == "error" or result.publish_status in {"rejected", "unknown"}: diff --git a/scripts/verify_deployed_runtime_target_admission.py b/scripts/verify_deployed_runtime_target_admission.py index 4a15cc4..f5cc7a0 100644 --- a/scripts/verify_deployed_runtime_target_admission.py +++ b/scripts/verify_deployed_runtime_target_admission.py @@ -37,6 +37,8 @@ class AdmissionError(ValueError): APPROVED_PAPER_HISTORY_CANDIDATE = "0b939723c1db3ef59175535998b470cbcd4b8824" # Reviewed PAPER HTTP snapshot image. Distinct from the natural-cycle history archive above. APPROVED_PAPER_HTTP_SNAPSHOT_CANDIDATE = "d8314a61df697cae1dd03a78ddc5c2fc4179ec67" +# Reviewed PAPER internal-probe snapshot producer. It carries no history or snapshot env update. +APPROVED_PAPER_PROBE_SNAPSHOT_CANDIDATE = "a2921d157efb887e9210fad6734ca040ea6e5293" # One already-staged PAPER revision that may hold history env while serving still runs an older image. _PAPER_HTTP_STAGED_SOURCE_REVISION = "longbridge-quant-paper-service-r36423178119" _PAPER_HTTP_STAGED_SOURCE_COMMIT = APPROVED_PAPER_HISTORY_CANDIDATE @@ -413,7 +415,7 @@ def history_update(env: Mapping[str, str], *, workflow_target: str, project_id: if _PROJECT_ID.fullmatch(project_id) is None or _TARGET_ID.fullmatch(values["ACCOUNT_HISTORY_TARGET_ID"]) is None: raise AdmissionError("history settings are invalid") try: - prefix = _gcs_prefix(values["ACCOUNT_HISTORY_GCS_PREFIX"]) + prefix, _, _ = _gcs_prefix(values["ACCOUNT_HISTORY_GCS_PREFIX"]) except _Rejected: raise AdmissionError("history settings are invalid") from None if prefix != values["ACCOUNT_HISTORY_GCS_PREFIX"]: @@ -753,6 +755,13 @@ def _validate_image_only_source( elif image_commit == APPROVED_PAPER_HTTP_SNAPSHOT_CANDIDATE: if str(env.get("WORKFLOW_TARGET") or "") != "PAPER" or history is not None: raise AdmissionError("image source is not approved") + elif image_commit == APPROVED_PAPER_PROBE_SNAPSHOT_CANDIDATE: + if ( + str(env.get("WORKFLOW_TARGET") or "") != "PAPER" + or history is not None + or snapshot_value is not None + ): + raise AdmissionError("image source is not approved") elif _is_exact_main_image(image_commit, env): if history is not None: raise AdmissionError("image source is not approved") diff --git a/tests/test_daily_account_snapshot.py b/tests/test_daily_account_snapshot.py index ed32742..569389f 100644 --- a/tests/test_daily_account_snapshot.py +++ b/tests/test_daily_account_snapshot.py @@ -7,839 +7,529 @@ import pytest - ROOT = Path(__file__).resolve().parents[1] if str(ROOT) not in sys.path: sys.path.insert(0, str(ROOT)) -from scripts.record_daily_account_snapshot import ( # noqa: E402 - record_daily_account_snapshot, - main, -) +from scripts import record_daily_account_snapshot as snapshots # noqa: E402 -NOW = datetime(2026, 9, 28, 12, 0, tzinfo=timezone.utc) -BINDING_A = "a" * 64 -BINDING_B = "b" * 64 -SERVICE = "https://longbridge-quant-paper-service-ab12.asia-east1.run.app" -PREFIX = "gs://acct-history/account_snapshots" -QRS_SYNC_URL = "https://qrs.example.test/api/account-facts/sync" -SECRET_TOKEN = "synthetic-oidc-token" +T0 = datetime(2026, 9, 28, 12, 0, tzinfo=timezone.utc) +BINDING = "a" * 64 +OTHER_BINDING = "b" * 64 +SERVICE_URL = "https://longbridge-quant-paper-service-kcc3gcgmwq-de.a.run.app" +PREFIX = "gs://acct-history/longbridge/account_snapshots" +QRS_URL = "https://qrs.example.test/api/account-facts/sync" QRS_TOKEN = "synthetic-qrs-token" -LEAK = "raw-secret-value" +SECRET = "synthetic-secret-value" def _env(**overrides): - env = { + result = { "ACCOUNT_HISTORY_RECORDING_ENABLED": "true", - "ACCOUNT_HISTORY_SERVICE_URL": SERVICE, + "ACCOUNT_HISTORY_SERVICE_URL": SERVICE_URL, "ACCOUNT_HISTORY_GCS_PREFIX": PREFIX, "ACCOUNT_HISTORY_TARGET_ID": "paper", "ACCOUNT_HISTORY_EXPECTED_SCOPE": "PAPER", + "ACCOUNT_HISTORY_EXPECTED_SOURCE_BINDING_ID": BINDING, "GOOGLE_CLOUD_PROJECT": "longbridgequant", } - env.update(overrides) - return env + result.update(overrides) + return result + + +def _job(**overrides): + job = { + "name": "projects/longbridgequant/locations/asia-east1/jobs/longbridge-quant-paper-service-probe-scheduler", + "state": "ENABLED", + "httpTarget": { + "httpMethod": "POST", + "uri": f"{SERVICE_URL}/probe", + "oidcToken": { + "serviceAccountEmail": "longbridge-platform-scheduler@longbridgequant.iam.gserviceaccount.com", + "audience": SERVICE_URL, + }, + }, + "retryConfig": {"retryCount": 0, "maxRetryDuration": "0s"}, + } + job.update(overrides) + return job -def _payload(**overrides): - started = NOW - timedelta(minutes=2) - finished = NOW - timedelta(minutes=1) - payload = { - "schema_version": "longbridge_account_snapshot.v1", - "status": "partial", +def _history(started=None, finished=None, *, binding=BINDING, balances=None, cash=None): + started = started or T0 + timedelta(seconds=1) + finished = finished or T0 + timedelta(seconds=2) + return { + "schema_version": "longbridge_account_snapshot_history.v1", + "snapshot_schema_version": "longbridge_account_snapshot.v1", "account_scope": "PAPER", - "positions_complete": True, - "cash_complete": True, - "no_order": True, - "live_authority_granted": False, - "snapshot_atomic": False, + "target_id": "paper", "source_binding": { "kind": "deployment_scope_token_version", "status": "bound", - "id": BINDING_A, + "id": binding, }, "observed_started_at": started.isoformat(), "observed_finished_at": finished.isoformat(), - "broker_reported_balances": [ - {"currency": "USD", "net_assets": "10", "total_cash": "-2.5"}, - {"currency": "HKD", "net_assets": "3", "total_cash": "1"}, - ], - "cash": [ - { - "currency": "USD", - "available_cash": "-2.5", - "frozen_cash": "0", - "settling_cash": "0", - }, + "snapshot_atomic": False, + "observation_date": started.date().isoformat(), + "broker_reported_balances": balances + or [{"currency": "USD", "net_assets": "10.25", "total_cash": "3"}], + "cash": cash + or [ { "currency": "HKD", "available_cash": "1", "frozen_cash": "0", "settling_cash": "0", - }, + } ], - "positions": [{"symbol": "SOXL.US", "quantity": "1"}], - "token": LEAK, } - payload.update(overrides) - return payload + + +def _object(payload, *, path_date=None, raw=None, generation=7, size=None, error=None): + started = datetime.fromisoformat(payload["observed_started_at"]) + finished = datetime.fromisoformat(payload["observed_finished_at"]) + day = path_date or started.date().isoformat() + object_name = ( + f"longbridge/account_snapshots/paper/{BINDING}/{day}/" + f"{finished.astimezone(timezone.utc).strftime('%H%M%S%fZ.json')}" + ) + raw_bytes = raw if raw is not None else json.dumps(payload, separators=(",", ":")).encode() + return _Blob(object_name, generation, len(raw_bytes) if size is None else size, raw_bytes, error) + + +class _Blob: + def __init__(self, name, generation, size, raw, error=None): + self.name = name + self.generation = generation + self.size = size + self.raw = raw + self.error = error + self.reads = [] + + def download_as_bytes(self, **kwargs): + self.reads.append(kwargs) + if self.error: + raise self.error + if kwargs.get("if_generation_match") != self.generation: + raise RuntimeError("generation condition mismatch") + return self.raw + + +class _PageIterator: + def __init__(self, objects, page_size, before_page=None): + self.objects = objects + self.page_size = page_size + self.before_page = before_page + self.next_page_token = "next" if len(objects) > page_size else None + + @property + def pages(self): + def iterate(): + if self.before_page: + self.before_page() + yield self.objects[: self.page_size] + + return iterate() + + +class _Storage: + def __init__(self, objects=()): + self.client = self + self.objects = list(objects) + self.list_calls = [] + self.blob_calls = [] + + def list_blobs(self, bucket, **kwargs): + self.list_calls.append((bucket, kwargs)) + prefix = kwargs["prefix"] + return _PageIterator( + [obj for obj in self.objects if obj.name.startswith(prefix)], + kwargs["page_size"], + ) + + def bucket(self, bucket): + assert bucket == "acct-history" + return self + + def blob(self, name, *, generation): + self.blob_calls.append((name, generation)) + return next(obj for obj in self.objects if obj.name == name) class _Response: - def __init__(self, status_code, payload=None, content=None): + def __init__(self, status_code=200, payload=None): self.status_code = status_code self.is_redirect = 300 <= status_code < 400 - if content is None: - content = json.dumps(payload if payload is not None else _payload()).encode() - self.content = content + self._payload = payload + self.closed = False + + def json(self): + return self._payload def iter_content(self, chunk_size=8192): - yield self.content + if self._payload is not None: + yield json.dumps(self._payload).encode() def close(self): - pass + self.closed = True -class _Spies: - def __init__(self, response=None, created=True, explode_store=False, post_response=None): +class _Session: + def __init__(self, job=None, run_response=None, run_error=None): + self.job = job or _job() + self.run_response = run_response or _Response() + self.run_error = run_error self.calls = [] - self.response = response or _Response(200) - self.created = created - self.explode_store = explode_store - self.stored = [] - self.post_response = post_response or _Response( - 200, - { - "ok": True, - "stored": True, - "target_id": "paper", - "observation_date": "2026-09-28", - "observed_finished_at": (NOW - timedelta(minutes=1)).isoformat(), - }, - ) + self.closed = False - def fetch_id_token(self, audience): - self.calls.append(("token", audience)) - return SECRET_TOKEN + def get(self, url, **kwargs): + self.calls.append(("get", url, kwargs)) + return _Response(payload=self.job) - def http_get(self, url, *, headers, timeout): - self.calls.append(("http", url, headers, timeout)) - if isinstance(self.response, Exception): - raise self.response - return self.response + def post(self, url, **kwargs): + self.calls.append(("post", url, kwargs)) + if self.run_error: + raise self.run_error + return self.run_response - def http_post(self, url, *, headers, data, timeout, allow_redirects, stream): - self.calls.append(("post", url, headers, data, timeout, allow_redirects, stream)) - if isinstance(self.post_response, Exception): - raise self.post_response - return self.post_response + def close(self): + self.closed = True - def open_store(self, project_id): - self.calls.append(("store", project_id)) - if self.explode_store: - raise RuntimeError(f"gs://secret {LEAK}") - spies = self - class _Store: - def create_text(self, uri, data, content_type="text/plain"): - spies.calls.append(("create", uri, data, content_type)) - spies.stored.append((uri, data, content_type)) - if isinstance(spies.created, Exception): - raise spies.created - return spies.created +class _Clock: + def __init__(self, values): + self.values = iter(values) + self.last = values[-1] - return _Store() + def __call__(self): + return next(self.values, self.last) -def _record(env=None, spies=None, now=NOW, now_reader=None): - spies = spies or _Spies() - result = record_daily_account_snapshot( - env if env is not None else _env(), - fetch_id_token=spies.fetch_id_token, - http_get=spies.http_get, - http_post=spies.http_post, - open_store=spies.open_store, - now_reader=now_reader or (lambda: now), - ) - return result, spies - - -def test_disabled_switch_does_not_touch_token_http_or_store(): - result, spies = _record( - _env( - ACCOUNT_HISTORY_RECORDING_ENABLED="false", - ACCOUNT_FACTS_SYNC_ENABLED="true", - ACCOUNT_FACTS_SYNC_URL="", - ACCOUNT_FACTS_SYNC_TOKEN="", +class _Spies: + def __init__(self, *, objects=(), job=None, run_error=None, post_error=None, post_response=None): + self.storage = _Storage(objects) + self.session = _Session(job=job, run_error=run_error) + self.open_calls = [] + self.posts = [] + self.post_error = post_error + self.post_response = post_response + + def open_store(self, project): + self.open_calls.append(project) + return self.storage + + def session_factory(self): + return self.session + + def http_post(self, url, **kwargs): + self.posts.append((url, kwargs)) + if self.post_error: + raise self.post_error + if self.post_response is not None: + return self.post_response + sent = json.loads(kwargs["data"]) + return _Response( + payload={ + "ok": True, + "stored": True, + "target_id": sent["target_id"], + "observation_date": sent["observation_date"], + "observed_finished_at": sent["observed_finished_at"], + } ) - ) - - assert result.status == "disabled" - assert spies.calls == [] -@pytest.mark.parametrize( - "bad_env", - [ - {"ACCOUNT_HISTORY_SERVICE_URL": "http://svc.a.run.app"}, - {"ACCOUNT_HISTORY_SERVICE_URL": "https://user:pass@svc.a.run.app"}, - {"ACCOUNT_HISTORY_SERVICE_URL": "https://svc.a.run.app/account-snapshot"}, - {"ACCOUNT_HISTORY_SERVICE_URL": "https://svc.a.run.app?q=1"}, - {"ACCOUNT_HISTORY_SERVICE_URL": "https://svc.a.run.app#frag"}, - {"ACCOUNT_HISTORY_SERVICE_URL": "https://example.com"}, - {"ACCOUNT_HISTORY_SERVICE_URL": "https://svc.a.run.app:99999"}, - {"ACCOUNT_HISTORY_SERVICE_URL": "https://svc.a.run.app:abc"}, - {"ACCOUNT_HISTORY_SERVICE_URL": "https://[svc.a.run.app]"}, - {"ACCOUNT_HISTORY_SERVICE_URL": ""}, - {"ACCOUNT_HISTORY_GCS_PREFIX": "gs://[acct-history]/account_snapshots"}, - {"ACCOUNT_HISTORY_GCS_PREFIX": "gs://acct-history/execution-reports"}, - {"ACCOUNT_HISTORY_GCS_PREFIX": "gs://acct-history/execution_reports/account_snapshots"}, - {"ACCOUNT_HISTORY_GCS_PREFIX": "https://acct-history/account_snapshots"}, - {"ACCOUNT_HISTORY_GCS_PREFIX": ""}, - {"ACCOUNT_HISTORY_TARGET_ID": ""}, - {"ACCOUNT_HISTORY_TARGET_ID": "../paper"}, - {"ACCOUNT_HISTORY_EXPECTED_SCOPE": "HK"}, - {"ACCOUNT_HISTORY_EXPECTED_SCOPE": ""}, - {"GOOGLE_CLOUD_PROJECT": ""}, - ], -) -def test_invalid_config_matrix(bad_env): - result, spies = _record(_env(**bad_env)) +def _record(env=None, spies=None, *, times=None, monotonic=None, sleep=None): + spies = spies or _Spies() + now_reader = _Clock(times or [T0, T0 + timedelta(seconds=5)]) + clock = {"now": 0.0} - assert result.status == "error" - assert result.category == "config_invalid" - assert spies.calls == [] - assert LEAK not in result.category - assert "run.app" not in result.category + def monotonic_fn(): + return clock["now"] if monotonic is None else monotonic() + def sleep_fn(seconds): + if sleep is not None: + sleep(seconds) + else: + clock["now"] += seconds -@pytest.mark.parametrize("status", [302, 503, 500]) -def test_non_200_and_redirect_matrix(status): - spies = _Spies(_Response(status, _payload())) - result, spies = _record(spies=spies) + result = snapshots.record_daily_account_snapshot( + env or _env(), + open_store=spies.open_store, + session_factory=spies.session_factory, + http_post=spies.http_post, + now_reader=now_reader, + monotonic=monotonic_fn, + sleep=sleep_fn, + ) + return result, spies - assert result.status == "error" - assert result.category == "http_failed" - assert [name for name, *_rest in spies.calls] == ["token", "http"] - assert spies.stored == [] +def test_disabled_and_invalid_binding_do_not_make_cloud_calls(): + result, spies = _record(_env(ACCOUNT_HISTORY_RECORDING_ENABLED="false")) + assert result.status == "disabled" + assert spies.session.calls == [] and spies.open_calls == [] -def test_timeout_does_not_write(): - spies = _Spies(TimeoutError("https://svc.a.run.app timed out")) - result, spies = _record(spies=spies) + result, spies = _record(_env(ACCOUNT_HISTORY_SERVICE_URL="http://longbridge-quant-paper-service-kcc3gcgmwq-de.a.run.app")) + assert result.category == "config_invalid" + assert spies.session.calls == [] and spies.open_calls == [] - assert result.status == "error" - assert result.category == "http_failed" - assert spies.stored == [] - assert "run.app" not in result.category + result, spies = _record(_env(ACCOUNT_HISTORY_EXPECTED_SOURCE_BINDING_ID="")) + assert result.category == "config_invalid" + assert spies.session.calls == [] and spies.open_calls == [] @pytest.mark.parametrize( - "payload", + "mutate", [ - _payload(account_scope="SG"), - _payload(status="blocked"), - _payload(source_binding={"kind": "deployment_scope_token_version", "status": "unavailable", "id": None}), - _payload(source_binding={"kind": "deployment_scope_token_version", "status": "bound", "id": "LATEST"}), - _payload(snapshot_atomic=True), - _payload(no_order=False), - _payload(live_authority_granted=True), - _payload(positions_complete=False), - _payload(cash_complete=False), - _payload(observed_finished_at=(NOW - timedelta(minutes=3)).isoformat()), - _payload(observed_started_at=(NOW + timedelta(minutes=1)).isoformat(), observed_finished_at=(NOW + timedelta(minutes=2)).isoformat()), - _payload(observed_started_at=(NOW - timedelta(minutes=16)).isoformat()), - _payload(observed_started_at="2026-09-28T11:58:00"), - _payload(broker_reported_balances=[{"currency": "USD", "net_assets": 10, "total_cash": "1"}]), - _payload(broker_reported_balances=[{"currency": "USD", "net_assets": True, "total_cash": "1"}]), - _payload(cash=[{"currency": "USD", "available_cash": "NaN", "frozen_cash": "0", "settling_cash": "0"}]), - _payload(cash=[]), - _payload( - cash=[ - {"currency": "USD", "available_cash": "1", "frozen_cash": "0", "settling_cash": "0"}, - {"currency": "USD", "available_cash": "2", "frozen_cash": "0", "settling_cash": "0"}, - ] - ), + lambda j: j.update(name="projects/other/locations/asia-east1/jobs/wrong"), + lambda j: j.update(state="PAUSED"), + lambda j: j["httpTarget"].update(httpMethod="GET"), + lambda j: j["httpTarget"].update(uri=f"{SERVICE_URL}/run"), + lambda j: j["httpTarget"].update(body="e30="), + lambda j: j["httpTarget"]["oidcToken"].update(serviceAccountEmail="other@example.com"), + lambda j: j["httpTarget"]["oidcToken"].update(audience="https://other.run.app"), + lambda j: j["retryConfig"].update(retryCount=1), + lambda j: j["retryConfig"].update(maxRetryDuration="1s"), ], ) -def test_invalid_snapshot_matrix(payload): - spies = _Spies(_Response(200, payload)) - result, spies = _record(spies=spies) +def test_scheduler_job_mismatch_fails_before_run_or_gcs(mutate): + job = _job() + mutate(job) + result, spies = _record(spies=_Spies(job=job)) + assert result.category == "scheduler_job_mismatch" + assert [call[0] for call in spies.session.calls] == ["get"] + assert spies.open_calls == [] + + +def test_scheduler_zero_retry_protobuf_defaults_are_accepted(): + job = _job() + job.pop("retryConfig") + result, spies = _record(spies=_Spies(objects=[_object(_history())])) + assert result.status == "recorded" + assert [call[0] for call in spies.session.calls] == ["get", "post"] - assert result.status == "error" - assert result.category == "response_invalid" - assert spies.stored == [] - assert LEAK not in result.category - assert "-2.5" not in result.category +def test_unknown_scheduler_trigger_is_attempted_once_and_never_lists(): + result, spies = _record(spies=_Spies(run_error=TimeoutError(SECRET))) + assert result.category == "scheduler_run_unknown" + assert [call[0] for call in spies.session.calls] == ["get", "post"] + assert spies.open_calls == [] and spies.posts == [] -def test_record_stores_only_whitelisted_multi_currency_fields(): - result, spies = _record() - assert result.status == "recorded" - assert [name for name, *_rest in spies.calls] == ["token", "http", "store", "create"] - audience = spies.calls[0][1] - url = spies.calls[1][1] - headers = spies.calls[1][2] - assert audience == SERVICE - assert url == f"{SERVICE}/account-snapshot" - assert headers["Authorization"] == f"Bearer {SECRET_TOKEN}" - assert spies.calls[1][3] > 0 - uri, data, content_type = spies.stored[0] - assert uri == f"{PREFIX}/paper/{BINDING_A}/2026-09-28.json" - assert content_type == "application/json" - saved = json.loads(data) - assert saved["schema_version"] == "longbridge_account_snapshot_history.v1" - assert saved["snapshot_schema_version"] == "longbridge_account_snapshot.v1" - assert saved["account_scope"] == "PAPER" - assert saved["target_id"] == "paper" - assert saved["source_binding"] == { - "kind": "deployment_scope_token_version", - "status": "bound", - "id": BINDING_A, - } - assert saved["snapshot_atomic"] is False - assert saved["observation_date"] == "2026-09-28" - assert saved["broker_reported_balances"] == [ - {"currency": "USD", "net_assets": "10", "total_cash": "-2.5"}, - {"currency": "HKD", "net_assets": "3", "total_cash": "1"}, - ] - assert saved["cash"][0]["available_cash"] == "-2.5" - assert saved["cash"][1]["currency"] == "HKD" - assert "positions" not in saved - assert "token" not in saved - assert "equity" not in saved - assert LEAK not in data - assert "net_assets_total" not in saved - - -def test_publishes_exact_created_history_body_once_to_qrs(): +def test_record_selects_latest_valid_object_and_posts_exact_original_bytes(): + older = _history(T0 + timedelta(seconds=1), T0 + timedelta(seconds=2)) + latest = _history(T0 + timedelta(seconds=3), T0 + timedelta(seconds=4)) + raw = json.dumps(latest, indent=2).encode() + spies = _Spies(objects=[_object(older), _object(latest, raw=raw)]) result, spies = _record( - env=_env( - ACCOUNT_FACTS_SYNC_ENABLED="true", - ACCOUNT_FACTS_SYNC_URL=QRS_SYNC_URL, - ACCOUNT_FACTS_SYNC_TOKEN=QRS_TOKEN, - ) + _env(ACCOUNT_FACTS_SYNC_ENABLED="true", ACCOUNT_FACTS_SYNC_URL=QRS_URL, ACCOUNT_FACTS_SYNC_TOKEN=QRS_TOKEN), + spies, ) - - assert result.status == "recorded" assert result.publish_status == "published" - assert [name for name, *_ in spies.calls].count("http") == 1 - posts = [call for call in spies.calls if call[0] == "post"] - assert len(posts) == 1 - _, url, headers, body, timeout, allow_redirects, stream = posts[0] - assert url == QRS_SYNC_URL - assert headers == { - "Authorization": f"Bearer {QRS_TOKEN}", - "Accept": "application/json", - "Content-Type": "application/json", - } - assert body.decode("utf-8") == spies.stored[0][1] - assert timeout > 0 - assert allow_redirects is False - assert stream is True - posted_body = body.decode("utf-8") - assert SECRET_TOKEN not in posted_body - assert QRS_TOKEN not in posted_body - assert LEAK not in posted_body - - -def test_qrs_publish_is_off_by_default(): - result, spies = _record() - - assert result.status == "recorded" - assert result.publish_status == "disabled" - assert [name for name, *_ in spies.calls].count("post") == 0 + assert len(spies.posts) == 1 + url, kwargs = spies.posts[0] + assert url == QRS_URL and kwargs["data"] == raw + assert kwargs["allow_redirects"] is False + assert kwargs["headers"]["Authorization"] == f"Bearer {QRS_TOKEN}" + assert all(call[0] == "post" or call[0] == "get" for call in spies.session.calls) + assert all(call[2]["timeout"] <= snapshots.HTTP_TIMEOUT_SECONDS for call in spies.session.calls) + assert all(kwargs["timeout"] <= snapshots.GCS_TIMEOUT_SECONDS for _bucket, kwargs in spies.storage.list_calls) + assert all(blob.reads[0]["retry"] is None for blob in spies.storage.objects if blob.reads) -def test_store_unknown_never_posts_to_qrs(): - spies = _Spies(created=None) +@pytest.mark.parametrize( + ("status", "receipt", "expected"), + [ + (200, {"target_id": "other", "observation_date": "2026-09-28", "observed_finished_at": (T0 + timedelta(seconds=2)).isoformat()}, "unknown"), + (200, {"target_id": "paper", "observation_date": "2026-09-27", "observed_finished_at": (T0 + timedelta(seconds=2)).isoformat()}, "unknown"), + (200, {"target_id": "paper", "observation_date": "2026-09-28", "observed_finished_at": (T0 + timedelta(seconds=3)).isoformat()}, "unknown"), + (302, {"target_id": "paper", "observation_date": "2026-09-28", "observed_finished_at": (T0 + timedelta(seconds=2)).isoformat()}, "rejected"), + (401, {"target_id": "paper", "observation_date": "2026-09-28", "observed_finished_at": (T0 + timedelta(seconds=2)).isoformat()}, "rejected"), + (500, {"target_id": "paper", "observation_date": "2026-09-28", "observed_finished_at": (T0 + timedelta(seconds=2)).isoformat()}, "unknown"), + ], +) +def test_qrs_receipt_identity_and_failure_matrix_is_single_post(status, receipt, expected): + response = _Response(status, {"ok": True, "stored": True, **receipt}) + spies = _Spies(objects=[_object(_history())], post_response=response) result, spies = _record( - env=_env( - ACCOUNT_FACTS_SYNC_ENABLED="true", - ACCOUNT_FACTS_SYNC_URL=QRS_SYNC_URL, - ACCOUNT_FACTS_SYNC_TOKEN=QRS_TOKEN, - ), - spies=spies, + _env(ACCOUNT_FACTS_SYNC_ENABLED="true", ACCOUNT_FACTS_SYNC_URL=QRS_URL, ACCOUNT_FACTS_SYNC_TOKEN=QRS_TOKEN), + spies, ) - - assert result.category == "store_unknown" - assert [name for name, *_ in spies.calls].count("post") == 0 + assert result.status == "recorded" + assert result.publish_status == expected + assert len(spies.posts) == 1 -def test_already_recorded_does_not_publish_newly_fetched_body(): - spies = _Spies(created=False) +def test_invalid_qrs_config_is_rejected_before_scheduler_or_gcs(): result, spies = _record( - env=_env( - ACCOUNT_FACTS_SYNC_ENABLED="true", - ACCOUNT_FACTS_SYNC_URL=QRS_SYNC_URL, - ACCOUNT_FACTS_SYNC_TOKEN=QRS_TOKEN, - ), - spies=spies, + _env(ACCOUNT_FACTS_SYNC_ENABLED="true", ACCOUNT_FACTS_SYNC_URL="https://qrs.example.test/other", ACCOUNT_FACTS_SYNC_TOKEN=QRS_TOKEN) ) + assert result.category == "qrs_config_invalid" + assert spies.session.calls == [] and spies.open_calls == [] - assert result.status == "already_recorded" - assert result.publish_status == "skipped_already_recorded" - assert [name for name, *_ in spies.calls].count("post") == 0 - - -def test_qrs_redirect_is_not_followed_or_treated_as_published(): - spies = _Spies(post_response=_Response(302, {"Location": "https://other.example.test"})) - result, spies = _record( - env=_env( - ACCOUNT_FACTS_SYNC_ENABLED="true", - ACCOUNT_FACTS_SYNC_URL=QRS_SYNC_URL, - ACCOUNT_FACTS_SYNC_TOKEN=QRS_TOKEN, - ), - spies=spies, - ) +def test_listing_uses_expected_source_and_both_utc_days_at_midnight(): + t0 = datetime(2026, 9, 28, 23, 59, 59, tzinfo=timezone.utc) + started = datetime(2026, 9, 29, 0, 0, 0, tzinfo=timezone.utc) + finished = started + timedelta(seconds=1) + obj = _object(_history(started, finished)) + result, spies = _record(spies=_Spies(objects=[obj]), times=[t0, finished + timedelta(seconds=3)]) + prefixes = {kwargs["prefix"] for _bucket, kwargs in spies.storage.list_calls} assert result.status == "recorded" - assert result.publish_status == "rejected" - assert [name for name, *_ in spies.calls].count("post") == 1 - assert spies.calls[-1][5] is False - - -def test_qrs_timeout_is_unknown_and_never_retried(): - spies = _Spies(post_response=TimeoutError(f"{QRS_TOKEN} {LEAK}")) + assert any(prefix.endswith("/2026-09-28/") for prefix in prefixes) + assert any(prefix.endswith("/2026-09-29/") for prefix in prefixes) + + +def test_wrong_source_and_observation_started_before_trigger_are_ignored(monkeypatch): + monkeypatch.setattr(snapshots, "WAIT_SECONDS", 1) + wrong_source = _history(binding=OTHER_BINDING) + old = _history(T0 - timedelta(seconds=2), T0 - timedelta(seconds=1)) + result, spies = _record(spies=_Spies(objects=[_object(wrong_source), _object(old)])) + assert result.category == "observation_timeout" + assert spies.posts == [] + + t0 = datetime(2026, 9, 28, 23, 59, 59, tzinfo=timezone.utc) + started = t0 + timedelta(microseconds=100) + finished = t0 + timedelta(microseconds=200) + wrong_path = _object(_history(started, finished), path_date="2026-09-29") result, spies = _record( - env=_env( - ACCOUNT_FACTS_SYNC_ENABLED="true", - ACCOUNT_FACTS_SYNC_URL=QRS_SYNC_URL, - ACCOUNT_FACTS_SYNC_TOKEN=QRS_TOKEN, - ), - spies=spies, + spies=_Spies(objects=[wrong_path]), + times=[t0, datetime(2026, 9, 29, 0, 0, 5, tzinfo=timezone.utc)], ) + assert result.category == "observation_timeout" + assert spies.posts == [] - assert result.status == "recorded" - assert result.publish_status == "unknown" - assert [name for name, *_ in spies.calls].count("post") == 1 - assert QRS_TOKEN not in result.publish_category - assert LEAK not in result.publish_category +def test_started_time_controls_fifteen_minute_freshness(): + now = T0 + timedelta(minutes=20) + record = _history(T0, now) + assert not snapshots._is_fresh(record, now) -def test_qrs_error_response_is_sanitized_and_record_stays_recorded(): - spies = _Spies(post_response=_Response(401, {"error": f"bad token {QRS_TOKEN} {LEAK}"})) - result, spies = _record( - env=_env( - ACCOUNT_FACTS_SYNC_ENABLED="true", - ACCOUNT_FACTS_SYNC_URL=QRS_SYNC_URL, - ACCOUNT_FACTS_SYNC_TOKEN=QRS_TOKEN, - ), - spies=spies, - ) - assert result.status == "recorded" - assert result.publish_status == "rejected" - assert result.publish_category == "qrs_http_rejected" - assert QRS_TOKEN not in result.publish_category - assert LEAK not in result.publish_category +def test_fake_storage_read_crossing_deadline_cannot_be_accepted(monkeypatch): + obj = _object(_history()) + spies = _Spies(objects=[obj]) + ticks = {"now": 0.0} + def mono(): + return ticks["now"] -def test_qrs_oversized_response_is_unknown_without_retry(): - spies = _Spies(post_response=_Response(200, content=b"x" * (64 * 1024 + 1))) - result, spies = _record( - env=_env( - ACCOUNT_FACTS_SYNC_ENABLED="true", - ACCOUNT_FACTS_SYNC_URL=QRS_SYNC_URL, - ACCOUNT_FACTS_SYNC_TOKEN=QRS_TOKEN, - ), - spies=spies, - ) + original = obj.download_as_bytes - assert result.status == "recorded" - assert result.publish_status == "unknown" - assert result.publish_category == "qrs_response_too_large" - assert [name for name, *_ in spies.calls].count("post") == 1 + def slow_read(**kwargs): + ticks["now"] += snapshots.WAIT_SECONDS + 1 + return original(**kwargs) + obj.download_as_bytes = slow_read + result, spies = _record(spies=spies, monotonic=mono) + assert result.category == "observation_timeout" + assert spies.posts == [] -def test_qrs_oversized_payload_is_not_sent(): - large_payload = _payload( - cash=[ - { - "currency": "USD", - "available_cash": "1" * (64 * 1024), - "frozen_cash": "0", - "settling_cash": "0", - } - ] - ) - spies = _Spies(response=_Response(200, large_payload)) - result, spies = _record( - env=_env( - ACCOUNT_FACTS_SYNC_ENABLED="true", - ACCOUNT_FACTS_SYNC_URL=QRS_SYNC_URL, - ACCOUNT_FACTS_SYNC_TOKEN=QRS_TOKEN, - ), - spies=spies, - ) +def test_scheduler_and_gcs_calls_are_bounded_and_generation_pinned(): + obj = _object(_history()) + result, spies = _record(spies=_Spies(objects=[obj])) assert result.status == "recorded" - assert result.publish_status == "rejected" - assert result.publish_category == "qrs_payload_too_large" - assert [name for name, *_ in spies.calls].count("post") == 0 - - -def test_cli_reports_qrs_failure_without_erasing_record_success(capsys): - spies = _Spies(post_response=TimeoutError(f"{QRS_TOKEN} {LEAK}")) - code = main( - [], - environ=_env( - ACCOUNT_FACTS_SYNC_ENABLED="true", - ACCOUNT_FACTS_SYNC_URL=QRS_SYNC_URL, - ACCOUNT_FACTS_SYNC_TOKEN=QRS_TOKEN, - ), - fetch_id_token=spies.fetch_id_token, - http_get=spies.http_get, - http_post=spies.http_post, - open_store=spies.open_store, - now_reader=lambda: NOW, - ) + assert spies.session.calls[0][2]["timeout"] <= snapshots.HTTP_TIMEOUT_SECONDS + assert spies.session.calls[1][2]["timeout"] <= snapshots.HTTP_TIMEOUT_SECONDS + assert spies.storage.list_calls[0][1]["retry"] is None + assert spies.storage.list_calls[0][1]["max_results"] == snapshots.MAX_OBJECTS_PER_DAY + 1 + read = obj.reads[0] + assert read["if_generation_match"] == obj.generation + assert read["retry"] is None and read["end"] == snapshots.MAX_OBJECT_BYTES - 1 + + +def test_generation_change_is_fail_closed(): + obj = _object(_history(), error=RuntimeError("conditionNotMet generation changed")) + result, spies = _record(spies=_Spies(objects=[obj])) + assert result.category == "generation_changed" + assert spies.posts == [] + + +def test_oversized_or_truncated_listing_never_posts(): + oversized = _object(_history(), size=snapshots.MAX_OBJECT_BYTES + 1) + result, spies = _record(spies=_Spies(objects=[oversized])) + assert result.category == "gcs_object_too_large" + assert spies.posts == [] + + many = [] + for index in range(snapshots.MAX_OBJECTS_PER_DAY + 1): + finished = T0 + timedelta(seconds=2, microseconds=index) + many.append(_object(_history(T0 + timedelta(seconds=1), finished))) + result, spies = _record(spies=_Spies(objects=many)) + assert result.category == "gcs_listing_truncated" + assert spies.posts == [] + + many = [] + for index in range(snapshots.MAX_OBJECTS_PER_DAY + 2): + finished = T0 + timedelta(seconds=2, microseconds=index) + many.append(_object(_history(T0 + timedelta(seconds=1), finished))) + result, spies = _record(spies=_Spies(objects=many)) + assert result.category == "gcs_listing_truncated" + assert spies.posts == [] + + +def test_slow_gcs_page_crossing_deadline_is_not_consumed_or_posted(): + obj = _object(_history()) + storage = _Storage([obj]) + spies = _Spies() + spies.storage = storage + ticks = {"now": 0.0} + original_list = storage.list_blobs - output = capsys.readouterr().out - assert code == 1 - assert output.strip() == "record=recorded account_facts_publish=unknown" - assert QRS_TOKEN not in output - assert LEAK not in output - - -def test_qrs_sync_url_must_be_exact_https_endpoint(): - for bad_url in ( - "http://qrs.example.test/api/account-facts/sync", - "https://qrs.example.test/api/account-facts/sync/", - "https://qrs.example.test/api/account-facts/sync?next=https://evil.test", - "https://user:pass@qrs.example.test/api/account-facts/sync", - "https://qrs.example.test:443/api/account-facts/sync", - "https://qrs.example.test/other", - ): - spies = _Spies() - result, spies = _record( - env=_env( - ACCOUNT_FACTS_SYNC_ENABLED="true", - ACCOUNT_FACTS_SYNC_URL=bad_url, - ACCOUNT_FACTS_SYNC_TOKEN=QRS_TOKEN, - ), - spies=spies, - ) - assert result.status == "error" - assert result.category == "qrs_config_invalid" - assert result.publish_status == "rejected" - assert spies.calls == [] - assert spies.stored == [] + def slow_list(bucket, **kwargs): + result = original_list(bucket, **kwargs) + result.before_page = lambda: ticks.update(now=snapshots.WAIT_SECONDS + 1) + return result + storage.list_blobs = slow_list + result, spies = _record(spies=spies, monotonic=lambda: ticks["now"]) + assert result.category == "observation_timeout" + assert obj.reads == [] and spies.posts == [] -@pytest.mark.parametrize( - "bad_config", - [ - {"ACCOUNT_FACTS_SYNC_URL": ""}, - {"ACCOUNT_FACTS_SYNC_URL": "https://qrs.example.test/other"}, - {"ACCOUNT_FACTS_SYNC_TOKEN": ""}, - {"ACCOUNT_FACTS_SYNC_TOKEN": " synthetic-qrs-token"}, - ], -) -def test_enabled_qrs_invalid_config_is_rejected_before_snapshot_io(bad_config): - env = _env( - ACCOUNT_FACTS_SYNC_ENABLED="true", - ACCOUNT_FACTS_SYNC_URL=QRS_SYNC_URL, - ACCOUNT_FACTS_SYNC_TOKEN=QRS_TOKEN, - ) - env.update(bad_config) - result, spies = _record(env=env) - assert result.status == "error" - assert result.category == "qrs_config_invalid" - assert result.publish_status == "rejected" - assert result.publish_category == "qrs_config_invalid" - assert spies.calls == [] - assert spies.stored == [] - - -def test_new_source_or_utc_day_uses_a_new_object_segment(): - late = datetime(2026, 9, 28, 0, 5, tzinfo=timezone.utc) - previous_start = datetime(2026, 9, 27, 23, 55, tzinfo=timezone.utc) - spies = _Spies( - _Response( - 200, - _payload( - source_binding={ - "kind": "deployment_scope_token_version", - "status": "bound", - "id": BINDING_B, - "token": LEAK, - }, - observed_started_at=previous_start.isoformat(), - observed_finished_at=late.isoformat(), - ), - ) +def test_cash_currency_set_may_differ_from_balance_currency_set(): + payload = _history( + balances=[{"currency": "USD", "net_assets": "10", "total_cash": "3"}], + cash=[{"currency": "HKD", "available_cash": "1", "frozen_cash": "0", "settling_cash": "0"}], ) - result, spies = _record(spies=spies, now=late) - + result, spies = _record(spies=_Spies(objects=[_object(payload)])) assert result.status == "recorded" - assert spies.stored[0][0] == f"{PREFIX}/paper/{BINDING_B}/2026-09-27.json" - assert LEAK not in spies.stored[0][1] - + assert spies.posts == [] -def test_existing_object_is_already_recorded_without_another_fetch(): - spies = _Spies(created=False) - result, spies = _record(spies=spies) - assert result.status == "already_recorded" - assert [name for name, *_rest in spies.calls] == ["token", "http", "store", "create"] - - -def test_unknown_store_result_is_not_retried(): - spies = _Spies(created=None) - result, spies = _record(spies=spies) - - assert result.status == "error" - assert result.category == "store_unknown" - assert [name for name, *_rest in spies.calls].count("create") == 1 - assert [name for name, *_rest in spies.calls].count("http") == 1 - assert LEAK not in result.category - - spies = _Spies(explode_store=True) - result, spies = _record(spies=spies) - assert result.status == "error" - assert result.category == "store_unknown" - assert "secret" not in result.category - assert [name for name, *_rest in spies.calls].count("http") == 1 - - -def test_cli_disabled_and_success_use_injected_dependencies(capsys): - disabled = main( - [], - environ=_env(ACCOUNT_HISTORY_RECORDING_ENABLED=""), - fetch_id_token=lambda *_args, **_kwargs: pytest.fail("token"), - http_get=lambda *_args, **_kwargs: pytest.fail("http"), - open_store=lambda *_args, **_kwargs: pytest.fail("store"), - now_reader=lambda: NOW, - ) - assert disabled == 0 - assert capsys.readouterr().out.strip() == "record=disabled account_facts_publish=disabled" - - spies = _Spies() - code = main( +def test_qrs_unknown_is_reported_once_and_secrets_are_redacted(capsys): + spies = _Spies(objects=[_object(_history())], post_error=TimeoutError(f"{SECRET} {QRS_TOKEN}")) + code = snapshots.main( [], - environ=_env(), - fetch_id_token=spies.fetch_id_token, - http_get=spies.http_get, + environ=_env(ACCOUNT_FACTS_SYNC_ENABLED="true", ACCOUNT_FACTS_SYNC_URL=QRS_URL, ACCOUNT_FACTS_SYNC_TOKEN=QRS_TOKEN), open_store=spies.open_store, - now_reader=lambda: NOW, - ) - captured = capsys.readouterr() - assert code == 0 - assert captured.out.strip() == "record=recorded account_facts_publish=disabled" - assert SECRET_TOKEN not in captured.out - - -def test_cli_error_is_a_short_category(capsys): - code = main( - [], - environ=_env(ACCOUNT_HISTORY_SERVICE_URL="https://example.com"), - fetch_id_token=lambda *_args, **_kwargs: pytest.fail("token"), - http_get=lambda *_args, **_kwargs: pytest.fail("http"), - open_store=lambda *_args, **_kwargs: pytest.fail("store"), - now_reader=lambda: NOW, + session_factory=spies.session_factory, + http_post=spies.http_post, + now_reader=_Clock([T0, T0 + timedelta(seconds=5)]), ) + output = capsys.readouterr().out assert code == 1 - assert capsys.readouterr().out.strip() == "record=error:config_invalid account_facts_publish=disabled" - - -def test_gcloud_token_wrapper_success_failure_and_timeout(monkeypatch, capsys): - import subprocess - - from scripts.record_daily_account_snapshot import _Rejected, _fetch_id_token - - calls = [] - - def succeed(args, **kwargs): - calls.append((args, kwargs)) - return subprocess.CompletedProcess( - args, - 0, - stdout=b"synthetic-oidc-token\n", - stderr=b"debug " + LEAK.encode(), - ) - - monkeypatch.setattr(subprocess, "run", succeed) - assert _fetch_id_token(SERVICE) == "synthetic-oidc-token" - args, kwargs = calls[0] - assert args == [ - "gcloud", - "auth", - "print-identity-token", - f"--audiences={SERVICE}", - "--quiet", - ] - assert kwargs["capture_output"] is True - assert kwargs["shell"] is False - assert kwargs["timeout"] > 0 - assert "synthetic-oidc-token" not in " ".join(args) - assert capsys.readouterr().out == "" - - def fail(args, **kwargs): - calls.append((args, kwargs)) - return subprocess.CompletedProcess(args, 1, stdout=b"", stderr=LEAK.encode()) - - monkeypatch.setattr(subprocess, "run", fail) - with pytest.raises(_Rejected) as raised: - _fetch_id_token(SERVICE) - assert raised.value.category == "token_unavailable" - assert LEAK not in str(raised.value) - assert len(calls) == 2 - - def time_out(args, **kwargs): - calls.append((args, kwargs)) - raise subprocess.TimeoutExpired( - args, - kwargs["timeout"], - output=b"synthetic-oidc-token", - stderr=LEAK.encode(), - ) - - monkeypatch.setattr(subprocess, "run", time_out) - with pytest.raises(_Rejected) as raised: - _fetch_id_token(SERVICE) - assert raised.value.category == "token_unavailable" - assert "synthetic-oidc-token" not in str(raised.value) - assert LEAK not in str(raised.value) - assert len(calls) == 3 - assert capsys.readouterr().out == "" - - -def test_cli_uses_the_gcloud_wrapper_and_keeps_stderr_private(monkeypatch, capsys): - import subprocess - - def succeed(args, **kwargs): - return subprocess.CompletedProcess( - args, - 0, - stdout=b"synthetic-oidc-token\n", - stderr=LEAK.encode(), - ) - - monkeypatch.setattr(subprocess, "run", succeed) - spies = _Spies() - code = main( - [], - environ=_env(), - http_get=spies.http_get, - open_store=spies.open_store, - now_reader=lambda: NOW, - ) - captured = capsys.readouterr() - assert code == 0 - assert captured.out.strip() == "record=recorded account_facts_publish=disabled" - assert LEAK not in captured.out - assert LEAK not in captured.err - assert spies.calls[0][2]["Authorization"] == "Bearer synthetic-oidc-token" - - -def test_observation_finished_during_the_request_is_recorded(): - request_started = NOW - clock = {"now": request_started} - reads = [] - - def fetch_id_token(_audience): - clock["now"] += timedelta(seconds=4) - return SECRET_TOKEN - - def http_get(_url, *, headers, timeout): - assert headers["Authorization"] == f"Bearer {SECRET_TOKEN}" - assert timeout > 0 - finish = clock["now"] + timedelta(seconds=3) - start = finish - timedelta(seconds=30) - clock["finish"] = finish - clock["now"] = finish + timedelta(milliseconds=200) - return _Response( - 200, - _payload( - observed_started_at=start.isoformat(), - observed_finished_at=finish.isoformat(), - ), - ) - - def now_reader(): - reads.append(clock["now"]) - return clock["now"] - - spies = _Spies() - result = record_daily_account_snapshot( - _env(), - fetch_id_token=fetch_id_token, - http_get=http_get, - open_store=spies.open_store, - now_reader=now_reader, - ) - - assert result.status == "recorded" - assert clock["finish"] > request_started - assert reads[-1] > clock["finish"] - assert spies.stored[0][0].endswith("/2026-09-28.json") - - -def test_finish_after_the_received_clock_is_still_rejected(): - def http_get(_url, *, headers, timeout): - return _Response( - 200, - _payload( - observed_started_at=(NOW - timedelta(minutes=1)).isoformat(), - observed_finished_at=(NOW + timedelta(minutes=1)).isoformat(), - ), - ) + assert output.strip() == "record=recorded account_facts_publish=unknown" + assert SECRET not in output and QRS_TOKEN not in output + assert len(spies.posts) == 1 - spies = _Spies() - result = record_daily_account_snapshot( - _env(), - fetch_id_token=spies.fetch_id_token, - http_get=http_get, - open_store=spies.open_store, - now_reader=lambda: NOW, - ) - assert result.status == "error" - assert result.category == "response_invalid" - assert spies.stored == [] +def test_authorized_session_disables_transport_retries(monkeypatch): + import google.auth + from google.auth.transport.requests import AuthorizedSession + class Credentials: + pass -def test_malformed_url_cli_stays_a_short_category(capsys): - cases = ( - {"ACCOUNT_HISTORY_SERVICE_URL": "https://svc.a.run.app:99999"}, - {"ACCOUNT_HISTORY_SERVICE_URL": "https://[svc.a.run.app]"}, - {"ACCOUNT_HISTORY_GCS_PREFIX": "gs://[acct-history]/account_snapshots"}, - ) - for bad in cases: - code = main( - [], - environ=_env(**bad), - fetch_id_token=lambda *_args, **_kwargs: pytest.fail("token"), - http_get=lambda *_args, **_kwargs: pytest.fail("http"), - open_store=lambda *_args, **_kwargs: pytest.fail("store"), - now_reader=lambda: NOW, - ) - captured = capsys.readouterr() - assert code == 1 - assert captured.out.strip() == "record=error:config_invalid account_facts_publish=disabled" - assert "Traceback" not in captured.err - assert "Port out of range" not in captured.err - assert "IPv6" not in captured.err - - -def test_heartbeat_workflow_adds_an_optional_paper_step_without_a_new_scheduler(): - workflow = (ROOT / ".github/workflows/execution-report-heartbeat.yml").read_text(encoding="utf-8") - heartbeat = workflow.index("Check recent execution report") - record = workflow.index("Record daily paper account snapshot") - assert heartbeat < record - assert 'cron: "20 22 * * *"' in workflow - assert workflow.count("schedule:") == 1 - assert "gcloud scheduler" not in workflow - step = workflow[record:] - assert "success()" in step - assert "matrix.target.label == 'PAPER'" in step - assert "vars.ACCOUNT_HISTORY_RECORDING_ENABLED == 'true'" in step - assert "ACCOUNT_HISTORY_TARGET_ID: ${{ matrix.target.id }}" in step - assert "ACCOUNT_HISTORY_EXPECTED_SCOPE: ${{ matrix.target.label }}" in step - assert "GOOGLE_CLOUD_PROJECT: ${{ env.GCP_PROJECT_ID }}" in step - assert "uv run --no-sync python scripts/record_daily_account_snapshot.py" in step - assert "ACCOUNT_HISTORY_RECORDING_ENABLED: true" not in workflow + monkeypatch.setattr(google.auth, "default", lambda **_kwargs: (Credentials(), "project")) + session = snapshots._authorized_session() + try: + assert isinstance(session, AuthorizedSession) + assert session._max_refresh_attempts == 0 + assert session.adapters["https://"].max_retries.total == 0 + finally: + session.close() diff --git a/tests/test_deployed_runtime_target_admission.py b/tests/test_deployed_runtime_target_admission.py index 1f06a43..fe9f904 100644 --- a/tests/test_deployed_runtime_target_admission.py +++ b/tests/test_deployed_runtime_target_admission.py @@ -83,6 +83,7 @@ def test_verify_service_rejects_target_drift(target, profile, message): } MAIN_SHA = "a" * 40 HTTP_SNAPSHOT_CANDIDATE = admission.APPROVED_PAPER_HTTP_SNAPSHOT_CANDIDATE +PROBE_SNAPSHOT_CANDIDATE = admission.APPROVED_PAPER_PROBE_SNAPSHOT_CANDIDATE STAGED_SOURCE_REVISION = "longbridge-quant-paper-service-r36423178119" STAGED_SOURCE_IMAGE = ( "asia-east1-docker.pkg.dev/synthetic-project/images/longbridgeplatform/paper-service" @@ -752,3 +753,84 @@ def run(command): entry["value"] = "gs://paper-bucket/tampered" if key.endswith("GCS_PREFIX") else "false" with pytest.raises(admission.AdmissionError): confirm(changed) + + +def test_probe_snapshot_candidate_uses_strict_serving_configuration(): + service = _service() + calls = [] + + def run(command): + calls.append(list(command)) + if command[:3] == ["gcloud", "run", "revisions"]: + return json.dumps(_revision("serving-rev", SERVING, "serving-image")) + sha, name = command[2].split(":", 1) + assert sha in {PROBE_SNAPSHOT_CANDIDATE, SERVING} + return _declaration(name, UES) + + plan = admission.prepare_image_only_staging( + service="paper-service", project="synthetic-project", region="synthetic-region", + service_json=service, + env={**_main_env(), "SOURCE_COMMIT": PROBE_SNAPSHOT_CANDIDATE}, + image_commit=PROBE_SNAPSHOT_CANDIDATE, + run=run, + ) + assert plan["image_commit"] == PROBE_SNAPSHOT_CANDIDATE + assert plan["history_arg"] == plan["snapshot_arg"] == "" + assert plan["history_values"] == plan["retained_history_values"] == {} + assert plan["template_digest"] == admission._config_digest(admission._configuration(service)) + assert not any( + command[:4] == ["gcloud", "run", "revisions", "describe"] + and command[4] == STAGED_SOURCE_REVISION + for command in calls + ) + + +@pytest.mark.parametrize( + "env_overrides", + [ + {"WORKFLOW_TARGET": "HK"}, + {"WORKFLOW_TARGET": "SG"}, + {"SOURCE_COMMIT": "c" * 40}, + {key: value for key, value in HISTORY.items() if key != "WORKFLOW_TARGET"}, + {"ACCOUNT_SNAPSHOT_ENABLED_INPUT": "true"}, + {"ACCOUNT_SNAPSHOT_ENABLED_INPUT": "false"}, + ], +) +def test_probe_snapshot_source_requires_exact_candidate_paper_and_empty_settings(env_overrides): + env = {**_main_env(), "SOURCE_COMMIT": PROBE_SNAPSHOT_CANDIDATE} + env.update(env_overrides) + calls = [] + with pytest.raises(admission.AdmissionError, match="not approved"): + admission.prepare_image_only_staging( + service="paper-service", project="synthetic-project", region="synthetic-region", + service_json={}, env=env, image_commit=env["SOURCE_COMMIT"], + run=lambda command: calls.append(command) or "", + ) + assert calls == [] + + +def test_probe_snapshot_candidate_cannot_use_http_staged_template_exception(): + service = _http_snapshot_service( + history={key: value for key, value in HISTORY.items() if key != "WORKFLOW_TARGET"} + ) + calls = [] + def run(command): + calls.append(list(command)) + if command[:3] == ["gcloud", "run", "revisions"]: + assert command[4] == "serving-rev" + return json.dumps(_revision("serving-rev", SERVING, "serving-image")) + raise AssertionError(command) + + with pytest.raises(admission.AdmissionError, match="template does not match"): + admission.prepare_image_only_staging( + service="paper-service", project="synthetic-project", region="synthetic-region", + service_json=service, + env={**_main_env(), "SOURCE_COMMIT": PROBE_SNAPSHOT_CANDIDATE}, + image_commit=PROBE_SNAPSHOT_CANDIDATE, + run=run, + ) + assert not any( + command[:4] == ["gcloud", "run", "revisions", "describe"] + and command[4] == STAGED_SOURCE_REVISION + for command in calls + ) diff --git a/tests/test_runtime_monitor_workflows.py b/tests/test_runtime_monitor_workflows.py index 92617b4..9703dea 100644 --- a/tests/test_runtime_monitor_workflows.py +++ b/tests/test_runtime_monitor_workflows.py @@ -121,3 +121,25 @@ def test_heartbeat_script_does_not_import_project_runtime_dependencies() -> None script = (ROOT / "scripts/execution_report_heartbeat.py").read_text() assert "from runtime_config_support import" not in script + + +def test_paper_snapshot_sync_uses_internal_probe_and_preserves_heartbeat_failure(): + workflow = (ROOT / ".github/workflows/execution-report-heartbeat.yml").read_text() + script = (ROOT / "scripts/record_daily_account_snapshot.py").read_text() + account_step = workflow.index("id: account_history") + gcloud_setup = workflow.index("- name: Set up gcloud") + heartbeat = workflow.index("- name: Check recent execution report") + final_failure = workflow.index("- name: Fail after completing heartbeat checks") + + assert 'cron: "20 22 * * *"' in workflow + assert gcloud_setup < account_step < heartbeat < final_failure + assert "continue-on-error: true" in workflow[account_step:heartbeat] + assert "if: ${{ !cancelled() && matrix.target.label == 'PAPER' && vars.ACCOUNT_HISTORY_RECORDING_ENABLED == 'true' }}" in workflow + assert "ACCOUNT_HISTORY_EXPECTED_SOURCE_BINDING_ID: ${{ vars.ACCOUNT_HISTORY_EXPECTED_SOURCE_BINDING_ID }}" in workflow + assert "matrix.target.service" not in workflow + assert "matrix.target.region" not in workflow + assert "steps.account_history.outcome == 'failure'" in workflow[final_failure:] + assert "https://cloudscheduler.googleapis.com/v1/" in script + assert "{config.service_url}/probe" in script + assert "config.scheduler_resource}:run" in script + assert 'ACCOUNT_HISTORY_SERVICE_URL' in workflow diff --git a/tests/test_sync_cloud_run_env_workflow.sh b/tests/test_sync_cloud_run_env_workflow.sh index 9feb691..c7163a0 100644 --- a/tests/test_sync_cloud_run_env_workflow.sh +++ b/tests/test_sync_cloud_run_env_workflow.sh @@ -447,6 +447,7 @@ assert 'git archive "${SOURCE_COMMIT}"' not in job assert "git archive ${{" not in job assert "0b939723c1db3ef59175535998b470cbcd4b8824" in job assert "d8314a61df697cae1dd03a78ddc5c2fc4179ec67" in job +assert "a2921d157efb887e9210fad6734ca040ea6e5293" in job assert '[ "${GITHUB_REPOSITORY:-}" != "QuantStrategyLab/LongBridgePlatform" ]' in job assert '[ "${SOURCE_COMMIT}" = "${GITHUB_SHA}" ] || [ "${SOURCE_COMMIT}" = "${approved_candidate}" ]' not in job for forbidden in ("sync_plan", "scheduler", "cleanup", "retire", "update-traffic"): @@ -474,6 +475,7 @@ with open(os.environ["STUB_LOG"], "a") as stream: stream.write(json.dumps([command, *args]) + "\\n") approved = "0b939723c1db3ef59175535998b470cbcd4b8824" http_snapshot_candidate = "d8314a61df697cae1dd03a78ddc5c2fc4179ec67" +probe_snapshot_candidate = "a2921d157efb887e9210fad6734ca040ea6e5293" staged_source_revision = "longbridge-quant-paper-service-r36423178119" staged_source_image = ( "asia-east1-docker.pkg.dev/synthetic-project/images/longbridgeplatform/synthetic-paper" @@ -602,11 +604,11 @@ def staged_revision(): if command == "git" and args == ["rev-parse", "HEAD"]: print(os.environ["CHECKOUT_SHA"]) -elif command == "git" and args[:4] == ["fetch", "--depth", "1", "origin"] and args[4] in (approved, http_snapshot_candidate): +elif command == "git" and args[:4] == ["fetch", "--depth", "1", "origin"] and args[4] in (approved, http_snapshot_candidate, probe_snapshot_candidate): pass -elif command == "git" and args[:2] == ["cat-file", "-t"] and args[2] in (approved, http_snapshot_candidate): +elif command == "git" and args[:2] == ["cat-file", "-t"] and args[2] in (approved, http_snapshot_candidate, probe_snapshot_candidate): print("commit") -elif command == "git" and args[:1] == ["rev-parse"] and len(args) == 2 and args[1] in (approved + "^{commit}", http_snapshot_candidate + "^{commit}"): +elif command == "git" and args[:1] == ["rev-parse"] and len(args) == 2 and args[1] in (approved + "^{commit}", http_snapshot_candidate + "^{commit}", probe_snapshot_candidate + "^{commit}"): print(args[1].split("^", 1)[0]) elif command == "git" and args == ["archive", "HEAD"]: print("synthetic tracked source archive") @@ -614,11 +616,13 @@ elif command == "git" and args == ["archive", approved]: print("synthetic candidate archive") elif command == "git" and args == ["archive", http_snapshot_candidate]: print("synthetic HTTP snapshot candidate archive") +elif command == "git" and args == ["archive", probe_snapshot_candidate]: + print("synthetic internal probe snapshot candidate archive") elif command == "git" and args[:1] == ["show"] and len(args) == 2 and ":" in args[1]: sha, name = args[1].split(":", 1) if name not in ("uv.lock", "pyproject.toml", "qsl.toml"): raise SystemExit("unexpected git show") - if sha not in ("a" * 40, serving_sha, approved, http_snapshot_candidate): + if sha not in ("a" * 40, serving_sha, approved, http_snapshot_candidate, probe_snapshot_candidate): raise SystemExit("admission read an unapproved source lock") pin = ("f" * 40) if sha == serving_sha and os.environ.get("BAD_SERVING_LOCK") == "1" else ues print(declaration(name, pin), end="") @@ -698,6 +702,7 @@ else: } candidate = "0b939723c1db3ef59175535998b470cbcd4b8824" http_snapshot_candidate = "d8314a61df697cae1dd03a78ddc5c2fc4179ec67" + probe_snapshot_candidate = "a2921d157efb887e9210fad6734ca040ea6e5293" staged_source_revision = "longbridge-quant-paper-service-r36423178119" staged_source_image = ( "asia-east1-docker.pkg.dev/synthetic-project/images/longbridgeplatform/" @@ -836,6 +841,34 @@ else: code, calls = execute(**overrides) assert code != 0 and calls == [], overrides cases += 1 + code, calls = execute(SOURCE_COMMIT=probe_snapshot_candidate, CHECKOUT_SHA="a" * 40) + assert code == 0, (code, calls) + updates = [call for call in calls if call[:4] == ["gcloud", "run", "services", "update"]] + assert len(updates) == 1 + assert "--no-traffic" in updates[0] + assert "--update-env-vars=" not in " ".join(updates[0]) + assert updates[0][-1] == "--revision-suffix=r123" + assert f"--image=registry.invalid/synthetic-project/synthetic-images/longbridgeplatform/synthetic-paper@{base['IMAGE_DIGEST']}" in updates[0] + assert f"--update-labels=commit-sha={probe_snapshot_candidate},github-run-id=123" in updates[0] + assert ["docker", "build", "--pull", "-t", + f"registry.invalid/synthetic-project/synthetic-images/longbridgeplatform/synthetic-paper:{probe_snapshot_candidate}-123", "-"] in calls + assert ["git", "fetch", "--depth", "1", "origin", probe_snapshot_candidate] in calls + assert ["git", "archive", probe_snapshot_candidate] in calls + assert ["git", "archive", "HEAD"] not in calls + assert not any(call[:4] == ["gcloud", "run", "revisions", "describe"] + and call[4] == staged_source_revision for call in calls) + cases += 1 + for overrides in ( + {"SOURCE_COMMIT": probe_snapshot_candidate, "WORKFLOW_TARGET": "HK"}, + {"SOURCE_COMMIT": probe_snapshot_candidate, "WORKFLOW_TARGET": "SG"}, + {"SOURCE_COMMIT": "c" * 40}, + {"SOURCE_COMMIT": probe_snapshot_candidate, **history}, + {"SOURCE_COMMIT": probe_snapshot_candidate, "ACCOUNT_SNAPSHOT_ENABLED_INPUT": "true"}, + {"SOURCE_COMMIT": probe_snapshot_candidate, "ACCOUNT_SNAPSHOT_ENABLED_INPUT": "false"}, + ): + code, calls = execute(**overrides) + assert code != 0 and calls == [], overrides + cases += 1 for overrides in ( { "GITHUB_REF": "refs/heads/codex/natural-cycle-history-20260928",