140 lines
6.2 KiB
Python
140 lines
6.2 KiB
Python
|
|
"""Publish complete enumerations only; retain staging checkpoints and old versions."""
|
||
|
|
|
||
|
|
import asyncio
|
||
|
|
import math
|
||
|
|
import re
|
||
|
|
from urllib.parse import parse_qs, urlparse
|
||
|
|
|
||
|
|
from sqlalchemy import select
|
||
|
|
|
||
|
|
from ..models import CatalogBatch, CatalogDataset, CatalogEntry, CatalogNote, CatalogScope, Job, now
|
||
|
|
from ..worldquant import WqError
|
||
|
|
from .contracts import Scope
|
||
|
|
|
||
|
|
|
||
|
|
def identifier(value):
|
||
|
|
if not isinstance(value, str) or not re.fullmatch(r"[A-Za-z0-9_.-]{1,200}", value):
|
||
|
|
raise WqError("平台目录包含无法识别的 ID,已保留进度", "invalid_response")
|
||
|
|
return value
|
||
|
|
|
||
|
|
|
||
|
|
def label(value):
|
||
|
|
if isinstance(value, dict):
|
||
|
|
value = value.get("name") or value.get("id")
|
||
|
|
return value if isinstance(value, str) and value else None
|
||
|
|
|
||
|
|
|
||
|
|
def number(value, integer=False):
|
||
|
|
if (
|
||
|
|
isinstance(value, bool)
|
||
|
|
or not isinstance(value, (int, float))
|
||
|
|
or not math.isfinite(value)
|
||
|
|
or value < 0
|
||
|
|
):
|
||
|
|
return None
|
||
|
|
return int(value) if integer and value == int(value) else None if integer else value
|
||
|
|
|
||
|
|
|
||
|
|
def normalize(raw, dataset_id):
|
||
|
|
if not isinstance(raw, dict):
|
||
|
|
raise WqError("平台目录记录格式无法识别", "invalid_response")
|
||
|
|
item_id = identifier(raw.get("id"))
|
||
|
|
owner = raw.get("dataset")
|
||
|
|
owner = owner.get("id") if isinstance(owner, dict) else owner
|
||
|
|
if dataset_id and owner != dataset_id:
|
||
|
|
raise WqError("平台返回了其他数据集的字段", "invalid_response")
|
||
|
|
coverage = number(raw.get("coverage"))
|
||
|
|
# BRAIN coverage is a fraction. Never guess that a value >1 means percent.
|
||
|
|
# Real-account schema/units still require read-only integration verification.
|
||
|
|
if coverage is not None and coverage > 1:
|
||
|
|
raise WqError("平台覆盖率单位无法确认,应为 0–1", "invalid_response")
|
||
|
|
return dict(
|
||
|
|
id=item_id,
|
||
|
|
name=label(raw.get("name")) or item_id,
|
||
|
|
category=label(raw.get("category")),
|
||
|
|
subcategory=label(raw.get("subcategory")),
|
||
|
|
field_type=label(raw.get("type")) if dataset_id else None,
|
||
|
|
coverage=coverage,
|
||
|
|
user_count=number(raw.get("userCount"), True),
|
||
|
|
alpha_count=number(raw.get("alphaCount"), True),
|
||
|
|
field_count=number(raw.get("fieldCount"), True),
|
||
|
|
description=label(raw.get("description")),
|
||
|
|
unit=label(raw.get("unit")),
|
||
|
|
)
|
||
|
|
|
||
|
|
|
||
|
|
async def sync_catalog(runner, job_id, payload):
|
||
|
|
scope = Scope.model_validate(payload["scope"])
|
||
|
|
dataset_id = payload.get("dataset_id")
|
||
|
|
async with runner.sessions() as db:
|
||
|
|
checkpoint = (await db.get(Job, job_id)).checkpoint
|
||
|
|
if checkpoint.get("done"):
|
||
|
|
return
|
||
|
|
offset = checkpoint.get("offset", 0)
|
||
|
|
while True:
|
||
|
|
await runner.checkpoint(job_id, {"next_retry_at": None})
|
||
|
|
raw = await runner.client.catalog_page(scope.model_dump(), dataset_id, offset)
|
||
|
|
rows = raw.get("results")
|
||
|
|
if not isinstance(rows, list):
|
||
|
|
raise WqError("平台目录缺少 results,已保留进度", "invalid_response")
|
||
|
|
entries = [normalize(r, dataset_id) for r in rows]
|
||
|
|
# Always probe to exhaustion if next is absent; count alone cannot prove completeness.
|
||
|
|
next_page = raw.get("next")
|
||
|
|
if "next" in raw and next_page is not None:
|
||
|
|
if not isinstance(next_page, str) or not next_page:
|
||
|
|
raise WqError("平台 next 分页格式无法识别", "invalid_response")
|
||
|
|
parsed = urlparse(next_page)
|
||
|
|
expected_path = "/data-fields" if dataset_id else "/data-sets"
|
||
|
|
offsets = parse_qs(parsed.query).get("offset", [])
|
||
|
|
if parsed.path.rstrip("/") != expected_path or offsets != [str(offset + len(rows))]:
|
||
|
|
raise WqError("平台 next 分页未按预期前进", "invalid_response")
|
||
|
|
more = next_page is not None if "next" in raw else bool(rows)
|
||
|
|
count = number(raw.get("count"), True)
|
||
|
|
if (more and not rows) or (not more and count is not None and offset + len(rows) < count):
|
||
|
|
raise WqError("平台分页提前结束,未发布不完整集合", "invalid_response")
|
||
|
|
async with runner.sessions() as db:
|
||
|
|
job = await db.get(Job, job_id)
|
||
|
|
if job.cancel_requested:
|
||
|
|
raise asyncio.CancelledError()
|
||
|
|
batch = await db.get(CatalogBatch, job_id)
|
||
|
|
added = 0
|
||
|
|
for entry in entries:
|
||
|
|
if await db.get(CatalogEntry, (job_id, entry["id"])):
|
||
|
|
continue
|
||
|
|
db.add(CatalogEntry(batch_id=job_id, **entry))
|
||
|
|
await db.flush()
|
||
|
|
added += 1
|
||
|
|
owner = dataset_id or entry["id"]
|
||
|
|
field_id = entry["id"] if dataset_id else ""
|
||
|
|
if not await db.get(CatalogNote, (scope.key(), owner, field_id)):
|
||
|
|
db.add(CatalogNote(scope_key=scope.key(), dataset_id=owner, field_id=field_id))
|
||
|
|
if rows and not added:
|
||
|
|
raise WqError("平台分页重复且未前进,已保留进度", "invalid_response")
|
||
|
|
batch.count += added
|
||
|
|
job.processed = batch.count
|
||
|
|
offset += len(rows)
|
||
|
|
job.checkpoint = dict(offset=offset, done=not more)
|
||
|
|
job.updated_at = now()
|
||
|
|
if not more:
|
||
|
|
batch.complete, batch.completed_at = True, now()
|
||
|
|
job.total = batch.count
|
||
|
|
if dataset_id:
|
||
|
|
dataset = await db.scalar(
|
||
|
|
select(CatalogDataset)
|
||
|
|
.where(CatalogDataset.scope_key == scope.key(), CatalogDataset.id == dataset_id)
|
||
|
|
.with_for_update()
|
||
|
|
)
|
||
|
|
dataset.field_version = job_id
|
||
|
|
else:
|
||
|
|
scope_row = await db.get(CatalogScope, scope.key())
|
||
|
|
scope_row.catalog_version, scope_row.synced_at = job_id, now()
|
||
|
|
ids = (
|
||
|
|
await db.scalars(select(CatalogEntry.id).where(CatalogEntry.batch_id == job_id))
|
||
|
|
).all()
|
||
|
|
for item_id in ids:
|
||
|
|
if not await db.get(CatalogDataset, (scope.key(), item_id)):
|
||
|
|
db.add(CatalogDataset(scope_key=scope.key(), id=item_id))
|
||
|
|
await db.commit()
|
||
|
|
if not more:
|
||
|
|
return
|