from __future__ import annotations

import json
import csv
from dataclasses import replace
from datetime import UTC, date, datetime
from pathlib import Path

import pandas as pd
import pytest

from psx_signal.config import Settings
from psx_signal.data.providers.yahoo_provider import YahooDownload
from psx_signal.data.storage import DatasetStore
from psx_signal.pipelines.daily_operations import (
    OFFICIAL_PROVIDER,
    BenchmarkAppendService,
    BenchmarkRefreshResult,
    ComparisonThresholds,
    DailyEquityRefreshService,
    DailyProviderReconciler,
    EquityRefreshResult,
    OfficialIngestResult,
    OfficialDailyIngestService,
    ProspectiveDailyService,
    SessionTiming,
    membership_status,
)


def _settings(tmp_path: Path) -> Settings:
    root = tmp_path / "data"
    return replace(
        Settings(),
        data=replace(
            Settings().data, root=str(root), inbox=str(root / "inbox/psx"),
            yahoo_rate_limit_seconds=0, yahoo_request_retries=1,
        ),
    )


def _equity(symbol: str, values: list[tuple[str, float]]) -> pd.DataFrame:
    return pd.DataFrame([{
        "symbol": symbol, "date": pd.Timestamp(day), "open": close, "high": close + 1,
        "low": close - 1, "close": close, "volume": 1000, "adjusted_close": close,
        "provider": "yahoo_finance", "provider_trust": "RESEARCH_SECONDARY",
        "source_reference": f"Yahoo Finance daily history: {symbol}.KA",
    } for day, close in values])


class FakeYahoo:
    def __init__(self, frame: pd.DataFrame) -> None:
        self.frame = frame

    @staticmethod
    def resolve_symbol(symbol: str) -> str:
        return f"{symbol}.KA"

    def download_symbol(self, symbol: str, start: date, end: date) -> YahooDownload:
        available = self.frame[self.frame["symbol"] == symbol].copy().set_index("date")
        if available.empty:
            return YahooDownload(symbol, f"{symbol}.KA", "NOT_FOUND", pd.DataFrame(), 1, "missing")
        raw = available.rename(columns={
            "open": "Open", "high": "High", "low": "Low", "close": "Close",
            "volume": "Volume", "adjusted_close": "Adj Close",
        })[["Open", "High", "Low", "Close", "Volume", "Adj Close"]]
        return YahooDownload(symbol, f"{symbol}.KA", "SUCCESS", raw, 1)

    def normalize_download(self, download: YahooDownload) -> pd.DataFrame:
        raw = download.frame.reset_index()
        return _equity(download.psx_symbol, [
            (str(pd.Timestamp(row.date).date()), float(row.Close)) for row in raw.itertuples()
        ])

    @staticmethod
    def get_actions(symbol: str, start: date, end: date) -> pd.DataFrame:
        return pd.DataFrame()


def test_incremental_refresh_rejects_current_session_before_finality(tmp_path: Path) -> None:
    settings = _settings(tmp_path); store = DatasetStore(settings.data.root)
    DatasetStore._atomic_parquet(_equity("AAA", [("2026-08-17", 100)]), store.equities_path)
    provider = FakeYahoo(_equity("AAA", [("2026-08-18", 101), ("2026-08-19", 102)]))
    result = DailyEquityRefreshService(
        settings, provider=provider, now=datetime(2026, 8, 19, 10, 0, tzinfo=UTC),
    ).run()
    assert result.status == "SUCCESS"
    assert result.rows_added == 1 and result.rejected_current_rows == 1
    assert pd.Timestamp(store.read(store.equities_path)["date"].max()).date() == date(2026, 8, 18)


def test_current_session_requires_time_and_cross_symbol_coverage(tmp_path: Path) -> None:
    settings = _settings(tmp_path); store = DatasetStore(settings.data.root)
    existing = pd.concat([_equity("AAA", [("2026-08-18", 100)]), _equity("BBB", [("2026-08-18", 200)])])
    DatasetStore._atomic_parquet(existing, store.equities_path)
    provider = FakeYahoo(_equity("AAA", [("2026-08-19", 101)]))
    result = DailyEquityRefreshService(
        settings, provider=provider, now=datetime(2026, 8, 19, 13, 0, tzinfo=UTC),
        timing=SessionTiming(minimum_coverage=.8),
    ).run()
    assert result.status == "SESSION_NOT_FINAL"
    assert result.failed == 1
    assert result.rejected_current_rows == 1 and result.rows_added == 0


