2026-09-08 10:40:49 +08:00
|
|
|
"""Behavioral coverage for submission scopes, daily recovery and local correlation."""
|
|
|
|
|
|
|
|
|
|
import csv
|
|
|
|
|
import io
|
|
|
|
|
from datetime import date, datetime, timedelta
|
|
|
|
|
|
|
|
|
|
import httpx
|
|
|
|
|
import pytest
|
|
|
|
|
from sqlalchemy import select
|
|
|
|
|
|
|
|
|
|
from app.alphas import upsert_alpha
|
|
|
|
|
from app.correlation import calculate_correlation, daily_changes
|
|
|
|
|
from app.jobs import Runner
|
|
|
|
|
from app.models import Alpha, JobItem, Pnl, Research, SelfCorrelation
|
|
|
|
|
from app.worldquant import WqClient, WqError
|
|
|
|
|
from tests.conftest import alpha
|
|
|
|
|
from tests.test_jobs import FakePlatform, ready_runner, result
|
|
|
|
|
|
|
|
|
|
PREFIX = "/api/v1"
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def points(changes, start=date(2025, 1, 1), multiplier=1):
|
|
|
|
|
value, data = 100, [{"date": start.isoformat(), "value": 100}]
|
|
|
|
|
for index, change in enumerate(changes, 1):
|
|
|
|
|
value += change * multiplier
|
|
|
|
|
data.append({"date": (start + timedelta(days=index)).isoformat(), "value": value})
|
|
|
|
|
return data
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def reference(alpha_id, data):
|
|
|
|
|
return {"alpha_id": alpha_id, "points": data, "fetched_at": "2025-03-01T00:00:00Z"}
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def test_signed_pearson_and_incomplete_coverage():
|
|
|
|
|
changes = [((i * 7) % 19) - 8 for i in range(60)]
|
|
|
|
|
target = points(changes)
|
|
|
|
|
positive = reference("positive", points(changes, multiplier=2))
|
|
|
|
|
negative = reference("negative", points(changes, multiplier=-1))
|
|
|
|
|
computed = calculate_correlation(target, [positive, negative])
|
|
|
|
|
assert computed["status"] == "high" and computed["compared_count"] == 2
|
|
|
|
|
assert computed["most_correlated_alpha_id"] == "positive"
|
|
|
|
|
assert computed["matches"][0]["sample_count"] == 60
|
|
|
|
|
assert computed["max_correlation"] == pytest.approx(1)
|
|
|
|
|
computed = calculate_correlation(target, [negative])
|
|
|
|
|
assert computed["status"] == "low" and computed["max_correlation"] == pytest.approx(-1)
|
|
|
|
|
incomplete = calculate_correlation(target, [negative, {"alpha_id": "missing", "error": "无权访问"}])
|
|
|
|
|
assert incomplete["status"] == "partial" and incomplete["skipped_count"] == 1
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
@pytest.mark.parametrize("candidate", [[], points([1] * 60), points([1, 3, 2] * 9)])
|
|
|
|
|
def test_insufficient_and_constant_series_are_never_passed(candidate):
|
|
|
|
|
report = calculate_correlation(points([1, 3, 2] * 20), [reference("invalid", candidate)])
|
|
|
|
|
assert report["status"] == "insufficient_data"
|
|
|
|
|
assert report["max_correlation"] is None and report["compared_count"] == 0
|
|
|
|
|
assert report["skipped_count"] == 1
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def test_missing_points_break_intervals_and_dates_align_before_comparison():
|
|
|
|
|
original = points([1, 3, 2] * 20)
|
|
|
|
|
missing = [dict(p) for p in original]
|
|
|
|
|
missing[4]["value"] = None
|
|
|
|
|
changes, _ = daily_changes(missing)
|
|
|
|
|
assert date(2025, 1, 5) not in changes and date(2025, 1, 6) not in changes
|
|
|
|
|
report = calculate_correlation(original, [reference("gap", list(reversed(missing)))])
|
|
|
|
|
assert report["matches"][0]["sample_count"] == 58
|
|
|
|
|
# A missing row must not pair a two-day increment with a one-day increment.
|
|
|
|
|
removed = original[:4] + original[5:]
|
|
|
|
|
report = calculate_correlation(original, [reference("gap", removed)])
|
|
|
|
|
assert report["matches"][0]["sample_count"] == 58
|
|
|
|
|
with pytest.raises(ValueError, match="同一天"):
|
|
|
|
|
daily_changes(original + [original[0]])
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def test_common_four_year_window_is_anchored_to_target():
|
|
|
|
|
changes = [1, 3, 2] * 20
|
|
|
|
|
target = points(changes, date(2025, 1, 1))
|
|
|
|
|
historical = points(changes, date(2019, 1, 1))
|
|
|
|
|
report = calculate_correlation(target, [reference("old", historical)])
|
|
|
|
|
assert report["status"] == "insufficient_data" and report["max_correlation"] is None
|
|
|
|
|
assert report["window_from"] == "2021-03-02"
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
async def test_submission_tabs_and_export_share_the_same_scope(app, logged_in):
|
|
|
|
|
async with app.state.sessions() as db:
|
|
|
|
|
for raw in (
|
|
|
|
|
alpha("pending"),
|
|
|
|
|
alpha("active", status="ACTIVE", stage="IS"),
|
|
|
|
|
alpha("retired", status="DECOMMISSIONED"),
|
|
|
|
|
alpha("unknown", status=None),
|
|
|
|
|
):
|
|
|
|
|
await upsert_alpha(db, raw)
|
|
|
|
|
await db.commit()
|
|
|
|
|
pending = (await logged_in.get(f"{PREFIX}/alphas?submission=UNSUBMITTED")).json()
|
|
|
|
|
submitted = (await logged_in.get(f"{PREFIX}/alphas?submission=SUBMITTED")).json()
|
|
|
|
|
assert [a["id"] for a in pending["items"]] == ["pending"]
|
|
|
|
|
assert {a["id"] for a in submitted["items"]} == {"active", "retired"}
|
|
|
|
|
exported = await logged_in.get(f"{PREFIX}/alphas/export?submission=SUBMITTED")
|
|
|
|
|
assert {r["id"] for r in csv.DictReader(io.StringIO(exported.text.lstrip("\ufeff")))} == {
|
|
|
|
|
"active",
|
|
|
|
|
"retired",
|
|
|
|
|
}
|
|
|
|
|
assert (await logged_in.get(f"{PREFIX}/alphas?submission=BAD")).status_code == 422
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
async def test_daily_job_validation_scope_deduplication_and_history(app, logged_in):
|
|
|
|
|
await ready_runner(app)
|
|
|
|
|
for payload in (
|
|
|
|
|
{"kind": "daily_sync"},
|
|
|
|
|
{"kind": "full_sync", "submission": "UNSUBMITTED"},
|
|
|
|
|
{"kind": "full_sync", "date_from": "2025-01-01"},
|
|
|
|
|
{"kind": "daily_sync", "submission": "SUBMITTED", "date_from": "2025-01-02", "date_to": "2025-01-01"},
|
|
|
|
|
{"kind": "daily_sync", "submission": "SUBMITTED", "date_from": "2025-01-01", "date_to": "9999-01-01"},
|
|
|
|
|
{"kind": "alpha_refresh", "alpha_ids": ["a"], "submission": "SUBMITTED"},
|
|
|
|
|
):
|
|
|
|
|
assert (await logged_in.post(f"{PREFIX}/sync-jobs", json=payload)).status_code == 422
|
|
|
|
|
payload = {
|
|
|
|
|
"kind": "daily_sync",
|
|
|
|
|
"submission": "UNSUBMITTED",
|
|
|
|
|
"date_from": "2025-01-01",
|
|
|
|
|
"date_to": "2025-01-02",
|
|
|
|
|
}
|
|
|
|
|
first = (await logged_in.post(f"{PREFIX}/sync-jobs", json=payload)).json()
|
|
|
|
|
duplicate = (await logged_in.post(f"{PREFIX}/sync-jobs", json=payload)).json()
|
|
|
|
|
other = (await logged_in.post(f"{PREFIX}/sync-jobs", json={**payload, "submission": "SUBMITTED"})).json()
|
|
|
|
|
assert first["id"] == duplicate["id"] != other["id"]
|
|
|
|
|
assert first["payload"]["date_from"] == "2025-01-01"
|
|
|
|
|
full = (await logged_in.post(f"{PREFIX}/sync-jobs", json={"kind": "full_sync"})).json()
|
|
|
|
|
assert full["payload"]["submission"] == "SUBMITTED"
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
class DailyPlatform(FakePlatform):
|
|
|
|
|
def __init__(self):
|
|
|
|
|
super().__init__()
|
|
|
|
|
self.daily_calls = []
|
|
|
|
|
self.fail_daily = True
|
|
|
|
|
self.records = [
|
|
|
|
|
alpha("midnight", dateCreated="2025-01-01T00:00:00Z"),
|
|
|
|
|
alpha("end", dateCreated="2025-01-01T23:59:59.999999Z"),
|
|
|
|
|
alpha("next", dateCreated="2025-01-02T00:00:00Z"),
|
|
|
|
|
alpha("hidden1", hidden=True, dateCreated="2025-01-02T01:00:00Z"),
|
|
|
|
|
alpha("hidden2", hidden=True, dateCreated="2025-01-02T02:00:00Z"),
|
|
|
|
|
alpha(
|
|
|
|
|
"submitted",
|
|
|
|
|
status="ACTIVE",
|
|
|
|
|
dateCreated="2024-01-01T00:00:00Z",
|
|
|
|
|
dateSubmitted="2025-01-02T00:00:00Z",
|
|
|
|
|
),
|
|
|
|
|
]
|
|
|
|
|
|
|
|
|
|
async def alphas(self, submission, hidden, offset, before, *, date_from=None, date_to=None):
|
|
|
|
|
self.daily_calls.append((submission, hidden, offset, date_from, date_to))
|
|
|
|
|
if self.fail_daily and hidden and offset == 1:
|
|
|
|
|
raise WqError("模拟日内第二页失败", "network_error")
|
|
|
|
|
field = "dateCreated" if submission == "UNSUBMITTED" else "dateSubmitted"
|
|
|
|
|
matched = [
|
|
|
|
|
r
|
|
|
|
|
for r in self.records
|
|
|
|
|
if (r["status"] == "UNSUBMITTED") == (submission == "UNSUBMITTED") and r["hidden"] == hidden
|
|
|
|
|
]
|
|
|
|
|
if date_from:
|
|
|
|
|
matched = [
|
|
|
|
|
r
|
|
|
|
|
for r in matched
|
|
|
|
|
if datetime.fromisoformat(date_from)
|
|
|
|
|
<= datetime.fromisoformat(r[field].replace("Z", "+00:00"))
|
|
|
|
|
< datetime.fromisoformat(date_to)
|
|
|
|
|
]
|
|
|
|
|
return {"results": matched[offset : offset + 1], "count": len(matched)}
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
async def test_daily_sync_midnight_hidden_pages_restart_and_research_preservation(app, logged_in):
|
|
|
|
|
runner = await ready_runner(app)
|
|
|
|
|
runner.client = DailyPlatform()
|
|
|
|
|
payload = {
|
|
|
|
|
"kind": "daily_sync",
|
|
|
|
|
"submission": "UNSUBMITTED",
|
|
|
|
|
"date_from": "2025-01-01",
|
|
|
|
|
"date_to": "2025-01-02",
|
|
|
|
|
}
|
|
|
|
|
job_id = (await logged_in.post(f"{PREFIX}/sync-jobs", json=payload)).json()["id"]
|
|
|
|
|
await runner.execute(job_id)
|
|
|
|
|
failed = await result(runner, job_id)
|
|
|
|
|
assert failed.status == "failed" and failed.processed == 4
|
|
|
|
|
assert failed.checkpoint["date"] == "2025-01-02" and failed.checkpoint["offset"] == 1
|
|
|
|
|
async with runner.sessions() as db:
|
|
|
|
|
research = await db.get(Research, "midnight")
|
|
|
|
|
research.note, research.state = "preserved", "candidate"
|
|
|
|
|
await db.commit()
|
|
|
|
|
resumed = Runner(runner.sessions, runner.settings, DailyPlatform())
|
|
|
|
|
resumed.client.fail_daily = False
|
|
|
|
|
await resumed.execute(job_id)
|
|
|
|
|
complete = await result(resumed, job_id)
|
|
|
|
|
assert complete.status == "completed" and complete.processed == complete.total == 5
|
|
|
|
|
assert complete.checkpoint["dates_completed"] == complete.checkpoint["dates_total"] == 2
|
|
|
|
|
assert resumed.client.daily_calls[0][1:3] == (True, 1)
|
|
|
|
|
async with runner.sessions() as db:
|
|
|
|
|
assert len((await db.scalars(select(JobItem).where(JobItem.job_id == job_id))).all()) == 5
|
|
|
|
|
assert (await db.get(Research, "midnight")).note == "preserved"
|
|
|
|
|
assert await db.get(Alpha, "submitted") is None
|
|
|
|
|
# Daily submitted sync uses submission date, although its creation was a year earlier.
|
|
|
|
|
submitted_id = (
|
|
|
|
|
await logged_in.post(f"{PREFIX}/sync-jobs", json={**payload, "submission": "SUBMITTED"})
|
|
|
|
|
).json()["id"]
|
|
|
|
|
await resumed.execute(submitted_id)
|
|
|
|
|
assert (await result(resumed, submitted_id)).processed == 1
|
|
|
|
|
full_id = (await logged_in.post(f"{PREFIX}/sync-jobs", json={"kind": "full_sync"})).json()["id"]
|
|
|
|
|
resumed.client.daily_calls.clear()
|
|
|
|
|
await resumed.execute(full_id)
|
|
|
|
|
assert (await result(resumed, full_id)).processed == 1
|
|
|
|
|
assert {call[0] for call in resumed.client.daily_calls} == {"SUBMITTED"}
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
async def test_platform_ignoring_daily_filter_does_not_import_other_days(app, logged_in):
|
|
|
|
|
runner = await ready_runner(app)
|
|
|
|
|
|
|
|
|
|
async def wrong_day(*args, **kwargs):
|
|
|
|
|
return {"results": [alpha("wrong", dateCreated="2024-01-01T00:00:00Z")], "next": None}
|
|
|
|
|
|
|
|
|
|
runner.client.alphas = wrong_day
|
|
|
|
|
job = (
|
|
|
|
|
await logged_in.post(
|
|
|
|
|
f"{PREFIX}/sync-jobs",
|
|
|
|
|
json={
|
|
|
|
|
"kind": "daily_sync",
|
|
|
|
|
"submission": "UNSUBMITTED",
|
|
|
|
|
"date_from": "2025-01-01",
|
|
|
|
|
"date_to": "2025-01-01",
|
|
|
|
|
},
|
|
|
|
|
)
|
|
|
|
|
).json()
|
|
|
|
|
await runner.execute(job["id"])
|
|
|
|
|
assert (await result(runner, job["id"])).status == "failed"
|
|
|
|
|
async with runner.sessions() as db:
|
|
|
|
|
assert await db.get(Alpha, "wrong") is None
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
async def test_local_detection_uses_cache_excludes_self_and_keeps_research(app, logged_in):
|
|
|
|
|
runner = app.state.runner # No credentials: any upstream call fails this test.
|
|
|
|
|
data = points([1, 4, 2, -2] * 20)
|
|
|
|
|
async with runner.sessions() as db:
|
|
|
|
|
for raw in [
|
|
|
|
|
alpha("target", status="ACTIVE"),
|
|
|
|
|
alpha("peer", status="ACTIVE"),
|
|
|
|
|
alpha("pending"),
|
|
|
|
|
alpha("other-region", status="ACTIVE", settings={"region": "CHN"}),
|
|
|
|
|
alpha("unknown", status=None),
|
|
|
|
|
]:
|
|
|
|
|
await upsert_alpha(db, raw)
|
|
|
|
|
db.add(Pnl(alpha_id=raw["id"], raw={}, points=data))
|
|
|
|
|
await db.flush()
|
|
|
|
|
research = await db.get(Research, "target")
|
|
|
|
|
research.note, research.state = "hypothesis", "candidate"
|
|
|
|
|
await db.commit()
|
|
|
|
|
job = (
|
|
|
|
|
await logged_in.post(
|
|
|
|
|
f"{PREFIX}/sync-jobs", json={"kind": "self_correlation", "alpha_ids": ["target"]}
|
|
|
|
|
)
|
|
|
|
|
).json()
|
|
|
|
|
await runner.execute(job["id"])
|
|
|
|
|
assert (await result(runner, job["id"])).status == "completed"
|
|
|
|
|
report = (await logged_in.get(f"{PREFIX}/alphas/target/self-correlation")).json()["result"]
|
|
|
|
|
assert report["candidate_count"] == report["compared_count"] == 1
|
|
|
|
|
assert report["matches"][0]["alpha_id"] == "peer" and report["status"] == "high"
|
|
|
|
|
for timestamp in (
|
|
|
|
|
report["calculated_at"],
|
|
|
|
|
report["target_pnl_fetched_at"],
|
|
|
|
|
report["matches"][0]["pnl_fetched_at"],
|
|
|
|
|
):
|
|
|
|
|
assert datetime.fromisoformat(timestamp).utcoffset() == timedelta(0)
|
|
|
|
|
async with runner.sessions() as db:
|
|
|
|
|
assert (await db.get(Research, "target")).state == "candidate"
|
|
|
|
|
assert (await db.get(Research, "target")).note == "hypothesis"
|
|
|
|
|
assert (await db.get(Alpha, "target")).checks[0]["result"] == "FAIL"
|
|
|
|
|
# A new submitted reference invalidates the old result without deleting it.
|
|
|
|
|
await upsert_alpha(db, alpha("new-peer", status="ACTIVE"))
|
|
|
|
|
await db.commit()
|
|
|
|
|
assert (await db.get(SelfCorrelation, "target")).stale
|
|
|
|
|
summary = (await logged_in.get(f"{PREFIX}/alphas?q=target")).json()["items"][0]["local_correlation"]
|
|
|
|
|
assert summary["stale"] and summary["max_correlation"] == pytest.approx(1)
|
|
|
|
|
assert datetime.fromisoformat(summary["calculated_at"]).utcoffset() == timedelta(0)
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
async def test_missing_pnl_filled_once_and_partial_data_reported(app, logged_in):
|
|
|
|
|
runner = await ready_runner(app)
|
|
|
|
|
calls = []
|
|
|
|
|
data = points([1, 3, -2] * 20)
|
|
|
|
|
|
|
|
|
|
async def pnl(alpha_id):
|
|
|
|
|
calls.append(alpha_id)
|
|
|
|
|
if alpha_id == "unavailable":
|
|
|
|
|
raise WqError("无权访问", "access_denied")
|
|
|
|
|
multiplier = 1 if alpha_id == "target" else -1
|
|
|
|
|
return {"records": [{"date": p["date"], "pnl": p["value"] * multiplier} for p in data]}
|
|
|
|
|
|
|
|
|
|
runner.client.pnl = pnl
|
|
|
|
|
async with runner.sessions() as db:
|
|
|
|
|
for raw in [alpha("target"), alpha("peer", status="ACTIVE"), alpha("unavailable", status="ACTIVE")]:
|
|
|
|
|
await upsert_alpha(db, raw)
|
|
|
|
|
await db.commit()
|
|
|
|
|
|
|
|
|
|
async def run():
|
|
|
|
|
job = (
|
|
|
|
|
await logged_in.post(
|
|
|
|
|
f"{PREFIX}/sync-jobs", json={"kind": "self_correlation", "alpha_ids": ["target"]}
|
|
|
|
|
)
|
|
|
|
|
).json()
|
|
|
|
|
await runner.execute(job["id"])
|
|
|
|
|
return (await logged_in.get(f"{PREFIX}/alphas/target/self-correlation")).json()["result"]
|
|
|
|
|
|
|
|
|
|
report = await run()
|
|
|
|
|
assert report["status"] == "partial" and report["skipped_count"] == 1
|
|
|
|
|
await run()
|
|
|
|
|
assert calls.count("target") == calls.count("peer") == 1
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
async def test_scoped_date_query_parameters_and_no_platform_check(settings):
|
|
|
|
|
requests = []
|
|
|
|
|
|
|
|
|
|
def handler(request):
|
|
|
|
|
requests.append(request)
|
2026-09-09 15:12:45 +08:00
|
|
|
# BRAIN parses the extra '=' as part of the datetime, returning HTTP 400.
|
|
|
|
|
if any(key.endswith(">=") for key in request.url.params):
|
|
|
|
|
return httpx.Response(400, json=["Expected ISO 8601 datetime with timezone"])
|
|
|
|
|
field = "dateCreated" if "status" in request.url.params else "dateSubmitted"
|
|
|
|
|
lower = datetime.fromisoformat(request.url.params[field + ">"])
|
|
|
|
|
upper = datetime.fromisoformat(request.url.params[field + "<"])
|
|
|
|
|
records = [
|
|
|
|
|
{"id": "midnight", field: "2025-01-01T00:00:00+00:00"},
|
|
|
|
|
{"id": "last", field: "2025-01-01T23:59:59.999999+00:00"},
|
|
|
|
|
{"id": "next", field: "2025-01-02T00:00:00+00:00"},
|
|
|
|
|
]
|
|
|
|
|
return httpx.Response(200, json={"results": [
|
|
|
|
|
row for row in records if lower <= datetime.fromisoformat(row[field]) <= upper
|
|
|
|
|
]})
|
2026-09-08 10:40:49 +08:00
|
|
|
|
|
|
|
|
client = WqClient(settings, transport=httpx.MockTransport(handler))
|
|
|
|
|
client.credentials, client.authenticated = ("test@example.com", "test"), True
|
|
|
|
|
for submission in ("UNSUBMITTED", "SUBMITTED"):
|
2026-09-09 15:12:45 +08:00
|
|
|
result = await client.alphas(
|
2026-09-08 10:40:49 +08:00
|
|
|
submission,
|
|
|
|
|
True,
|
|
|
|
|
100,
|
|
|
|
|
"2025-03-01T00:00:00+00:00",
|
|
|
|
|
date_from="2025-01-01T00:00:00+00:00",
|
|
|
|
|
date_to="2025-01-02T00:00:00+00:00",
|
|
|
|
|
)
|
2026-09-09 15:12:45 +08:00
|
|
|
assert [row["id"] for row in result["results"]] == ["midnight", "last"]
|
2026-09-08 10:40:49 +08:00
|
|
|
first, second = [dict(r.url.params) for r in requests]
|
2026-09-09 15:12:45 +08:00
|
|
|
assert first["dateCreated>"] == "2025-01-01T00:00:00+00:00"
|
|
|
|
|
assert first["dateCreated<"] == "2025-01-01T23:59:59.999999+00:00"
|
|
|
|
|
assert second["dateSubmitted>"] == first["dateCreated>"]
|
2026-09-08 10:40:49 +08:00
|
|
|
assert second["dateSubmitted<"] == first["dateCreated<"]
|
|
|
|
|
assert "status!" in second and second["hidden"] == "true"
|
|
|
|
|
assert all(r.method == "GET" and r.url.path == "/users/self/alphas" for r in requests)
|
|
|
|
|
await client.close()
|