2026-09-07 14:54:20 +08:00
|
|
|
"""Normalize upstream business data, without making platform or research decisions."""
|
|
|
|
|
|
|
|
|
|
import math
|
|
|
|
|
import re
|
|
|
|
|
from datetime import datetime
|
|
|
|
|
|
2026-09-08 10:40:49 +08:00
|
|
|
from sqlalchemy import or_, select, update
|
|
|
|
|
|
|
|
|
|
from .models import Alpha, Research, ResearchTag, SelfCorrelation, now
|
2026-09-12 22:34:14 +08:00
|
|
|
from .platform_checks import check_result, split_checks, submission_limits
|
2026-09-08 12:43:00 +08:00
|
|
|
from .research.provenance import source_alpha_ids
|
2026-09-08 10:40:49 +08:00
|
|
|
|
2026-09-09 19:02:37 +08:00
|
|
|
METRIC_FIELDS = (
|
|
|
|
|
"sharpe", "fitness", "returns", "turnover", "margin", "drawdown",
|
|
|
|
|
"sub_universe_sharpe", "robust_universe_sharpe", "two_year_sharpe", "prod_correlation", "pnl",
|
|
|
|
|
)
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def failed_checks(checks):
|
2026-09-12 22:11:11 +08:00
|
|
|
"""Return failed Alpha check names, excluding submission limits and local correlation."""
|
2026-09-09 19:02:37 +08:00
|
|
|
return [
|
|
|
|
|
check.get("name") if isinstance(check.get("name"), str) else "未命名检查"
|
2026-09-12 22:34:14 +08:00
|
|
|
for check in split_checks(checks)[0] if isinstance(check, dict) and check_result(check) == "FAIL"
|
2026-09-09 19:02:37 +08:00
|
|
|
] if isinstance(checks, list) else []
|
|
|
|
|
|
|
|
|
|
|
2026-09-12 22:34:14 +08:00
|
|
|
def snapshot_columns(settings, metrics, checks, *, checked=False):
|
2026-09-09 19:02:37 +08:00
|
|
|
"""Derive list fields from a platform snapshot, preserving missing metrics as null.
|
|
|
|
|
|
2026-09-12 22:11:11 +08:00
|
|
|
Submission limits are excluded. Only explicit Alpha FAIL results count.
|
2026-09-12 22:34:14 +08:00
|
|
|
Sync snapshots with no failures are PRE_CHECK; a completed explicit /check
|
|
|
|
|
with no failures is PASS. WARNING/PENDING do not count as failures, matching
|
|
|
|
|
the legacy workflow. Empty, malformed or unknown results remain PENDING.
|
2026-09-09 19:02:37 +08:00
|
|
|
No submission eligibility or activity eligibility is inferred here.
|
|
|
|
|
"""
|
|
|
|
|
settings = settings if isinstance(settings, dict) else {}
|
|
|
|
|
metrics = metrics if isinstance(metrics, dict) else {}
|
2026-09-12 22:11:11 +08:00
|
|
|
checks, _ = split_checks(checks)
|
2026-09-09 19:02:37 +08:00
|
|
|
valid = [check for check in checks if isinstance(check, dict)]
|
|
|
|
|
failures = len(failed_checks(checks))
|
|
|
|
|
by_name = {check["name"]: check for check in valid if isinstance(check.get("name"), str)}
|
|
|
|
|
if failures:
|
|
|
|
|
check_type = "FAIL_1" if failures == 1 else "FAIL_2"
|
2026-09-12 22:34:14 +08:00
|
|
|
elif not checks or len(valid) != len(checks) or any(check_result(check) not in ("PASS", "WARNING", "PENDING") for check in valid):
|
2026-09-09 19:02:37 +08:00
|
|
|
check_type = "PENDING"
|
|
|
|
|
else:
|
2026-09-12 22:34:14 +08:00
|
|
|
check_type = "PASS" if checked else "PRE_CHECK"
|
2026-09-11 13:14:27 +08:00
|
|
|
# /check values are freshest; submitted snapshots also expose a scalar in IS.
|
|
|
|
|
prod_correlation = number(by_name.get("PROD_CORRELATION", {}).get("value"))
|
|
|
|
|
if prod_correlation is None:
|
|
|
|
|
prod_correlation = number(metrics.get("prodCorrelation"))
|
2026-09-09 19:02:37 +08:00
|
|
|
neutralization = settings.get("neutralization")
|
|
|
|
|
return {
|
|
|
|
|
"check_type": check_type,
|
|
|
|
|
"neutralization": neutralization if isinstance(neutralization, str) else None,
|
|
|
|
|
"pnl": number(metrics.get("pnl")),
|
2026-09-11 13:14:27 +08:00
|
|
|
"prod_correlation": prod_correlation,
|
2026-09-09 19:02:37 +08:00
|
|
|
**{
|
|
|
|
|
field: number(by_name.get(name, {}).get("value"))
|
|
|
|
|
for field, name in (
|
|
|
|
|
("sub_universe_sharpe", "LOW_SUB_UNIVERSE_SHARPE"),
|
|
|
|
|
("robust_universe_sharpe", "LOW_ROBUST_UNIVERSE_SHARPE"),
|
|
|
|
|
("two_year_sharpe", "LOW_2Y_SHARPE"),
|
2026-09-11 13:14:27 +08:00
|
|
|
|
2026-09-09 19:02:37 +08:00
|
|
|
)
|
|
|
|
|
},
|
|
|
|
|
}
|
|
|
|
|
|
2026-09-08 10:40:49 +08:00
|
|
|
|
2026-09-12 22:34:14 +08:00
|
|
|
def check_summary(checks, *, check_type):
|
2026-09-12 22:11:11 +08:00
|
|
|
"""Separate cached Alpha findings from submission limits; infer no live eligibility."""
|
|
|
|
|
return {
|
2026-09-12 22:34:14 +08:00
|
|
|
"check_type": check_type,
|
2026-09-12 22:11:11 +08:00
|
|
|
"failed_checks": failed_checks(checks),
|
|
|
|
|
"submission_limits": submission_limits(checks),
|
2026-09-12 22:34:14 +08:00
|
|
|
"meaning": "PRE_CHECK 为同步无失败项;PASS 为主动检查完成且无失败项。PENDING/WARNING 不算失败,不代表全部检查项 PASS 或当前可提交",
|
2026-09-12 22:11:11 +08:00
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
2026-09-08 10:40:49 +08:00
|
|
|
def submission_condition(submission):
|
|
|
|
|
"""Match the platform list contract; a missing status is never assumed submitted."""
|
|
|
|
|
return Alpha.status == "UNSUBMITTED" if submission == "UNSUBMITTED" else Alpha.status != "UNSUBMITTED"
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
async def invalidate_correlations(db, alpha_id, regions=()):
|
|
|
|
|
"""A changed baseline or PnL invalidates local conclusions without touching platform checks."""
|
|
|
|
|
await db.execute(
|
|
|
|
|
update(SelfCorrelation)
|
|
|
|
|
.where(or_(SelfCorrelation.alpha_id == alpha_id, SelfCorrelation.region.in_(regions)))
|
|
|
|
|
.values(stale=True)
|
|
|
|
|
)
|
2026-09-07 14:54:20 +08:00
|
|
|
|
|
|
|
|
|
|
|
|
|
SENSITIVE_KEYS = {
|
|
|
|
|
"password",
|
|
|
|
|
"token",
|
|
|
|
|
"cookies",
|
|
|
|
|
"cookie",
|
|
|
|
|
"authorization",
|
|
|
|
|
"credentials",
|
|
|
|
|
"secret",
|
|
|
|
|
"accesstoken",
|
|
|
|
|
"refreshtoken",
|
|
|
|
|
"sessiontoken",
|
|
|
|
|
"clientsecret",
|
|
|
|
|
"authorizationheader",
|
|
|
|
|
"setcookie",
|
|
|
|
|
"apikey",
|
|
|
|
|
"csrftoken",
|
|
|
|
|
"xsrftoken",
|
|
|
|
|
"authentication",
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def sanitize(value):
|
|
|
|
|
if isinstance(value, dict):
|
|
|
|
|
return {
|
|
|
|
|
k: sanitize(v) for k, v in value.items() if re.sub(r"[^a-z]", "", k.lower()) not in SENSITIVE_KEYS
|
|
|
|
|
}
|
|
|
|
|
if isinstance(value, list):
|
|
|
|
|
return [sanitize(v) for v in value]
|
|
|
|
|
if isinstance(value, float) and not math.isfinite(value):
|
|
|
|
|
return None
|
|
|
|
|
return value
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def number(value):
|
|
|
|
|
if value is None or isinstance(value, bool):
|
|
|
|
|
return None
|
|
|
|
|
try:
|
|
|
|
|
result = float(value)
|
|
|
|
|
return result if math.isfinite(result) else None
|
|
|
|
|
except (ValueError, TypeError):
|
|
|
|
|
return None
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def date(value):
|
|
|
|
|
try:
|
|
|
|
|
return datetime.fromisoformat(value.replace("Z", "+00:00")) if value else None
|
|
|
|
|
except (ValueError, TypeError, AttributeError):
|
|
|
|
|
return None
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def code(value):
|
|
|
|
|
return value.get("code") if isinstance(value, dict) else value if isinstance(value, str) else None
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
async def upsert_alpha(db, raw: dict):
|
|
|
|
|
alpha_id = raw.get("id")
|
|
|
|
|
if not isinstance(alpha_id, str) or not alpha_id:
|
|
|
|
|
raise ValueError("Alpha 数据缺少 ID")
|
|
|
|
|
item = await db.get(Alpha, alpha_id)
|
2026-09-08 10:40:49 +08:00
|
|
|
previous_region = item.region if item else None
|
|
|
|
|
previous_status = item.status if item else None
|
2026-09-07 14:54:20 +08:00
|
|
|
if item is None:
|
|
|
|
|
item = Alpha(id=alpha_id)
|
|
|
|
|
db.add(item)
|
|
|
|
|
settings = raw.get("settings") or {}
|
|
|
|
|
metrics = raw.get("is") if isinstance(raw.get("is"), dict) else {}
|
|
|
|
|
item.name = raw.get("name")
|
|
|
|
|
item.expression = code(raw.get("regular"))
|
|
|
|
|
item.selection, item.combo = code(raw.get("selection")), code(raw.get("combo"))
|
|
|
|
|
item.alpha_type, item.language = raw.get("type"), settings.get("language")
|
|
|
|
|
item.stage, item.status, item.hidden = raw.get("stage"), raw.get("status"), raw.get("hidden") is True
|
|
|
|
|
item.region, item.universe = settings.get("region"), settings.get("universe")
|
2026-09-08 10:40:49 +08:00
|
|
|
if previous_region != item.region or previous_status != item.status:
|
|
|
|
|
regions = {
|
|
|
|
|
region
|
|
|
|
|
for region, status in ((previous_region, previous_status), (item.region, item.status))
|
|
|
|
|
if region and status and status != "UNSUBMITTED"
|
|
|
|
|
}
|
|
|
|
|
await invalidate_correlations(db, alpha_id, regions)
|
2026-09-07 14:54:20 +08:00
|
|
|
item.settings, item.is_metrics = sanitize(settings), sanitize(metrics)
|
|
|
|
|
item.os_metrics = sanitize(raw.get("os")) if isinstance(raw.get("os"), dict) else {}
|
|
|
|
|
item.checks = sanitize(metrics.get("checks") or raw.get("checks") or [])
|
2026-09-09 19:02:37 +08:00
|
|
|
for key, value in snapshot_columns(item.settings, item.is_metrics, item.checks).items():
|
|
|
|
|
setattr(item, key, value)
|
2026-09-07 14:54:20 +08:00
|
|
|
for key in ("sharpe", "fitness", "returns", "turnover", "margin", "drawdown"):
|
|
|
|
|
setattr(item, key, number(metrics.get(key)))
|
|
|
|
|
item.date_created, item.date_submitted = date(raw.get("dateCreated")), date(raw.get("dateSubmitted"))
|
|
|
|
|
item.synced_at, item.raw = now(), sanitize(raw)
|
|
|
|
|
await db.flush()
|
|
|
|
|
if await db.get(Research, alpha_id) is None:
|
|
|
|
|
db.add(Research(alpha_id=alpha_id))
|
|
|
|
|
return item
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def list_statement(filters):
|
|
|
|
|
query = select(Alpha, Research).join(Research, Research.alpha_id == Alpha.id)
|
2026-09-08 10:40:49 +08:00
|
|
|
if filters.submission:
|
|
|
|
|
query = query.where(submission_condition(filters.submission))
|
2026-09-08 12:43:00 +08:00
|
|
|
source_filters = {k: getattr(filters, k) for k in ("source", "source_reference", "research_id", "backtest_run_id")}
|
|
|
|
|
if any(source_filters.values()):
|
|
|
|
|
query = query.where(Alpha.id.in_(source_alpha_ids(**source_filters)))
|
2026-09-07 14:54:20 +08:00
|
|
|
q = filters.q
|
|
|
|
|
if q:
|
|
|
|
|
pattern = "%" + q.replace("\\", "\\\\").replace("%", "\\%").replace("_", "\\_") + "%"
|
|
|
|
|
query = query.where(
|
|
|
|
|
or_(
|
|
|
|
|
*(
|
|
|
|
|
getattr(Alpha, f).ilike(pattern, escape="\\")
|
|
|
|
|
for f in ("id", "name", "expression", "selection", "combo")
|
|
|
|
|
)
|
|
|
|
|
)
|
|
|
|
|
)
|
2026-09-09 19:02:37 +08:00
|
|
|
for name in ("region", "universe", "alpha_type", "language", "status", "stage", "hidden", "check_type", "neutralization"):
|
2026-09-07 14:54:20 +08:00
|
|
|
value = getattr(filters, name)
|
|
|
|
|
if value is not None:
|
|
|
|
|
query = query.where(getattr(Alpha, name) == value)
|
|
|
|
|
if filters.research_state:
|
|
|
|
|
query = query.where(Research.state == filters.research_state)
|
|
|
|
|
if filters.favorite is not None:
|
|
|
|
|
query = query.where(Research.favorite == filters.favorite)
|
|
|
|
|
if filters.tag:
|
|
|
|
|
query = query.where(Alpha.id.in_(select(ResearchTag.alpha_id).where(ResearchTag.tag == filters.tag)))
|
|
|
|
|
if filters.created_from:
|
|
|
|
|
query = query.where(Alpha.date_created >= filters.created_from)
|
|
|
|
|
if filters.created_to:
|
|
|
|
|
query = query.where(Alpha.date_created <= filters.created_to)
|
2026-09-09 19:02:37 +08:00
|
|
|
for name in METRIC_FIELDS:
|
2026-09-07 14:54:20 +08:00
|
|
|
for suffix, compare in (("min", "ge"), ("max", "le")):
|
|
|
|
|
value = getattr(filters, f"{name}_{suffix}")
|
|
|
|
|
if value is not None:
|
|
|
|
|
column = getattr(Alpha, name)
|
|
|
|
|
query = query.where(column >= value if compare == "ge" else column <= value)
|
|
|
|
|
return query
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def sorted_statement(query, sort, direction):
|
|
|
|
|
column = getattr(Alpha, sort)
|
|
|
|
|
order = column.desc() if direction == "desc" else column.asc()
|
|
|
|
|
return query.order_by(order.nullslast(), Alpha.id.asc())
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def summary(item: Alpha, research: Research):
|
|
|
|
|
keys = (
|
|
|
|
|
"id",
|
|
|
|
|
"name",
|
|
|
|
|
"alpha_type",
|
|
|
|
|
"language",
|
|
|
|
|
"stage",
|
|
|
|
|
"status",
|
|
|
|
|
"hidden",
|
|
|
|
|
"region",
|
|
|
|
|
"universe",
|
|
|
|
|
"sharpe",
|
|
|
|
|
"fitness",
|
|
|
|
|
"returns",
|
|
|
|
|
"turnover",
|
|
|
|
|
"margin",
|
|
|
|
|
"drawdown",
|
|
|
|
|
"date_created",
|
|
|
|
|
"date_submitted",
|
|
|
|
|
"synced_at",
|
2026-09-09 19:02:37 +08:00
|
|
|
"check_type",
|
|
|
|
|
"neutralization",
|
|
|
|
|
"sub_universe_sharpe",
|
|
|
|
|
"robust_universe_sharpe",
|
|
|
|
|
"two_year_sharpe",
|
|
|
|
|
"prod_correlation",
|
|
|
|
|
"pnl",
|
2026-09-07 14:54:20 +08:00
|
|
|
)
|
|
|
|
|
result = {k: getattr(item, k) for k in keys}
|
2026-09-09 19:02:37 +08:00
|
|
|
result["failed_checks"] = failed_checks(item.checks)
|
2026-09-07 14:54:20 +08:00
|
|
|
result["expression_preview"] = (item.expression or item.selection or "")[:240]
|
|
|
|
|
result["research"] = {
|
2026-09-07 23:02:55 +08:00
|
|
|
k: getattr(research, k) for k in ("note", "tags", "favorite", "state", "updated_at", "version")
|
2026-09-07 14:54:20 +08:00
|
|
|
}
|
|
|
|
|
return result
|
|
|
|
|
|
|
|
|
|
|
2026-09-10 14:09:41 +08:00
|
|
|
def pnl_points(raw, column=None):
|
2026-09-07 14:54:20 +08:00
|
|
|
"""Use schema column names, preserving missing values rather than creating zero PnL."""
|
|
|
|
|
records = raw.get("records")
|
|
|
|
|
schema = raw.get("schema") or {}
|
|
|
|
|
properties = schema.get("properties", []) if isinstance(schema, dict) else schema
|
|
|
|
|
if isinstance(properties, dict):
|
|
|
|
|
names = list(properties)
|
|
|
|
|
else:
|
|
|
|
|
names = [p.get("name", "") if isinstance(p, dict) else str(p) for p in properties]
|
|
|
|
|
normalized = [name.lower() for name in names]
|
2026-09-10 14:09:41 +08:00
|
|
|
value_names = (column,) if column else ("pnl", "value")
|
2026-09-07 14:54:20 +08:00
|
|
|
if not isinstance(records, list):
|
|
|
|
|
raise ValueError("PnL 缺少 records")
|
|
|
|
|
points = []
|
|
|
|
|
for row in records:
|
|
|
|
|
if isinstance(row, dict):
|
|
|
|
|
row = {str(k).lower(): v for k, v in row.items()}
|
|
|
|
|
timestamp = next((row[k] for k in ("date", "datetime", "timestamp") if k in row), None)
|
2026-09-10 14:09:41 +08:00
|
|
|
value = next((row[k] for k in value_names if k in row), None)
|
2026-09-07 14:54:20 +08:00
|
|
|
else:
|
|
|
|
|
date_i = next(
|
|
|
|
|
(i for i, n in enumerate(normalized) if n in ("date", "datetime", "timestamp")), None
|
|
|
|
|
)
|
2026-09-10 14:09:41 +08:00
|
|
|
pnl_i = next((i for i, n in enumerate(normalized) if n in value_names), None)
|
|
|
|
|
if date_i is None or pnl_i is None or not isinstance(row, list) or len(row) <= date_i:
|
2026-09-07 14:54:20 +08:00
|
|
|
raise ValueError("PnL schema 无法识别日期或数值列")
|
2026-09-10 14:09:41 +08:00
|
|
|
if len(row) <= pnl_i and column is None:
|
|
|
|
|
raise ValueError("PnL schema 无法识别日期或数值列")
|
|
|
|
|
timestamp, value = row[date_i], row[pnl_i] if len(row) > pnl_i else None
|
2026-09-07 14:54:20 +08:00
|
|
|
if timestamp is not None:
|
|
|
|
|
if isinstance(timestamp, (int, float)):
|
|
|
|
|
from datetime import timezone
|
|
|
|
|
|
|
|
|
|
timestamp = datetime.fromtimestamp(
|
|
|
|
|
timestamp / 1000 if timestamp > 1e11 else timestamp, tz=timezone.utc
|
|
|
|
|
).isoformat()
|
|
|
|
|
points.append({"date": str(timestamp), "value": number(value)})
|
|
|
|
|
return sorted(points, key=lambda p: p["date"])
|
2026-09-10 14:09:41 +08:00
|
|
|
|
|
|
|
|
|
|
|
|
|
def glb_pnl_series(raw, points):
|
|
|
|
|
"""Read GLB display series from cached raw data; keep the correlation baseline intact.
|
|
|
|
|
|
|
|
|
|
Missing columns are omitted, while missing values remain gaps. Legacy caches
|
|
|
|
|
containing only normalized points still return their overall PnL.
|
|
|
|
|
"""
|
|
|
|
|
series = [{"id": "pnl", "label": "总体 PnL", "points": points}]
|
|
|
|
|
schema = raw.get("schema") or {}
|
|
|
|
|
properties = schema.get("properties", []) if isinstance(schema, dict) else schema
|
|
|
|
|
names = properties if isinstance(properties, dict) else [
|
|
|
|
|
p.get("name", "") if isinstance(p, dict) else str(p) for p in properties
|
|
|
|
|
]
|
|
|
|
|
available = {name.lower() for name in names}
|
|
|
|
|
for row in raw.get("records", []):
|
|
|
|
|
if isinstance(row, dict):
|
|
|
|
|
available.update(str(key).lower() for key in row)
|
|
|
|
|
for column, label in (
|
|
|
|
|
("investability-constrained-pnl", "可投资性约束 PnL"),
|
|
|
|
|
("amer-pnl", "AMER PnL"),
|
|
|
|
|
("apac-pnl", "APAC PnL"),
|
|
|
|
|
("emea-pnl", "EMEA PnL"),
|
|
|
|
|
):
|
|
|
|
|
if column in available:
|
|
|
|
|
series.append({"id": column, "label": label, "points": pnl_points(raw, column)})
|
|
|
|
|
return series
|