def test_official_ingest_preserves_hash_url_and_read_only_original(tmp_path: Path) -> None:
    settings = _settings(tmp_path)
    folder = Path(settings.data.inbox) / "daily/2026-08-18"; folder.mkdir(parents=True)
    artifact = folder / "market_summary_2026-08-18.csv"
    artifact.write_text("symbol,open,high,low,close,volume\nAAA,10,11,9,10.5,1000\n", encoding="utf-8")
    Path(str(artifact) + ".source.json").write_text(json.dumps({
        "source_url": "https://dps.psx.com.pk/download/example",
        "downloaded_at": "2026-08-19T12:00:00+05:00",
    }), encoding="utf-8")
    result = OfficialDailyIngestService(settings).run()
    assert result.status == "SUCCESS" and result.equity_rows == 1
    manifest = pd.read_csv(result.manifest)
    preserved = Path(settings.data.root) / manifest.iloc[0]["raw_artifact"]
    assert manifest.iloc[0]["provider"] == OFFICIAL_PROVIDER
    assert len(manifest.iloc[0]["sha256"]) == 64
    assert preserved.stat().st_mode & 0o222 == 0


def test_official_ingest_refuses_missing_source_metadata(tmp_path: Path) -> None:
    settings = _settings(tmp_path)
    folder = Path(settings.data.inbox) / "daily/2026-08-18"; folder.mkdir(parents=True)
    (folder / "market_summary_2026-08-18.csv").write_text(
        "symbol,open,high,low,close,volume\nAAA,10,11,9,10.5,1000\n", encoding="utf-8"
    )
    result = OfficialDailyIngestService(settings).run()
    assert result.status == "FAILED"
    assert "Missing provenance sidecar" in result.failures[0]["error"]


def test_reconciliation_classifies_and_flags_selected_conflict(tmp_path: Path) -> None:
    settings = _settings(tmp_path); store = DatasetStore(settings.data.root)
    yahoo = _equity("AAA", [("2026-08-19", 100)])
    official = _equity("AAA", [("2026-08-19", 110)])
    official["provider"] = OFFICIAL_PROVIDER
    DatasetStore._atomic_parquet(yahoo, store.equities_path)
    ref = store.root / "normalized/official_reference/equities.parquet"
    DatasetStore._atomic_parquet(official, ref)
    study = tmp_path / "study"; study.mkdir()
    (study / "signals.jsonl").write_text(json.dumps({
        "signal_id": "S1", "as_of_date": "2026-08-19", "symbol": "AAA",
    }) + "\n", encoding="utf-8")
    (study / "operational_log.jsonl").touch()
    result = DailyProviderReconciler(settings, study, ComparisonThresholds()).run()
    assert result["status"] == "DATA_SOURCE_CONFLICT"
    assert result["MATERIAL_DIFFERENCE"] == 1
    assert "DATA_SOURCE_CONFLICT" in (study / "operational_log.jsonl").read_text()


def test_manual_benchmark_is_idempotent_and_conflict_is_append_only(tmp_path: Path) -> None:
    settings = _settings(tmp_path); study = tmp_path / "study"
    service = BenchmarkAppendService(settings, study)
    args = ("KSE100", date(2026, 8, 18), 151234.5, OFFICIAL_PROVIDER, "official daily file")
    assert service.add(*args)["status"] == "ADDED"
    assert service.add(*args)["status"] == "NO_CHANGE"
    assert service.add("KSE100", date(2026, 8, 18), 151000, OFFICIAL_PROVIDER, "correction")["status"] == "CONFLICT"
    frame = DatasetStore(settings.data.root).read(DatasetStore(settings.data.root).kse100_path)
    assert len(frame) == 1 and frame.iloc[0]["close"] == 151234.5
    assert "CORRECTION_CONFLICT" in (study / "benchmark_audit.jsonl").read_text()
    transition = pd.read_csv(study / "benchmark_transition.csv")
    assert transition.iloc[-1]["continuity_status"] == "OK"


