"""Local Pearson comparison of daily PnL changes; no platform eligibility decisions. The four-year window and signed 0.7 threshold follow the legacy research tool. Thirty paired changes is a local minimum, not a claim about BRAIN's own checks. """ import math from datetime import datetime, timezone from statistics import StatisticsError, correlation THRESHOLD = 0.7 MIN_SAMPLES = 30 WINDOW_YEARS = 4 def daily_changes(points): """Return dated changes and the final date, rejecting ambiguous daily records. Parameters are normalized PnL points. Missing values break a change interval; each change retains its start date so differently spaced samples never pair. Raises ValueError for malformed dates, duplicate days or non-finite values. """ daily = {} for point in points: try: timestamp = datetime.fromisoformat(point["date"].replace("Z", "+00:00")) day = ( (timestamp.replace(tzinfo=timezone.utc) if timestamp.tzinfo is None else timestamp) .astimezone(timezone.utc) .date() ) except (ValueError, TypeError, KeyError, AttributeError): raise ValueError("PnL 日期无法识别") from None if day in daily: raise ValueError("PnL 同一天存在多条记录") value = point.get("value") if value is not None and ( isinstance(value, bool) or not isinstance(value, (int, float)) or not math.isfinite(value) ): raise ValueError("PnL 包含无效数值") daily[day] = value changes, previous_day, previous = {}, None, None for day, value in sorted(daily.items()): if value is not None and previous is not None: delta = value - previous if not math.isfinite(delta): raise ValueError("PnL 变化超出有效数值范围") changes[day] = (previous_day, delta) previous_day, previous = day, value return changes, max(daily) if daily else None def calculate_correlation(target_points, references): """Compare a target with supplied reference caches and return a JSON-safe report. Each reference has alpha_id, points, fetched_at and optionally error. Callers select the same-region submitted set and exclude the target. Incomplete coverage can report high correlation, but can never report an all-clear. """ result = { "status": "insufficient_data", "threshold": THRESHOLD, "min_samples": MIN_SAMPLES, "window_years": WINDOW_YEARS, "window_from": None, "window_to": None, "max_correlation": None, "most_correlated_alpha_id": None, "candidate_count": len(references), "compared_count": 0, "skipped_count": 0, "matches": [], "skipped": [], "reason": None, } try: target, latest = daily_changes(target_points) except ValueError as exc: result["reason"] = str(exc) return result if latest is None or not target: result["reason"] = "目标 Alpha 没有可用的 PnL 日变化" return result try: cutoff = latest.replace(year=latest.year - WINDOW_YEARS) except ValueError: cutoff = latest.replace(year=latest.year - WINDOW_YEARS, day=28) result["window_from"], result["window_to"] = cutoff.isoformat(), latest.isoformat() matches, skipped = [], [] for reference in references: alpha_id, reason = reference["alpha_id"], reference.get("error") if not reason: try: changes, _ = daily_changes(reference["points"]) days = sorted( day for day in target.keys() & changes.keys() if cutoff < day <= latest and target[day][0] == changes[day][0] ) if len(days) < MIN_SAMPLES: reason = f"共同有效样本不足 {MIN_SAMPLES} 个(实际 {len(days)})" else: coefficient = correlation([target[d][1] for d in days], [changes[d][1] for d in days]) if not math.isfinite(coefficient): reason = "无法计算有效相关系数" else: matches.append( { "alpha_id": alpha_id, "correlation": coefficient, "sample_count": len(days), "date_from": days[0].isoformat(), "date_to": days[-1].isoformat(), "pnl_fetched_at": reference.get("fetched_at"), } ) except StatisticsError: reason = "目标或基准 PnL 日变化为常量" except (ValueError, OverflowError) as exc: reason = str(exc) if reason: skipped.append({"alpha_id": alpha_id, "reason": reason}) matches.sort(key=lambda row: (-row["correlation"], row["alpha_id"])) result.update( compared_count=len(matches), skipped_count=len(skipped), matches=matches[:10], skipped=skipped[:100] ) if matches: maximum = matches[0]["correlation"] result.update( max_correlation=maximum, most_correlated_alpha_id=matches[0]["alpha_id"], status="high" if maximum >= THRESHOLD else "partial" if skipped else "low", ) else: result["reason"] = ( "没有同地区已提交 Alpha 可供比较" if not references else "所有基准均缺少足够的有效样本" ) return result