Files
market_intel/apps/worker/settle.py
T
2026-08-01 10:33:19 +08:00

188 lines
5.7 KiB
Python

"""到期结算回填:从本地指数或 OKX 历史指数锚定 settle_index_px。"""
from __future__ import annotations
import logging
import time
from typing import Any
from apps.collector.okx_rest import OkxRestClient, safe_float
from packages.db.repository import Repository
from packages.domain.expiry import expiry_ms_from_ymd
log = logging.getLogger("worker.settle")
# 本地 index_ticks 与到期时刻的最大偏离
_LOCAL_MAX_DELTA_MS = 15 * 60 * 1000
def list_expiry_ymds_needing_settle(repo: Repository, *, now_ms: int | None = None) -> list[str]:
"""option_quotes 中已到期且尚未写入 settlements 的 expiry_ymd。"""
now = int(now_ms if now_ms is not None else time.time() * 1000)
rows = repo.conn.execute(
"""
SELECT DISTINCT expiry_ymd FROM option_quotes
WHERE expiry_ymd IS NOT NULL AND expiry_ymd != ''
ORDER BY expiry_ymd ASC
"""
).fetchall()
out: list[str] = []
for r in rows:
ymd = str(r["expiry_ymd"])
try:
settle_ts = expiry_ms_from_ymd(ymd)
except ValueError:
continue
if settle_ts > now:
continue
if repo.get_settlement(ymd) is not None:
continue
out.append(ymd)
return out
def resolve_settle_index(
repo: Repository,
*,
expiry_ymd: str,
underlying: str,
index_inst_id: str,
exchange: str = "okx",
client: OkxRestClient | None = None,
) -> dict[str, Any] | None:
"""
解析到期指数。优先本地 index_ticks 最近点;否则 OKX 历史指数 K 线。
"""
settle_ts = expiry_ms_from_ymd(expiry_ymd)
local = repo.nearest_index_tick(
underlying=underlying,
target_ts_ms=settle_ts,
max_delta_ms=_LOCAL_MAX_DELTA_MS,
)
if local is not None:
return {
"expiry_ymd": expiry_ymd,
"settle_ts_ms": settle_ts,
"settle_index_px": float(local["index_px"]),
"exchange": exchange,
"underlying": underlying,
"source": "index_ticks",
"source_ts_ms": int(local["ts_ms"]),
}
own_client = client is None
cli = client or OkxRestClient()
try:
px, src_ts = cli.fetch_index_at(index_inst_id, settle_ts)
if px is None:
return None
return {
"expiry_ymd": expiry_ymd,
"settle_ts_ms": settle_ts,
"settle_index_px": float(px),
"exchange": exchange,
"underlying": underlying,
"source": "okx_history_index",
"source_ts_ms": src_ts,
}
finally:
if own_client:
cli.close()
def backfill_settlements(
repo: Repository,
*,
underlying: str,
index_inst_id: str,
client: OkxRestClient | None = None,
ymds: list[str] | None = None,
now_ms: int | None = None,
) -> dict[str, Any]:
"""回填到期锚点;返回 {filled, skipped, pending, errors}。"""
targets = ymds if ymds is not None else list_expiry_ymds_needing_settle(repo, now_ms=now_ms)
filled: list[str] = []
skipped: list[str] = []
errors: list[str] = []
own_client = client is None
cli = client
try:
for ymd in targets:
if repo.get_settlement(ymd) is not None:
skipped.append(ymd)
continue
try:
if cli is None:
cli = OkxRestClient()
row = resolve_settle_index(
repo,
expiry_ymd=ymd,
underlying=underlying,
index_inst_id=index_inst_id,
client=cli,
)
if row is None:
errors.append(f"{ymd}: settle index unavailable")
continue
repo.upsert_settlement(
expiry_ymd=row["expiry_ymd"],
settle_ts_ms=int(row["settle_ts_ms"]),
settle_index_px=float(row["settle_index_px"]),
exchange=str(row["exchange"]),
underlying=str(row["underlying"]),
)
filled.append(ymd)
log.info(
"settled %s index=%.4f source=%s",
ymd,
row["settle_index_px"],
row.get("source"),
)
except Exception as e: # noqa: BLE001
errors.append(f"{ymd}: {e}")
log.exception("backfill %s failed", ymd)
finally:
if own_client and cli is not None:
cli.close()
return {"filled": filled, "skipped": skipped, "errors": errors, "targets": targets}
def ensure_settlements_for_ymds(
repo: Repository,
ymds: list[str],
*,
underlying: str,
index_inst_id: str,
client: OkxRestClient | None = None,
now_ms: int | None = None,
) -> dict[str, Any]:
"""对给定到期日尽量回填(未到期的跳过)。"""
now = int(now_ms if now_ms is not None else time.time() * 1000)
due = []
for ymd in sorted(set(ymds)):
try:
if expiry_ms_from_ymd(ymd) <= now:
due.append(ymd)
except ValueError:
continue
return backfill_settlements(
repo,
underlying=underlying,
index_inst_id=index_inst_id,
client=client,
ymds=due,
now_ms=now,
)
# re-export for typing clarity
__all__ = [
"backfill_settlements",
"ensure_settlements_for_ymds",
"list_expiry_ymds_needing_settle",
"resolve_settle_index",
"safe_float",
]