def test_membership_status_unknown_verified_and_conflict(tmp_path: Path) -> None:
    settings = _settings(tmp_path); store = DatasetStore(settings.data.root)
    assert membership_status(settings, date(2026, 8, 19)) == "UNKNOWN"
    path = store.constituents_path
    frame = pd.DataFrame({
        "index": ["KSE100"], "symbol": ["AAA"],
        "effective_from": pd.to_datetime(["2026-04-01"]), "effective_to": [pd.NaT],
        "verification_status": ["VERIFIED_MEMBER"],
    })
    DatasetStore._atomic_parquet(frame, path)
    assert membership_status(settings, date(2026, 8, 19)) == "VERIFIED"
    frame.loc[0, "verification_status"] = "CONFLICT"
    DatasetStore._atomic_parquet(frame, path)
    assert membership_status(settings, date(2026, 8, 19)) == "CONFLICT"


def test_orchestrator_waits_for_benchmark_without_calling_eod(tmp_path: Path, monkeypatch) -> None:
    settings = _settings(tmp_path); store = DatasetStore(settings.data.root)
    DatasetStore._atomic_parquet(_equity("AAA", [("2026-08-19", 100)]), store.equities_path)
    fake_result = type("R", (), {"run": lambda self: EquityRefreshResult(
        "UP_TO_DATE", "2026-08-20", "2026-08-19", "2026-08-19", 1, 0, 0, 0, 0,
    )})()
    benchmark_refresh = type("B", (), {"run": lambda self: BenchmarkRefreshResult(
        "WAITING_FOR_BENCHMARK_PROVIDER", None, None, None, 0, None,
    )})()
    result = ProspectiveDailyService(
        settings, tmp_path / "study", now=datetime(2026, 8, 19, 13, tzinfo=UTC),
        refresh=fake_result, benchmark_refresh=benchmark_refresh,
    ).run()
    assert result["status"] == "WAITING_FOR_BENCHMARK_PROVIDER"
    assert result["paper_eod"]["status"] == "NOT_RUN"
    assert (tmp_path / "study/provider_health.csv").exists()
    assert (tmp_path / "study/daily/2026-08-19.md").exists()


def test_orchestration_updates_data_before_positions_and_new_signal(tmp_path: Path, monkeypatch) -> None:
    settings = _settings(tmp_path); store = DatasetStore(settings.data.root)
    DatasetStore._atomic_parquet(_equity("AAA", [("2026-08-19", 100)]), store.equities_path)
    DatasetStore._atomic_parquet(pd.DataFrame({
        "index": ["KSE100"], "date": pd.to_datetime(["2026-08-19"]), "close": [150000],
        "provider": [OFFICIAL_PROVIDER],
    }), store.kse100_path)
    order: list[str] = []

    class Refresh:
        def run(self):
            order.append("refresh")
            return EquityRefreshResult("UP_TO_DATE", None, None, "2026-08-19", 0, 0, 0, 0, 0)

    class BenchmarkRefresh:
        def run(self):
            order.append("benchmark")
            return BenchmarkRefreshResult(
                "NO_NEW_DATA", "2026-08-19", None, None, 0, "2026-08-19",
                transition_status="VALIDATED_RESEARCH_TRANSITION", validation="PASS",
            )

    monkeypatch.setattr("psx_signal.pipelines.daily_operations.membership_status",
                        lambda *args: order.append("membership") or "VERIFIED")
    monkeypatch.setattr(ProspectiveDailyService, "_validation",
                        lambda *args: order.append("validate") or {"status": "PASS", "valid": True})
    monkeypatch.setattr("psx_signal.pipelines.daily_operations.ProspectivePaperService.update",
                        lambda self: order.append("update") or {"status": "NO_CHANGES"})
    monkeypatch.setattr("psx_signal.pipelines.daily_operations.ProspectivePaperService.eod",
                        lambda self, strategy, as_of: order.append("eod") or {"status": "NO_ELIGIBLE_SIGNAL"})
    ProspectiveDailyService(
        settings, tmp_path / "study", now=datetime(2026, 8, 19, 13, tzinfo=UTC),
        refresh=Refresh(), benchmark_refresh=BenchmarkRefresh(),
    ).run()
    assert order[:6] == ["refresh", "benchmark", "membership", "validate", "update", "eod"]


