diff --git a/.github/workflows/execution-report-heartbeat.yml b/.github/workflows/execution-report-heartbeat.yml index 909a800..ff9d40b 100644 --- a/.github/workflows/execution-report-heartbeat.yml +++ b/.github/workflows/execution-report-heartbeat.yml @@ -156,5 +156,8 @@ jobs: RUNTIME_DAILY_PROJECTION_ENABLED: ${{ vars.RUNTIME_DAILY_PROJECTION_ENABLED }} RUNTIME_DAILY_PROJECTION_GCS_PREFIX: ${{ vars.RUNTIME_DAILY_PROJECTION_GCS_PREFIX }} RUNTIME_DAILY_PROJECTION_TARGET_ID: ${{ vars.RUNTIME_DAILY_PROJECTION_TARGET_ID }} + RUNTIME_DAILY_SYNC_ENABLED: ${{ vars.RUNTIME_DAILY_SYNC_ENABLED }} + RUNTIME_DAILY_SYNC_URL: ${{ vars.RUNTIME_DAILY_SYNC_URL }} + 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 diff --git a/docs/daily_runtime_projection.md b/docs/daily_runtime_projection.md index 53d4739..b50fe58 100644 --- a/docs/daily_runtime_projection.md +++ b/docs/daily_runtime_projection.md @@ -82,6 +82,8 @@ projected = project_listed_reports( 对象路径是 `{prefix}/longbridge/paper/{business_date}.json`。`business_date` 和时区来自该目标的有效调度。文件内容就是 `project_listed_reports` 的原 JSON,不是资产曲线,也不是成交账本。`fills` 仍是未接通。上传使用现有存储客户端,`if_generation_match=0`、`timeout=20`、`retry=None`;同日对象已存在则不覆盖。列表或读取不完整时,对象里的 `read_errors` 和 `completeness` 保持不完整,不把这一天写成正常休市。 +可选的 QRS 日报同步另由 `RUNTIME_DAILY_SYNC_ENABLED` 精确等于 `true` 才开启,并要求 HTTPS 精确路径 `/api/runtime-daily/sync`。它复用既有 `EXECUTION_EVIDENCE_SYNC_TOKEN`,不增加 token 类型;workflow 仅在该开关开启时把 secret 送入当前 step。只允许固定 QRS PAPER 目标 `longbridge-quant-paper-service|russell_top50_leader_rotation|paper`,不由调用者指定或推导账户 key。只有本次 GCS 明确新建后才 POST 与 GCS 完全相同的 JSON;`already_recorded` 与写入结果未知均不 POST。POST 不跟随重定向、不重试,限制超时和响应大小。同步状态独立输出为 `qrs_sync=disabled|recorded|rejected|unknown`;网络结果未知或 QRS 拒绝不会撤销已保存的 GCS 记录,也不表示网页已经展示。成功响应必须回显同一 `business_date` 和 `target_key`,且包含服务端解析出的非空 `account_key` 才报告 `recorded`,但不会记录或输出该账户 key。 + heartbeat 的告警、返回码和交易链不变。这个步骤不调用 heartbeat 的 `main`。workflow 只在 PAPER、变量显式为 `true`、任务未取消且 Google 认证成功时运行;heartbeat 业务失败后仍可以写出异常投影。 -云端还没有设置这个变量,也没有新增 secret。真实报告前缀、桶权限和第一份对象要由主助手在获准身份上验收。bot 送达不在这一步。 +该同步仍默认关闭;本地只读核验发现 LongBridge PAPER 环境尚无 `RUNTIME_DAILY_SYNC_ENABLED` 或 `RUNTIME_DAILY_SYNC_URL` variable。它复用已有 execution-evidence secret 名称,不新增 secret。真实报告前缀、桶权限及第一份对象仍需在获准身份上验收。bot 送达不在这一步。 diff --git a/scripts/publish_daily_runtime_projection.py b/scripts/publish_daily_runtime_projection.py index 5995d2a..4112ccb 100644 --- a/scripts/publish_daily_runtime_projection.py +++ b/scripts/publish_daily_runtime_projection.py @@ -34,6 +34,12 @@ _MAX_REPORT_BYTES = 1_048_576 _MAX_LIST_BYTES = 262_144 _MAX_LIST_SCAN = 256 +_RUNTIME_DAILY_SYNC_PATH = "/api/runtime-daily/sync" +_RUNTIME_DAILY_SYNC_TIMEOUT_SECONDS = 20 +_RUNTIME_DAILY_SYNC_MAX_BODY_BYTES = 64 * 1024 +_RUNTIME_DAILY_SYNC_MAX_RESPONSE_BYTES = 64 * 1024 +_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])?$") +_RUNTIME_DAILY_SYNC_TARGET_KEY = "longbridge-quant-paper-service|russell_top50_leader_rotation|paper" class _Rejected(Exception): @@ -416,6 +422,121 @@ def _upload(client: Any, prefix: str, business_date: str, body: str) -> str: return uri +def _runtime_daily_sync_url(value: str) -> str: + try: + parsed = urlsplit(value.strip()) + port = parsed.port + except ValueError: + raise _Rejected("qrs_sync_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 != _RUNTIME_DAILY_SYNC_PATH + or _HTTPS_HOST.fullmatch(host) is None + ): + raise _Rejected("qrs_sync_config_invalid") + return f"https://{host}{_RUNTIME_DAILY_SYNC_PATH}" + + +def _bounded_sync_json(response: Any) -> Any: + chunks: list[bytes] = [] + total = 0 + iterator = getattr(response, "iter_content", None) + if callable(iterator): + for chunk in iterator(chunk_size=8192): + if not chunk: + continue + if not isinstance(chunk, bytes): + raise _Rejected("qrs_sync_response_invalid") + total += len(chunk) + if total > _RUNTIME_DAILY_SYNC_MAX_RESPONSE_BYTES: + raise _Rejected("qrs_sync_response_too_large") + chunks.append(chunk) + raw = b"".join(chunks) + else: + raw = getattr(response, "content", None) + if not isinstance(raw, bytes) or len(raw) > _RUNTIME_DAILY_SYNC_MAX_RESPONSE_BYTES: + raise _Rejected("qrs_sync_response_invalid") + try: + return json.loads(raw) + except (UnicodeError, json.JSONDecodeError): + raise _Rejected("qrs_sync_response_invalid") from None + + +def _sync_runtime_daily( + env: Mapping[str, str], + body: str, + *, + business_date: str, + target_key: str, + http_post: Callable[..., Any] | None = None, +) -> str: + if env.get("RUNTIME_DAILY_SYNC_ENABLED") != "true": + return "disabled" + if target_key != _RUNTIME_DAILY_SYNC_TARGET_KEY: + return "rejected" + try: + url = _runtime_daily_sync_url(str(env.get("RUNTIME_DAILY_SYNC_URL") or "")) + token = str(env.get("EXECUTION_EVIDENCE_SYNC_TOKEN") or "") + if not token or token != token.strip(): + raise _Rejected("qrs_sync_config_invalid") + except _Rejected: + return "rejected" + encoded = body.encode("utf-8") + if len(encoded) > _RUNTIME_DAILY_SYNC_MAX_BODY_BYTES: + return "rejected" + try: + if http_post is None: + import requests + + http_post = requests.post + response = http_post( + url, + headers={ + "Authorization": f"Bearer {token}", + "Accept": "application/json", + "Content-Type": "application/json", + }, + data=encoded, + timeout=_RUNTIME_DAILY_SYNC_TIMEOUT_SECONDS, + allow_redirects=False, + stream=True, + ) + except Exception: + return "unknown" + status_code = getattr(response, "status_code", None) + if not isinstance(status_code, int): + return "unknown" + if getattr(response, "is_redirect", False) or 300 <= status_code < 400: + return "rejected" + if 400 <= status_code < 500: + return "rejected" + if status_code >= 500 or not 200 <= status_code < 300: + return "unknown" + try: + payload = _bounded_sync_json(response) + except Exception: + return "unknown" + if isinstance(payload, dict) and payload.get("ok") is False: + return "rejected" + if ( + not isinstance(payload, dict) + or payload.get("ok") is not True + or payload.get("stored") is not True + or payload.get("business_date") != business_date + or payload.get("target_key") != target_key + or not isinstance(payload.get("account_key"), str) + or not payload["account_key"] + ): + return "unknown" + return "recorded" + + def publish( env: Mapping[str, str], *, @@ -425,11 +546,12 @@ def publish( list_objects: Callable[..., list[dict[str, Any]]] | None = None, read_payload: Callable[[str], Mapping[str, Any] | None] | None = None, report_globs: Callable[[dt.datetime, dt.datetime], list[str]] | None = None, -) -> tuple[str, str]: - """Return ``(status, business_date)``. Disabled returns ``("disabled", "")``.""" + http_post: Callable[..., Any] | None = None, +) -> tuple[str, str, str]: + """Return record status, date, and independently confirmed QRS sync status.""" if not _enabled(env): - return "disabled", "" + return "disabled", "", "disabled" prefix = _require_config(env) if now.tzinfo is None or now.utcoffset() is None: raise _Rejected("config_invalid") @@ -522,11 +644,23 @@ def read_payload( business_date = records[0].get("business_date") if not isinstance(business_date, str) or not business_date: raise RuntimeError("schedule_unavailable") + # QRS stores the fixed PAPER target with its canonical lower-case scope. + # Keep the actual scope check above case-insensitive and serialize this + # normalized representation once so GCS and the optional POST share bytes. + if records[0].get("target_key") == _RUNTIME_DAILY_SYNC_TARGET_KEY: + records[0]["target"]["account_scope"] = "paper" body = json.dumps(projected, ensure_ascii=True, sort_keys=True, separators=(",", ":")) uploaded = _upload(store, prefix, business_date, body) if uploaded == "already_recorded": - return "already_recorded", business_date - return "recorded", business_date + return "already_recorded", business_date, "skipped_existing" + sync_status = _sync_runtime_daily( + env, + body, + business_date=business_date, + target_key=str(records[0].get("target_key") or ""), + http_post=http_post, + ) + return "recorded", business_date, sync_status def _storage_client() -> Any: @@ -538,7 +672,7 @@ def _storage_client() -> Any: def main(argv: Sequence[str] | None = None) -> int: del argv try: - status, business_date = publish(os.environ, now=dt.datetime.now(dt.timezone.utc)) + status, business_date, sync_status = publish(os.environ, now=dt.datetime.now(dt.timezone.utc)) except _Rejected: print("daily runtime projection rejected") return 2 @@ -549,10 +683,10 @@ def main(argv: Sequence[str] | None = None) -> int: print("daily runtime projection disabled") return 0 if status == "recorded": - print(f"daily runtime projection recorded {business_date}") + print(f"daily runtime projection recorded {business_date}; qrs_sync={sync_status}") return 0 if status == "already_recorded": - print(f"daily runtime projection already recorded {business_date}") + print(f"daily runtime projection already recorded {business_date}; qrs_sync={sync_status}") return 0 print("daily runtime projection failed") return 1 diff --git a/tests/test_publish_daily_runtime_projection.py b/tests/test_publish_daily_runtime_projection.py index 5eb1027..daaa35d 100644 --- a/tests/test_publish_daily_runtime_projection.py +++ b/tests/test_publish_daily_runtime_projection.py @@ -7,7 +7,7 @@ from zoneinfo import ZoneInfo import pytest -from google.api_core.exceptions import PreconditionFailed +from google.api_core.exceptions import GoogleAPICallError, PreconditionFailed from scripts import execution_report_heartbeat as heartbeat from scripts import publish_daily_runtime_projection as publisher @@ -52,6 +52,17 @@ def list_blobs(self, bucket_or_name, **kwargs): return [] +class _SyncResponse: + def __init__(self, status_code=200, payload=None, *, redirect=False): + self.status_code = status_code + self.is_redirect = redirect + self.content = json.dumps(payload or {}).encode() + + def iter_content(self, *, chunk_size): + del chunk_size + yield self.content + + def _open_calendar(calendar: str, **kwargs) -> set[dt.date]: del calendar return {kwargs["start_date"]} @@ -89,6 +100,35 @@ def _env(**overrides: str) -> dict[str, str]: return env +def _sync_env(**overrides: str) -> dict[str, str]: + target = _target(service="longbridge-quant-paper-service") + target.update( + strategy_profile="russell_top50_leader_rotation", + market_timezone="America/New_York", + scheduler={"main_time": "5 16 * * 1-5", "timezone": "America/New_York"}, + ) + env = _env( + RUNTIME_TARGET_JSON=json.dumps(target), + RUNTIME_DAILY_SYNC_ENABLED="true", + RUNTIME_DAILY_SYNC_URL="https://qrs.example.com/api/runtime-daily/sync", + EXECUTION_EVIDENCE_SYNC_TOKEN="synthetic-test-token", + ) + env.update(overrides) + return env + + +def _sync_describe(service: str, *, project: str | None) -> dict: + del project + return _cloud_run( + { + "strategy_profile": "russell_top50_leader_rotation", + "account_scope": "PAPER", + "service_name": service, + "runtime_target_enabled": True, + } + ) + + def _report(**changes) -> dict: payload = { "platform": "longbridge", @@ -176,6 +216,8 @@ def _publish( list_error=False, limit=None, jobs="default", + http_post=None, + describe=_describe, ): blob = blob or _Blob() client = _Client(blob) @@ -203,9 +245,9 @@ def list_jobs(*, project: str | None) -> list: del project return selected - monkeypatch.setattr(heartbeat, "_describe_cloud_run_service", _describe) + monkeypatch.setattr(heartbeat, "_describe_cloud_run_service", describe) monkeypatch.setattr(heartbeat, "_list_scheduler_jobs", list_jobs) - status, business_date = publisher.publish( + status, business_date, sync_status = publisher.publish( env, now=OBSERVED, session_dates_loader=calendar, @@ -213,7 +255,9 @@ def list_jobs(*, project: str | None) -> list: list_objects=list_objects, read_payload=read_payload, report_globs=lambda since, now: ["gs://reports/longbridge/**/2026-09/*.json"], + http_post=http_post, ) + calls["sync_status"] = sync_status return status, business_date, blob, client, calls @@ -237,9 +281,10 @@ def boom(*args, **kwargs): monkeypatch.setattr(heartbeat, "_list_gcs_objects", boom) monkeypatch.setattr(heartbeat, "_cat_gcs_json", boom) monkeypatch.setattr(publisher, "_storage_client", boom) - status, business_date = publisher.publish({"RUNTIME_DAILY_PROJECTION_ENABLED": "false"}, now=OBSERVED) + status, business_date, sync_status = publisher.publish({"RUNTIME_DAILY_PROJECTION_ENABLED": "false"}, now=OBSERVED) assert status == "disabled" assert business_date == "" + assert sync_status == "disabled" def test_bad_config_and_cross_scope_do_not_write(monkeypatch) -> None: @@ -595,7 +640,7 @@ def blob(self, object_name: str): monkeypatch.setattr(heartbeat, "_list_scheduler_jobs", lambda *, project: _jobs_for_env(_env())) upload = _Blob() client, scanner, listed = bind(reports, upload) - status, _ = publisher.publish( + status, _, _ = publisher.publish( _env(), now=OBSERVED, session_dates_loader=_closed_calendar, @@ -622,7 +667,7 @@ def blob(self, object_name: str): huge_upload = _Blob() huge_client, _, _ = bind([huge], huge_upload) monkeypatch.setattr(heartbeat, "_list_scheduler_jobs", lambda *, project: _jobs_for_env(_env())) - huge_status, _ = publisher.publish( + huge_status, _, _ = publisher.publish( _env(), now=OBSERVED, session_dates_loader=_closed_calendar, @@ -644,7 +689,7 @@ def blob(self, object_name: str): ) broken_upload = _Blob() broken_client, _, _ = bind([broken], broken_upload) - broken_status, _ = publisher.publish( + broken_status, _, _ = publisher.publish( _env(), now=OBSERVED, session_dates_loader=_closed_calendar, @@ -865,7 +910,157 @@ def test_disabled_target_is_still_projected(monkeypatch) -> None: assert _stored(blob)["records"][0]["status"] == "missing_report" -def test_workflow_adds_a_paper_projection_after_auth_without_new_secrets() -> None: +def test_runtime_daily_sync_posts_only_new_projection_and_checks_qrs_contract(monkeypatch) -> None: + requests = [] + + def post(url, **kwargs): + requests.append((url, kwargs)) + payload = json.loads(kwargs["data"]) + record = payload["records"][0] + return _SyncResponse(payload={ + "ok": True, + "stored": True, + "business_date": record["business_date"], + "target_key": "longbridge-quant-paper-service|russell_top50_leader_rotation|paper", + "account_key": "synthetic-account-key", + }) + + status, business_date, blob, _, calls = _publish( + monkeypatch, + _sync_env(), + [], + http_post=post, + describe=_sync_describe, + ) + assert status == "recorded" + assert business_date == "2026-09-28" + assert calls["sync_status"] == "recorded" + assert len(requests) == 1 + stored_projection = json.loads(blob.uploads[0]["data"]) + assert stored_projection["records"][0]["target"]["account_scope"] == "paper" + url, kwargs = requests[0] + assert url == "https://qrs.example.com/api/runtime-daily/sync" + assert kwargs["allow_redirects"] is False + assert kwargs["timeout"] == publisher._RUNTIME_DAILY_SYNC_TIMEOUT_SECONDS + assert kwargs["stream"] is True + assert kwargs["data"] == blob.uploads[0]["data"].encode("utf-8") + assert kwargs["headers"]["Authorization"] == "Bearer synthetic-test-token" + + +@pytest.mark.parametrize( + "overrides", + [ + {"EXECUTION_EVIDENCE_SYNC_TOKEN": ""}, + {"RUNTIME_DAILY_SYNC_URL": ""}, + {"RUNTIME_DAILY_SYNC_URL": "https://qrs.example.com/wrong"}, + ], +) +def test_runtime_daily_sync_missing_or_invalid_auth_config_does_not_post(monkeypatch, overrides) -> None: + requests = [] + status, _, blob, _, calls = _publish( + monkeypatch, + _sync_env(**overrides), + [], + http_post=lambda *args, **kwargs: requests.append((args, kwargs)), + describe=_sync_describe, + ) + assert status == "recorded" + assert calls["sync_status"] == "rejected" + assert blob.uploads + assert requests == [] + + +def test_runtime_daily_sync_skips_existing_object_and_unknown_store(monkeypatch) -> None: + requests = [] + existing = _Blob(fail=PreconditionFailed("already exists")) + status, _, _, _, calls = _publish( + monkeypatch, + _sync_env(), + [], + blob=existing, + http_post=lambda *args, **kwargs: requests.append((args, kwargs)), + describe=_sync_describe, + ) + assert status == "already_recorded" + assert calls["sync_status"] == "skipped_existing" + assert requests == [] + + unknown = _Blob(fail=GoogleAPICallError("synthetic write failure")) + with pytest.raises(RuntimeError, match="write_unknown"): + _publish( + monkeypatch, + _sync_env(), + [], + blob=unknown, + http_post=lambda *args, **kwargs: requests.append((args, kwargs)), + describe=_sync_describe, + ) + assert requests == [] + + +@pytest.mark.parametrize( + ("response", "expected"), + [ + (TimeoutError("synthetic timeout"), "unknown"), + (_SyncResponse(status_code=302, redirect=True), "rejected"), + (_SyncResponse(payload={"ok": True, "stored": True, "business_date": "wrong", "target_key": "wrong", "account_key": "synthetic"}), "unknown"), + (_SyncResponse(status_code=401, payload={"error": "unauthorized"}), "rejected"), + ], +) +def test_runtime_daily_sync_failure_is_separate_and_never_retried(monkeypatch, response, expected) -> None: + requests = [] + + def post(*args, **kwargs): + requests.append((args, kwargs)) + if isinstance(response, BaseException): + raise response + return response + + status, _, blob, _, calls = _publish( + monkeypatch, + _sync_env(), + [], + http_post=post, + describe=_sync_describe, + ) + assert status == "recorded" + assert calls["sync_status"] == expected + assert blob.uploads + assert len(requests) == 1 + + +def test_runtime_daily_sync_stays_off_by_default(monkeypatch) -> None: + requests = [] + status, _, blob, _, calls = _publish( + monkeypatch, + _env(), + [], + http_post=lambda *args, **kwargs: requests.append((args, kwargs)), + ) + assert status == "recorded" + assert calls["sync_status"] == "disabled" + assert blob.uploads + assert requests == [] + + +def test_runtime_daily_sync_rejects_noncontract_paper_target_locally(monkeypatch) -> None: + requests = [] + target = _target(service="other-paper-service") + target.update(strategy_profile="russell_top50_leader_rotation", market_timezone="America/New_York") + status, _, blob, _, calls = _publish( + monkeypatch, + _sync_env(RUNTIME_TARGET_JSON=json.dumps(target)), + [], + http_post=lambda *args, **kwargs: requests.append((args, kwargs)), + describe=_sync_describe, + ) + assert status == "recorded" + assert calls["sync_status"] == "rejected" + assert blob.uploads + assert requests == [] + + +def test_workflow_adds_a_paper_projection_after_auth_with_existing_sync_token() -> None: workflow = (ROOT / ".github/workflows/execution-report-heartbeat.yml").read_text(encoding="utf-8") heartbeat_step = workflow.index("name: Check recent execution report") projection = workflow.index("name: Publish daily runtime projection") @@ -878,5 +1073,7 @@ def test_workflow_adds_a_paper_projection_after_auth_without_new_secrets() -> No assert "!cancelled()" in step assert "steps.gcp_auth_primary.outcome == 'success'" in step assert "steps.gcp_auth_retry.outcome == 'success'" in step - assert "secrets." not in step assert "RUNTIME_DAILY_PROJECTION_ENABLED: ${{ vars.RUNTIME_DAILY_PROJECTION_ENABLED }}" in step + assert "RUNTIME_DAILY_SYNC_ENABLED: ${{ vars.RUNTIME_DAILY_SYNC_ENABLED }}" in step + assert "RUNTIME_DAILY_SYNC_URL: ${{ vars.RUNTIME_DAILY_SYNC_URL }}" in step + assert "EXECUTION_EVIDENCE_SYNC_TOKEN: ${{ vars.RUNTIME_DAILY_SYNC_ENABLED == 'true' && secrets.EXECUTION_EVIDENCE_SYNC_TOKEN || '' }}" in step