From 159064c404ca4873abf83d6fa21014319f9dcd86 Mon Sep 17 00:00:00 2001 From: Pigbibi <20649888+Pigbibi@users.noreply.github.com> Date: Tue, 29 Sep 2026 04:13:00 +0800 Subject: [PATCH 1/2] feat: publish recorded account snapshots to QRS Co-Authored-By: Codex --- .../workflows/execution-report-heartbeat.yml | 3 + docs/account_snapshot_history.md | 6 +- scripts/record_daily_account_snapshot.py | 191 ++++++++++++- tests/test_daily_account_snapshot.py | 250 +++++++++++++++++- 4 files changed, 435 insertions(+), 15 deletions(-) diff --git a/.github/workflows/execution-report-heartbeat.yml b/.github/workflows/execution-report-heartbeat.yml index 9b4682d..909a800 100644 --- a/.github/workflows/execution-report-heartbeat.yml +++ b/.github/workflows/execution-report-heartbeat.yml @@ -145,6 +145,9 @@ jobs: 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 diff --git a/docs/account_snapshot_history.md b/docs/account_snapshot_history.md index 5cd8317..a246408 100644 --- a/docs/account_snapshot_history.md +++ b/docs/account_snapshot_history.md @@ -13,6 +13,8 @@ `PAPER` 只表示这份清单配置的范围,不能据此推断券商账户身份。脚本再要求期望 scope 精确为 `PAPER`,服务根必须是没有 userinfo、query、fragment 和额外 path 的 HTTPS `*.run.app`,并且只 GET 该源站的 `/account-snapshot`。GCS 前缀最后一段必须是 `account_snapshots`,不能落在 execution report 路径上。缺任何一项就失败,不补默认值。 +可选 QRS 发布仍默认关闭;只有 `ACCOUNT_FACTS_SYNC_ENABLED` 精确为 `true` 时才使用 `ACCOUNT_FACTS_SYNC_URL` 与专用 `ACCOUNT_FACTS_SYNC_TOKEN`。URL 必须是 HTTPS 的精确 `/api/account-facts/sync`,不接受 userinfo、端口、query、fragment 或其他路径;POST 不跟随重定向,设置超时并限制响应体大小。token 只作为该 workflow step 的环境变量和 Authorization header 使用,不打印。 + OIDC 使用 heartbeat 里已经配置的 gcloud:`gcloud auth print-identity-token --audiences=<服务根> --quiet`。stdout 只留在内存,不打印,也不放进参数;失败只报短类别。HTTP 有超时、不跟随重定向、不重试。观察时钟在响应收齐之后读取。 ## 保存什么 @@ -23,6 +25,8 @@ OIDC 使用 heartbeat 里已经配置的 gcloud:`gcloud auth print-identity-to 路径是 `{prefix}/{target_id}/{source_binding_id}/{YYYY-MM-DD}.json`。目标与来源绑定分成不同路径段,不把两段来源拼进同一个对象。`create_text` 只创建:已有对象时结果是 `already_recorded`,不覆盖、不重新拉取。存储结果不明则停止,不写占位点,也不重试。错误输出只有短类别。 +仅当本次 `create_text` 明确返回新建成功时,脚本才把完全相同的历史 JSON body POST 到 QRS。`already_recorded` 明确跳过发布,不能把这次新读取的内容冒充为已保存对象;`store_unknown` 不 POST。QRS 发布状态与历史记录状态分开输出:发布拒绝或结果未知不会撤销已写入历史;未知 POST 不自动重试。QRS `ok=true` 且回读的目标、观察日、观察结束时间匹配,只表示接收端确认保存,不证明页面已经展示或数据完成物理账户身份核验。接收端按其可信配置绑定目标与来源,调用方不传账户 key 或身份结论。 + 这份记录不是 TWR,不是收益率,也不授予 live 权限。 ## main 上的 PAPER 镜像暂存 @@ -33,4 +37,4 @@ OIDC 使用 heartbeat 里已经配置的 gcloud:`gcloud auth print-identity-to ## 尚未启用 -本轮没有打开 GitHub variable,没有新增 secret 或 IAM,也没有真实采样。以后若要启用,需要先确认 paper 服务上的部署版本、来源读取权限、bucket 保留策略和上述变量,并读回第一份对象。在那之前,仓库里的步骤保持关闭。 +本轮没有打开 GitHub variable,没有新增 secret 或 IAM,也没有真实采样。本次代码不配置生产 QRS URL/token,也不改变该关闭状态。以后若要启用,需要先确认 paper 服务上的部署版本、来源读取权限、bucket 保留策略和上述变量,并读回第一份对象。在那之前,仓库里的步骤保持关闭。 diff --git a/scripts/record_daily_account_snapshot.py b/scripts/record_daily_account_snapshot.py index 639540c..dc449cb 100644 --- a/scripts/record_daily_account_snapshot.py +++ b/scripts/record_daily_account_snapshot.py @@ -26,7 +26,12 @@ 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$") +_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]$") @@ -42,6 +47,8 @@ class DailyAccountRecordResult: status: str category: str = "" + publish_status: str = "disabled" + publish_category: str = "" class _Rejected(Exception): @@ -65,11 +72,12 @@ def record_daily_account_snapshot( http_get: Callable[..., Any], open_store: Callable[[str], Any], now_reader: Callable[[], datetime], + http_post: Callable[..., Any] | None = None, ) -> DailyAccountRecordResult: """Validate one snapshot response and create the first object for that day.""" if str(env.get("ACCOUNT_HISTORY_RECORDING_ENABLED") or "").strip() != "true": - return DailyAccountRecordResult("disabled") + return DailyAccountRecordResult("disabled", publish_status="disabled") try: config = _config(env) except _Rejected as rejected: @@ -98,12 +106,151 @@ def record_daily_account_snapshot( try: created = open_store(config.project_id).create_text(uri, body, "application/json") except Exception: - return DailyAccountRecordResult("error", "store_unknown") + return DailyAccountRecordResult( + "error", "store_unknown", "skipped_store_unknown" + ) if created is True: - return DailyAccountRecordResult("recorded") + publish_status, publish_category = _publish_account_facts( + env, body, http_post=http_post or _http_post + ) + return DailyAccountRecordResult( + "recorded", "", publish_status, publish_category + ) if created is False: - return DailyAccountRecordResult("already_recorded") - return DailyAccountRecordResult("error", "store_unknown") + return DailyAccountRecordResult( + "already_recorded", "", "skipped_already_recorded" + ) + return DailyAccountRecordResult( + "error", "store_unknown", "skipped_store_unknown" + ) + + +def _publish_account_facts( + env: Mapping[str, str], + body: str, + *, + http_post: Callable[..., Any], +) -> tuple[str, str]: + if str(env.get("ACCOUNT_FACTS_SYNC_ENABLED") or "").strip() != "true": + return "disabled", "" + try: + 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") + except _Rejected as rejected: + return "rejected", rejected.category + + 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, + ) + 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 _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 def _config(env: Mapping[str, str]) -> _Config: @@ -323,6 +470,27 @@ def _http_get(url: str, *, headers: Mapping[str, str], timeout: float): return requests.get(url, headers=dict(headers), timeout=timeout, allow_redirects=False) +def _http_post( + url: str, + *, + headers: Mapping[str, str], + data: bytes, + timeout: float, + allow_redirects: bool, + stream: bool, +): + import requests + + return requests.post( + url, + headers=dict(headers), + data=data, + timeout=timeout, + allow_redirects=allow_redirects, + stream=stream, + ) + + def _open_store(project_id: str): from quant_platform_kit.cloud import get_object_store @@ -335,6 +503,7 @@ def main( 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, now_reader: Callable[[], datetime] | None = None, ) -> int: @@ -348,13 +517,19 @@ def main( 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, now_reader=now_reader or (lambda: datetime.now(timezone.utc)), ) - if result.status == "error": - print(f"error: {result.category}") + 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"}: return 1 - print(result.status) return 0 diff --git a/tests/test_daily_account_snapshot.py b/tests/test_daily_account_snapshot.py index eed33b9..6833a39 100644 --- a/tests/test_daily_account_snapshot.py +++ b/tests/test_daily_account_snapshot.py @@ -23,7 +23,9 @@ 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" +QRS_TOKEN = "synthetic-qrs-token" LEAK = "raw-secret-value" @@ -92,14 +94,30 @@ def __init__(self, status_code, payload=None, content=None): content = json.dumps(payload if payload is not None else _payload()).encode() self.content = content + def iter_content(self, chunk_size=8192): + yield self.content + + def close(self): + pass + class _Spies: - def __init__(self, response=None, created=True, explode_store=False): + def __init__(self, response=None, created=True, explode_store=False, post_response=None): 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(), + }, + ) def fetch_id_token(self, audience): self.calls.append(("token", audience)) @@ -111,6 +129,12 @@ def http_get(self, url, *, headers, timeout): raise self.response return self.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 open_store(self, project_id): self.calls.append(("store", project_id)) if self.explode_store: @@ -134,6 +158,7 @@ def _record(env=None, spies=None, now=NOW, now_reader=None): 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), ) @@ -282,6 +307,219 @@ def test_record_stores_only_whitelisted_multi_currency_fields(): assert "net_assets_total" not in saved +def test_publishes_exact_created_history_body_once_to_qrs(): + result, spies = _record( + env=_env( + ACCOUNT_FACTS_SYNC_ENABLED="true", + ACCOUNT_FACTS_SYNC_URL=QRS_SYNC_URL, + ACCOUNT_FACTS_SYNC_TOKEN=QRS_TOKEN, + ) + ) + + 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 + + +def test_store_unknown_never_posts_to_qrs(): + spies = _Spies(created=None) + 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.category == "store_unknown" + assert [name for name, *_ in spies.calls].count("post") == 0 + + +def test_already_recorded_does_not_publish_newly_fetched_body(): + spies = _Spies(created=False) + 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 == "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, + ) + + 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}")) + 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 == "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_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_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, + ) + + 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 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, + ) + + 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, + ) + + 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 == "recorded" + assert result.publish_status == "rejected" + assert [name for name, *_ in spies.calls].count("post") == 0 + + 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) @@ -343,7 +581,7 @@ def test_cli_disabled_and_success_use_injected_dependencies(capsys): now_reader=lambda: NOW, ) assert disabled == 0 - assert capsys.readouterr().out.strip() == "disabled" + assert capsys.readouterr().out.strip() == "record=disabled account_facts_publish=disabled" spies = _Spies() code = main( @@ -356,7 +594,7 @@ def test_cli_disabled_and_success_use_injected_dependencies(capsys): ) captured = capsys.readouterr() assert code == 0 - assert captured.out.strip() == "recorded" + assert captured.out.strip() == "record=recorded account_facts_publish=disabled" assert SECRET_TOKEN not in captured.out @@ -370,7 +608,7 @@ def test_cli_error_is_a_short_category(capsys): now_reader=lambda: NOW, ) assert code == 1 - assert capsys.readouterr().out.strip() == "error: config_invalid" + assert capsys.readouterr().out.strip() == "record=error:config_invalid account_facts_publish=disabled" def test_gcloud_token_wrapper_success_failure_and_timeout(monkeypatch, capsys): @@ -457,7 +695,7 @@ def succeed(args, **kwargs): ) captured = capsys.readouterr() assert code == 0 - assert captured.out.strip() == "recorded" + 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" @@ -547,7 +785,7 @@ def test_malformed_url_cli_stays_a_short_category(capsys): ) captured = capsys.readouterr() assert code == 1 - assert captured.out.strip() == "error: config_invalid" + 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 From 8c7d7ca48a76371b4415d330948ac3c014a7eed2 Mon Sep 17 00:00:00 2001 From: Pigbibi <20649888+Pigbibi@users.noreply.github.com> Date: Tue, 29 Sep 2026 05:12:40 +0800 Subject: [PATCH 2/2] fix: publish stored daily runtime projections to QRS Co-Authored-By: Codex --- .../workflows/execution-report-heartbeat.yml | 3 + docs/daily_runtime_projection.md | 4 +- scripts/publish_daily_runtime_projection.py | 150 +++++++++++- .../test_publish_daily_runtime_projection.py | 215 +++++++++++++++++- 4 files changed, 354 insertions(+), 18 deletions(-) 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