def test_provider_health_uses_pkt_iso_and_null_lag_for_not_final(tmp_path: Path) -> None:
    settings = _settings(tmp_path); study = tmp_path / "study"; study.mkdir()
    path = study / "provider_health.csv"
    columns = [
        "recorded_at", "as_of_date", "yahoo_status", "yahoo_successful", "yahoo_failed",
        "yahoo_rows_added", "expected_session", "available_at", "publication_lag_minutes",
        "rejected_observations", "official_status", "official_artifacts_imported",
        "comparison_rows", "material_differences", "validation_status", "notes",
    ]
    pd.DataFrame([{
        "recorded_at": "46253.43373", "as_of_date": "2026-08-18",
        "yahoo_status": "SESSION_NOT_FINAL", "publication_lag_minutes": "1439",
    }], columns=columns).to_csv(path, index=False)
    original = path.read_bytes()
    service = ProspectiveDailyService(
        settings, study, now=datetime(2026, 8, 19, 10, 34, 22, tzinfo=UTC),
    )
    service._provider_health(
        EquityRefreshResult("SESSION_NOT_FINAL", None, None, "2026-08-18", 0, 0, 0, 0, 0),
        OfficialIngestResult("NO_NEW_FILES", 0, 0, 0, 0, []), {"rows": 0}, date(2026, 8, 18),
    )
    assert path.read_bytes().startswith(original)
    with path.open(newline="", encoding="utf-8") as handle:
        latest = list(csv.DictReader(handle))[-1]
    assert latest["recorded_at"] == "2026-08-19T15:34:22+05:00"
    assert latest["available_at"] == ""
    assert latest["publication_lag_minutes"] == ""
    normalized = pd.read_csv(study / "provider_health_normalized.csv", dtype="string")
    assert normalized.iloc[0]["recorded_at"].endswith("+05:00")
    assert pd.isna(normalized.iloc[0]["available_at"])
    assert pd.isna(normalized.iloc[0]["publication_lag_minutes"])


def test_completed_session_provider_health_calculates_publication_lag(tmp_path: Path) -> None:
    settings = _settings(tmp_path); study = tmp_path / "study"
    service = ProspectiveDailyService(
        settings, study, now=datetime(2026, 8, 19, 12, 0, tzinfo=UTC),
    )
    service._provider_health(
        EquityRefreshResult("SUCCESS", None, None, "2026-08-19", 97, 97, 0, 97, 0),
        OfficialIngestResult("NO_NEW_FILES", 0, 0, 0, 0, []), {"rows": 0}, date(2026, 8, 19),
    )
    with (study / "provider_health.csv").open(newline="", encoding="utf-8") as handle:
        row = next(csv.DictReader(handle))
    assert row["recorded_at"] == "2026-08-19T17:00:00+05:00"
    assert row["available_at"] == "2026-08-19T17:00:00+05:00"
    assert row["publication_lag_minutes"] == "90"


@pytest.mark.parametrize(
    ("now", "expected"),
    [
        (datetime(2026, 8, 19, 20, tzinfo=UTC), "LATE_EOD_GENERATION"),
        (datetime(2026, 8, 20, 5, tzinfo=UTC), "MISSED_DATA_CUTOFF"),
    ],
)
def test_timing_labels_are_injectable(tmp_path: Path, monkeypatch, now: datetime, expected: str) -> None:
    settings = _settings(tmp_path); store = DatasetStore(settings.data.root)
    DatasetStore._atomic_parquet(_equity("AAA", [("2026-08-19", 100)]), store.equities_path)
    benchmark = pd.DataFrame({"index": ["KSE100"], "date": pd.to_datetime(["2026-08-19"]),
                              "close": [150000], "provider": [OFFICIAL_PROVIDER]})
    DatasetStore._atomic_parquet(benchmark, store.kse100_path)
    refresh = type("Refresh", (), {"run": lambda self: EquityRefreshResult(
        "UP_TO_DATE", None, None, "2026-08-19", 1, 0, 0, 0, 0,
    )})()
    benchmark_refresh = type("BenchmarkRefresh", (), {"run": lambda self: BenchmarkRefreshResult(
        "NO_NEW_DATA", "2026-08-19", None, None, 0, "2026-08-19",
        transition_status="VALIDATED_RESEARCH_TRANSITION", validation="PASS",
    )})()
    monkeypatch.setattr("psx_signal.pipelines.daily_operations.ProspectivePaperService.update",
                        lambda self: {"status": "NO_CHANGES"})
    monkeypatch.setattr("psx_signal.pipelines.daily_operations.ProspectivePaperService.eod",
                        lambda self, strategy, as_of: {"status": "NO_ELIGIBLE_SIGNAL"})
    result = ProspectiveDailyService(
        settings, tmp_path / "study", now=now, refresh=refresh,
        benchmark_refresh=benchmark_refresh,
    ).run()
    assert result["status"] == expected
