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
3 changes: 3 additions & 0 deletions .github/workflows/execution-report-heartbeat.yml
Original file line number Diff line number Diff line change
Expand Up @@ -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
4 changes: 3 additions & 1 deletion docs/daily_runtime_projection.md
Original file line number Diff line number Diff line change
Expand Up @@ -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 送达不在这一步。
150 changes: 142 additions & 8 deletions scripts/publish_daily_runtime_projection.py
Original file line number Diff line number Diff line change
Expand Up @@ -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):
Expand Down Expand Up @@ -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],
*,
Expand All @@ -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")
Expand Down Expand Up @@ -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:
Expand All @@ -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
Expand All @@ -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
Expand Down
Loading
Loading