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 @@ -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
Expand Down
6 changes: 5 additions & 1 deletion docs/account_snapshot_history.md
Original file line number Diff line number Diff line change
Expand Up @@ -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 有超时、不跟随重定向、不重试。观察时钟在响应收齐之后读取。

## 保存什么
Expand All @@ -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 镜像暂存
Expand All @@ -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 保留策略和上述变量,并读回第一份对象。在那之前,仓库里的步骤保持关闭。
191 changes: 183 additions & 8 deletions scripts/record_daily_account_snapshot.py
Original file line number Diff line number Diff line change
Expand Up @@ -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]$")
Expand All @@ -42,6 +47,8 @@
class DailyAccountRecordResult:
status: str
category: str = ""
publish_status: str = "disabled"
publish_category: str = ""


class _Rejected(Exception):
Expand All @@ -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:
Expand Down Expand Up @@ -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:
Expand Down Expand Up @@ -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

Expand All @@ -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:
Expand All @@ -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


Expand Down
Loading
Loading