From 8db9df20f4d68547a59ba17dd433f51f00475a0f Mon Sep 17 00:00:00 2001 From: Pigbibi <20649888+Pigbibi@users.noreply.github.com> Date: Wed, 30 Sep 2026 03:00:19 +0800 Subject: [PATCH 1/2] fix(runtime): archive daily projections by observation time Co-Authored-By: Codex --- docs/daily_runtime_projection.md | 4 ++-- scripts/publish_daily_runtime_projection.py | 21 ++++++++++++++----- .../test_publish_daily_runtime_projection.py | 14 ++++++++++++- 3 files changed, 31 insertions(+), 8 deletions(-) diff --git a/docs/daily_runtime_projection.md b/docs/daily_runtime_projection.md index b50fe58..4272934 100644 --- a/docs/daily_runtime_projection.md +++ b/docs/daily_runtime_projection.md @@ -80,9 +80,9 @@ projected = project_listed_reports( 报告列表和单份读取走已安装的存储客户端,带 20 秒时限、字节上限,并且 `retry=None`。列表扫描上限是 256 条,和最终保留的 20 条分开;预算内按 `updated` 取最新 20 条。扫描达到上限,或已发现超过 20 条因而只保留最新 20 条时,都记 `listing truncated`,`completeness` 为 incomplete。调度身份沿用 heartbeat 对服务主机和策略 run URI 的既有判断,读取入口是现有的 scheduler job list。 -对象路径是 `{prefix}/longbridge/paper/{business_date}.json`。`business_date` 和时区来自该目标的有效调度。文件内容就是 `project_listed_reports` 的原 JSON,不是资产曲线,也不是成交账本。`fills` 仍是未接通。上传使用现有存储客户端,`if_generation_match=0`、`timeout=20`、`retry=None`;同日对象已存在则不覆盖。列表或读取不完整时,对象里的 `read_errors` 和 `completeness` 保持不完整,不把这一天写成正常休市。 +新观测对象路径是 `{prefix}/longbridge/paper/{business_date}/{observedUTC}.json`,其中 `observedUTC` 为 `YYYYMMDDTHHMMSSffffffZ`。既有 `{business_date}.json` 对象作为旧记录保留,不读取、覆盖或删除。`business_date` 和时区来自该目标的有效调度;观测版本使用本次投影的 UTC `observed_at`。文件内容就是 `project_listed_reports` 的原 JSON,不是资产曲线,也不是成交账本。`fills` 仍是未接通。上传使用现有存储客户端,`if_generation_match=0`、`timeout=20`、`retry=None`;同一观测版本冲突时不 POST。每次新观测独立保存,保留当时的 `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。 +可选的 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 与该对象完全相同的 JSON;同一观测版本冲突及 GCS 写入结果未知均不 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 业务失败后仍可以写出异常投影。 diff --git a/scripts/publish_daily_runtime_projection.py b/scripts/publish_daily_runtime_projection.py index 1ea7dbe..949f5a0 100644 --- a/scripts/publish_daily_runtime_projection.py +++ b/scripts/publish_daily_runtime_projection.py @@ -411,15 +411,26 @@ def _read_bounded(client: Any, uri: str) -> tuple[dict[str, Any] | None, str | N return parsed, None -def _object_uri(prefix: str, business_date: str) -> tuple[str, str, str]: +def _object_uri( + prefix: str, + business_date: str, + observed_at: dt.datetime, +) -> tuple[str, str, str]: root = prefix[len("gs://") :] bucket, _, base = root.partition("/") - name = f"{base}/longbridge/paper/{business_date}.json" + observation_version = observed_at.astimezone(dt.timezone.utc).strftime("%Y%m%dT%H%M%S%fZ") + name = f"{base}/longbridge/paper/{business_date}/{observation_version}.json" return bucket, name, f"gs://{bucket}/{name}" -def _upload(client: Any, prefix: str, business_date: str, body: str) -> str: - bucket_name, object_name, uri = _object_uri(prefix, business_date) +def _upload( + client: Any, + prefix: str, + business_date: str, + observed_at: dt.datetime, + body: str, +) -> str: + bucket_name, object_name, uri = _object_uri(prefix, business_date, observed_at) blob = client.bucket(bucket_name).blob(object_name) try: blob.upload_from_string( @@ -668,7 +679,7 @@ def read_payload( 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) + uploaded = _upload(store, prefix, business_date, observed, body) if uploaded == "already_recorded": return "already_recorded", business_date, "skipped_existing" sync_status = _sync_runtime_daily( diff --git a/tests/test_publish_daily_runtime_projection.py b/tests/test_publish_daily_runtime_projection.py index 4f5e480..295c7ab 100644 --- a/tests/test_publish_daily_runtime_projection.py +++ b/tests/test_publish_daily_runtime_projection.py @@ -321,7 +321,7 @@ def test_projections_cover_quiet_closed_missing_and_prior_unknown(monkeypatch) - assert status == "recorded" assert business_date == "2026-09-28" assert client.buckets == ["paper-bucket"] - assert client._bucket.names == ["runtime_daily/longbridge/paper/2026-09-28.json"] + assert client._bucket.names == ["runtime_daily/longbridge/paper/2026-09-28/20260928T084000000000Z.json"] assert stored["records"][0]["status"] == "no_submission" assert stored["records"][0]["runs"][0]["run_id"] == "run-1" assert "secret" not in json.dumps(stored) @@ -392,6 +392,18 @@ def test_list_read_and_truncation_stay_incomplete(monkeypatch, capsys) -> None: assert truncated["completeness"] == "incomplete" +def test_observation_paths_are_versioned_by_utc_timestamp() -> None: + first = dt.datetime(2026, 9, 28, 16, 40, tzinfo=HK) + next_observation = first + dt.timedelta(minutes=1) + + _, first_name, _ = publisher._object_uri("gs://paper-bucket/runtime_daily", "2026-09-28", first) + _, next_name, _ = publisher._object_uri("gs://paper-bucket/runtime_daily", "2026-09-28", next_observation) + + assert first_name == "runtime_daily/longbridge/paper/2026-09-28/20260928T084000000000Z.json" + assert next_name == "runtime_daily/longbridge/paper/2026-09-28/20260928T084100000000Z.json" + assert first_name != next_name + + def test_existing_object_is_not_overwritten_and_unknown_write_is_not_retried(monkeypatch) -> None: exists = _Blob(fail=PreconditionFailed("already present")) status, business_date, blob, _, _ = _publish(monkeypatch, _env(), [], blob=exists) From 5e46cd9aabdb117122874f87a17eda7e6a0b7f7d Mon Sep 17 00:00:00 2001 From: Pigbibi <20649888+Pigbibi@users.noreply.github.com> Date: Wed, 30 Sep 2026 03:21:37 +0800 Subject: [PATCH 2/2] fix(accounts): validate scheduler state and isolate manual sampling Co-Authored-By: Codex --- .../workflows/execution-report-heartbeat.yml | 14 ++++++++++- .../workflows/runtime-target-lifecycle.yml | 2 +- docs/account_snapshot_history.md | 8 ++++--- scripts/record_daily_account_snapshot.py | 2 +- scripts/render_runtime_target_matrix.py | 12 ++++++++++ tests/test_daily_account_snapshot.py | 23 +++++++++---------- tests/test_runtime_monitor_workflows.py | 15 +++++++++++- tests/test_runtime_target_matrix.py | 19 ++++++++++++++- 8 files changed, 75 insertions(+), 20 deletions(-) diff --git a/.github/workflows/execution-report-heartbeat.yml b/.github/workflows/execution-report-heartbeat.yml index e493439..722311d 100644 --- a/.github/workflows/execution-report-heartbeat.yml +++ b/.github/workflows/execution-report-heartbeat.yml @@ -3,6 +3,16 @@ name: Execution Report Heartbeat on: workflow_dispatch: inputs: + target: + description: "Target account to check; scheduled runs always use all targets." + required: false + type: choice + default: "all" + options: + - "all" + - "paper" + - "sg" + - "hk" lookback_hours: description: "Report lookback window in hours." required: false @@ -41,9 +51,11 @@ jobs: - name: Render validated matrix from public manifest id: render + env: + MATRIX_TARGET: ${{ inputs.target || 'all' }} run: | set -euo pipefail - python3 scripts/render_runtime_target_matrix.py --profile heartbeat --github-output + python3 scripts/render_runtime_target_matrix.py --profile heartbeat --target "$MATRIX_TARGET" --github-output heartbeat: needs: resolve-matrix diff --git a/.github/workflows/runtime-target-lifecycle.yml b/.github/workflows/runtime-target-lifecycle.yml index b48e260..e5049da 100644 --- a/.github/workflows/runtime-target-lifecycle.yml +++ b/.github/workflows/runtime-target-lifecycle.yml @@ -310,7 +310,7 @@ jobs: - name: Publish lifecycle to the unified control plane if: ${{ steps.selection.outputs.selected == 'true' && steps.lifecycle_ingress.outputs.enabled == 'true' }} - uses: QuantStrategyLab/QuantRuntimeSettings/actions/publish-runtime-target-lifecycle@c87355fed83dc60e051244f7c7fbd5c9c3968f2b + uses: QuantStrategyLab/QuantRuntimeSettings/actions/publish-runtime-target-lifecycle@88cbbfdde1e41f51f0a617324f8ff56c2ceebe02 with: source-id: longbridge.${{ matrix.target.id }} target-id: longbridge.${{ matrix.target.id }} diff --git a/docs/account_snapshot_history.md b/docs/account_snapshot_history.md index f096e6b..eb534ca 100644 --- a/docs/account_snapshot_history.md +++ b/docs/account_snapshot_history.md @@ -4,14 +4,16 @@ ## 何时会写 -四个条件同时成立才调用 `scripts/record_daily_account_snapshot.py`: +以下条件同时成立才调用 `scripts/record_daily_account_snapshot.py`: -- 当前 matrix 目标的 `label` 是 `PAPER` +- 当前 matrix 目标的 `id` 是 `paper`、`hk` 或 `sg`(目标 `label` 分别作为预期 scope) - GitHub variable `ACCOUNT_HISTORY_RECORDING_ENABLED` 精确等于 `true` - 同一步提供 `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` 只表示这份清单配置的范围,不能据此推断券商账户身份。脚本从已校验的 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 请求结果未知时不重触发。缺任何配置就失败,不补默认值。 +`execution-report-heartbeat` 的每日定时运行始终选择完整 manifest matrix。手动 `workflow_dispatch` 可将 `target` 选为 `paper`、`hk` 或 `sg`,只检查该 manifest target;默认 `all` 保持完整 matrix。单目标选择只缩小这次 heartbeat 的账户范围,不启用账户记录或修改任何 target variable;对应 GitHub Environment 的 `ACCOUNT_HISTORY_RECORDING_ENABLED` 仍须单独精确为 `true`,否则采样步骤跳过。 + +target label 只表示这份清单配置的 scope,不能据此推断券商账户身份。脚本从已校验的 runtime target manifest 读取准确 service/region,并先读回 `{service}-probe-scheduler` 的完整 job。只有 job 完整资源名、`ENABLED`、`POST {service_url}/probe`、空 body、Scheduler OIDC service account/audience 及零重试配置全部匹配时,才调用一次 Cloud Scheduler `jobs:run`。暂停的 job 会在触发前拒绝;不直连 internal Cloud Run,也不使用 `/account-snapshot`。Scheduler 请求结果未知时不重触发。缺任何配置就失败,不补默认值。 `ACCOUNT_HISTORY_EXPECTED_SOURCE_BINDING_ID` 必须由部署后的可信读回提供,并与 QRS 的预期绑定完全相同;不能从第一个 GCS 对象反推信任。GCS 前缀最后一段必须是 `account_snapshots`,不能落在 execution report 路径上。脚本只查本次触发 UTC 日期与当前 UTC 日期下准确的 target/source 前缀,限单页和 64KiB 对象。列举与固定 generation 读取都带超时、`retry=None`;若分页截断、对象变大或 generation 改变则失败,不把部分结果当完整结果。 diff --git a/scripts/record_daily_account_snapshot.py b/scripts/record_daily_account_snapshot.py index c847045..339fd7f 100644 --- a/scripts/record_daily_account_snapshot.py +++ b/scripts/record_daily_account_snapshot.py @@ -298,7 +298,7 @@ def _validate_scheduler_job(job: Mapping[str, Any], config: _Config) -> None: body_empty = body is None if ( job.get("name") != config.scheduler_resource - or job.get("state") not in ({"ENABLED", "PAUSED"} if config.target_id == "hk" else {"ENABLED"}) + or job.get("state") != "ENABLED" or target.get("httpMethod") != "POST" or target.get("uri") != f"{config.service_url}/probe" or not body_empty diff --git a/scripts/render_runtime_target_matrix.py b/scripts/render_runtime_target_matrix.py index 5c64d12..fc6626e 100644 --- a/scripts/render_runtime_target_matrix.py +++ b/scripts/render_runtime_target_matrix.py @@ -41,6 +41,11 @@ def main(argv: list[str] | None = None) -> int: choices=sorted(MATRIX_PROFILES), help="Workflow matrix shape to render.", ) + parser.add_argument( + "--target", + default="all", + help="Render all targets (default) or exactly one manifest target ID.", + ) parser.add_argument( "--github-output", action="store_true", @@ -51,6 +56,13 @@ def main(argv: list[str] | None = None) -> int: try: manifest = load_runtime_target_manifest(args.path or default_manifest_path()) matrix = build_github_actions_matrix(manifest, profile=args.profile) + if args.target != "all": + selected = [row for row in matrix["target"] if row["id"] == args.target] + if not selected: + raise RuntimeTargetManifestError( + f"Unknown target {args.target!r}; expected 'all' or a manifest target ID" + ) + matrix["target"] = selected except RuntimeTargetManifestError as exc: print(f"runtime-target matrix render failed: {exc}", file=sys.stderr) return 1 diff --git a/tests/test_daily_account_snapshot.py b/tests/test_daily_account_snapshot.py index f037512..dae4e1d 100644 --- a/tests/test_daily_account_snapshot.py +++ b/tests/test_daily_account_snapshot.py @@ -325,7 +325,7 @@ def test_scheduler_zero_retry_protobuf_defaults_are_accepted(): assert [call[0] for call in spies.session.calls] == ["get", "post"] -@pytest.mark.parametrize(("target_id", "state"), [("hk", "PAUSED"), ("sg", "ENABLED")]) +@pytest.mark.parametrize(("target_id", "state"), [("hk", "ENABLED"), ("sg", "ENABLED")]) def test_sghk_targets_use_exact_manifest_identity_and_publish_matching_history(target_id, state): payload = _history(target_id=target_id) job = _job(target_id, state=state) @@ -355,17 +355,16 @@ def test_sghk_targets_use_exact_manifest_identity_and_publish_matching_history(t assert json.loads(spies.posts[0][1]["data"])["account_scope"] == {"hk": "HK", "sg": "SG"}[target_id] -def test_paused_hk_probe_is_the_only_paused_scheduler_accepted(): - hk_job = _job("hk", state="PAUSED") - hk_result, hk_spies = _record(_env("hk"), _Spies(objects=[_object(_history(target_id="hk"))], job=hk_job)) - assert hk_result.status == "recorded" - assert [call[0] for call in hk_spies.session.calls] == ["get", "post"] +@pytest.mark.parametrize("target_id", ["paper", "sg", "hk"]) +def test_paused_scheduler_is_rejected_before_run_or_gcs_for_every_target(target_id): + result, spies = _record( + _env(target_id), + _Spies(job=_job(target_id, state="PAUSED")), + ) - sg_job = _job("sg", state="PAUSED") - sg_result, sg_spies = _record(_env("sg"), _Spies(job=sg_job)) - assert sg_result.category == "scheduler_job_mismatch" - assert [call[0] for call in sg_spies.session.calls] == ["get"] - assert sg_spies.open_calls == [] + assert result.category == "scheduler_job_mismatch" + assert [call[0] for call in spies.session.calls] == ["get"] + assert spies.open_calls == [] @pytest.mark.parametrize(("target_id", "wrong_scope"), [("hk", "SG"), ("sg", "HK")]) @@ -374,7 +373,7 @@ def test_cross_scope_snapshot_is_never_published(monkeypatch, target_id, wrong_s payload = _history(target_id=target_id, scope=wrong_scope) result, spies = _record( _env(target_id), - _Spies(objects=[_object(payload)], job=_job(target_id, state="PAUSED" if target_id == "hk" else "ENABLED")), + _Spies(objects=[_object(payload)], job=_job(target_id, state="ENABLED")), ) assert result.category == "observation_timeout" assert spies.posts == [] diff --git a/tests/test_runtime_monitor_workflows.py b/tests/test_runtime_monitor_workflows.py index c448af3..84850ae 100644 --- a/tests/test_runtime_monitor_workflows.py +++ b/tests/test_runtime_monitor_workflows.py @@ -17,6 +17,19 @@ def test_execution_report_heartbeat_has_market_neutral_daily_schedule() -> None: assert "pandas-market-calendars" in (ROOT / "uv.lock").read_text() +def test_heartbeat_dispatch_selects_one_manifest_target_and_schedule_uses_all() -> None: + workflow = (ROOT / ".github/workflows/execution-report-heartbeat.yml").read_text() + + assert 'cron: "20 22 * * *"' in workflow + dispatch = workflow[workflow.index(" workflow_dispatch:"):workflow.index(" schedule:")] + assert "target:" in dispatch + assert "type: choice" in dispatch + assert 'default: "all"' in dispatch + assert all(f'- "{target}"' in dispatch for target in ("all", "paper", "sg", "hk")) + assert "MATRIX_TARGET: ${{ inputs.target || 'all' }}" in workflow + assert 'scripts/render_runtime_target_matrix.py --profile heartbeat --target "$MATRIX_TARGET" --github-output' in workflow + + def test_runtime_monitor_workflows_retry_gcp_authentication() -> None: for name in ("execution-report-heartbeat.yml", "runtime-guard.yml"): workflow = (ROOT / ".github/workflows" / name).read_text() @@ -90,7 +103,7 @@ def test_runtime_monitor_workflows_use_frozen_runtime_environment() -> None: heartbeat = (ROOT / ".github/workflows/execution-report-heartbeat.yml").read_text() assert "resolve-matrix:" in heartbeat - assert "render_runtime_target_matrix.py --profile heartbeat --github-output" in heartbeat + assert 'render_runtime_target_matrix.py --profile heartbeat --target "$MATRIX_TARGET" --github-output' in heartbeat assert "matrix: ${{ fromJSON(needs.resolve-matrix.outputs.matrix) }}" in heartbeat diff --git a/tests/test_runtime_target_matrix.py b/tests/test_runtime_target_matrix.py index 5dd5c49..d24e86e 100644 --- a/tests/test_runtime_target_matrix.py +++ b/tests/test_runtime_target_matrix.py @@ -240,6 +240,20 @@ def test_render_script_writes_github_output_and_rejects_bad_manifest( assert written.startswith("matrix=") assert json.loads(written.removeprefix("matrix=")) == {"target": EXPECTED_GUARD} + assert module.main(["--profile", "heartbeat"]) == 0 + complete = json.loads(capsys.readouterr().out) + assert complete == {"target": EXPECTED_HEARTBEAT} + + assert module.main(["--profile", "heartbeat", "--target", "hk"]) == 0 + selected = json.loads(capsys.readouterr().out) + assert selected == {"target": [EXPECTED_HEARTBEAT[1]]} + + assert module.main(["--profile", "heartbeat", "--target", "unknown"]) == 1 + assert "unknown" in capsys.readouterr().err + + assert module.main(["--profile", "heartbeat", "--target", ""]) == 1 + assert capsys.readouterr().err + bad = tmp_path / "bad.json" bad.write_text( json.dumps({"schema_version": 1, "platform_id": "longbridge", "targets": []}), @@ -262,7 +276,10 @@ def test_monitor_and_deploy_workflows_consume_validated_matrix_output(): text = (REPO_ROOT / ".github" / "workflows" / name).read_text(encoding="utf-8") assert "jobs:" in text assert "resolve-matrix:" in text - assert f"render_runtime_target_matrix.py --profile {profile} --github-output" in text + if name == "execution-report-heartbeat.yml": + assert 'render_runtime_target_matrix.py --profile heartbeat --target "$MATRIX_TARGET" --github-output' in text + else: + assert f"render_runtime_target_matrix.py --profile {profile} --github-output" in text assert "needs: resolve-matrix" in text assert "matrix: ${{ fromJSON(needs.resolve-matrix.outputs.matrix) }}" in text assert "vars.RUNTIME_TARGET_ENABLED" in text