"""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) # 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 ]}) client = WqClient(settings, transport=httpx.MockTransport(handler)) client.credentials, client.authenticated = ("test@example.com", "test"), True for submission in ("UNSUBMITTED", "SUBMITTED"): result = await client.alphas( 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", ) assert [row["id"] for row in result["results"]] == ["midnight", "last"] first, second = [dict(r.url.params) for r in requests] 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>"] 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() @pytest.mark.parametrize("region", ["USA", None]) async def test_empty_region_reference_set_passes_with_zero_but_missing_region_does_not( app, logged_in, region ): runner = app.state.runner # No credentials or PnL; this branch must stay local. async with runner.sessions.begin() as db: for raw in ( alpha("target", settings={"region": region}), alpha("pending", settings={"region": "USA"}), alpha("other-region", status="ACTIVE", settings={"region": "CHN"}), ): await upsert_alpha(db, raw) 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"] == report["skipped_count"] == 0 assert report["status"] == ("low" if region else "insufficient_data") assert report["max_correlation"] == (0.0 if region else None) assert not report["stale"] and report["matches"] == [] summary = (await logged_in.get(f"{PREFIX}/alphas?q=target")).json()["items"][0]["local_correlation"] assert summary["status"] == report["status"] assert summary["max_correlation"] == report["max_correlation"] state = (await logged_in.get(f"{PREFIX}/alphas/target/submission")).json() assert state["can_check"] is bool(region) if region: async with runner.sessions.begin() as db: await upsert_alpha(db, alpha("new-peer", status="ACTIVE")) report = (await logged_in.get(f"{PREFIX}/alphas/target/self-correlation")).json()["result"] assert report["stale"] assert not (await logged_in.get(f"{PREFIX}/alphas/target/submission")).json()["can_check"]