Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
14 changes: 13 additions & 1 deletion .github/workflows/execution-report-heartbeat.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down
2 changes: 1 addition & 1 deletion .github/workflows/runtime-target-lifecycle.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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 }}
Expand Down
8 changes: 5 additions & 3 deletions docs/account_snapshot_history.md
Original file line number Diff line number Diff line change
Expand Up @@ -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 改变则失败,不把部分结果当完整结果。

Expand Down
4 changes: 2 additions & 2 deletions docs/daily_runtime_projection.md
Original file line number Diff line number Diff line change
Expand Up @@ -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 业务失败后仍可以写出异常投影。

Expand Down
21 changes: 16 additions & 5 deletions scripts/publish_daily_runtime_projection.py
Original file line number Diff line number Diff line change
Expand Up @@ -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(
Expand Down Expand Up @@ -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(
Expand Down
2 changes: 1 addition & 1 deletion scripts/record_daily_account_snapshot.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
12 changes: 12 additions & 0 deletions scripts/render_runtime_target_matrix.py
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand All @@ -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
Expand Down
23 changes: 11 additions & 12 deletions tests/test_daily_account_snapshot.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down Expand Up @@ -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")])
Expand All @@ -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 == []
Expand Down
14 changes: 13 additions & 1 deletion tests/test_publish_daily_runtime_projection.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down Expand Up @@ -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)
Expand Down
15 changes: 14 additions & 1 deletion tests/test_runtime_monitor_workflows.py
Original file line number Diff line number Diff line change
Expand Up @@ -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()
Expand Down Expand Up @@ -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


Expand Down
19 changes: 18 additions & 1 deletion tests/test_runtime_target_matrix.py
Original file line number Diff line number Diff line change
Expand Up @@ -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": []}),
Expand All @@ -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
Expand Down
Loading