#!/usr/bin/env python3
"""Audit Freqtrade research Feather data without modifying the source files.

The script is intentionally independent of a running bot.  It reads the four
datasets needed by the systematic-trading validation work, records provenance
and coverage, and writes machine-readable audit artifacts.

Run inside the Freqtrade container when the host does not have ``pyarrow``:

    python3 /tmp/audit_research_data.py \
      --data-dir /freqtrade/user_data/data/binance/futures \
      --config /freqtrade/user_data/config.json \
      --output-dir /tmp/systematic-trading-data-audit \
      --source-role production_runtime_export

The earliest observation is a coverage boundary, not proof of a listing date.
Pass ``--listing-dates`` when authoritative point-in-time listing metadata is
available.  Funding rates are event observations whose schedule can change;
the audit only flags intervals larger than ``--funding-max-gap-hours``.
"""

from __future__ import annotations

import argparse
import csv
import hashlib
import json
import platform
import socket
import sys
from collections import Counter
from dataclasses import dataclass
from datetime import datetime, timezone
from pathlib import Path
from typing import Any, Iterable

import pandas as pd

SCRIPT_VERSION = "1.1"
DEFAULT_PAIRS = (
    "BTC/USDT:USDT",
    "ETH/USDT:USDT",
    "SOL/USDT:USDT",
    "BNB/USDT:USDT",
    "XRP/USDT:USDT",
    "DOGE/USDT:USDT",
    "ADA/USDT:USDT",
    "ZEC/USDT:USDT",
    "EDGE/USDT:USDT",
    "POWER/USDT:USDT",
    "VANRY/USDT:USDT",
    "UNI/USDT:USDT",
)


@dataclass(frozen=True)
class DatasetSpec:
    key: str
    suffix: str
    expected_seconds: int | None
    grid_kind: str


DATASETS = (
    DatasetSpec("futures_5m", "5m-futures", 5 * 60, "fixed"),
    DatasetSpec("futures_1d", "1d-futures", 24 * 60 * 60, "fixed"),
    DatasetSpec("funding_1h", "1h-funding_rate", None, "funding_event"),
    DatasetSpec("mark_1h", "1h-mark", 60 * 60, "fixed"),
)

COVERAGE_FIELDS = (
    "pair",
    "dataset",
    "path",
    "status",
    "classification",
    "rows",
    "bytes",
    "sha256",
    "min_time",
    "max_time",
    "timezone",
    "is_monotonic",
    "duplicate_timestamps",
    "invalid_timestamps",
    "nan_cells",
    "columns_with_nan",
    "expected_grid",
    "expected_seconds",
    "off_grid_timestamps",
    "internal_gap_count",
    "missing_grid_points",
    "largest_delta_seconds",
    "listing_time",
    "pre_listing",
    "xsmom_28d_first_feature_asof",
    "xsmom_28d_first_eligible_at",
    "xsmom_28d_latest_endpoint_usable",
    "funding_coverage_feasibility",
)

GAP_FIELDS = (
    "pair",
    "dataset",
    "classification",
    "previous_time",
    "next_time",
    "delta_seconds",
    "expected_seconds",
    "missing_grid_points",
)


def utc_now() -> str:
    return datetime.now(timezone.utc).isoformat()


def sha256_file(path: Path, chunk_size: int = 1024 * 1024) -> str:
    digest = hashlib.sha256()
    with path.open("rb") as handle:
        for chunk in iter(lambda: handle.read(chunk_size), b""):
            digest.update(chunk)
    return digest.hexdigest()


def pair_file(data_dir: Path, pair: str, suffix: str) -> Path:
    contract, *settle_parts = pair.split(":", 1)
    base, quote = contract.split("/", 1)
    settle = settle_parts[0] if settle_parts else quote
    return data_dir / f"{base}_{quote}_{settle}-{suffix}.feather"


def parse_pairs(raw_values: Iterable[str]) -> list[str]:
    pairs: list[str] = []
    for raw in raw_values:
        for pair in raw.split(","):
            pair = pair.strip()
            if pair and pair not in pairs:
                pairs.append(pair)
    return pairs


def load_pairs(config: Path | None, raw_pairs: list[str]) -> list[str]:
    if raw_pairs:
        return parse_pairs(raw_pairs)
    if config is not None:
        payload = json.loads(config.read_text(encoding="utf-8"))
        pairs = payload.get("exchange", {}).get("pair_whitelist")
        if not isinstance(pairs, list) or not all(isinstance(item, str) for item in pairs):
            raise ValueError(f"{config}: exchange.pair_whitelist must be a string list")
        return parse_pairs(pairs)
    return list(DEFAULT_PAIRS)


def load_listing_dates(path: Path | None) -> dict[str, pd.Timestamp]:
    if path is None:
        return {}
    payload = json.loads(path.read_text(encoding="utf-8"))
    if not isinstance(payload, dict):
        raise ValueError("--listing-dates must contain a JSON object mapping pair to timestamp")
    result: dict[str, pd.Timestamp] = {}
    for pair, value in payload.items():
        stamp = pd.Timestamp(value)
        result[str(pair)] = stamp.tz_localize("UTC") if stamp.tz is None else stamp.tz_convert("UTC")
    return result


def iso_or_none(value: pd.Timestamp | None) -> str | None:
    if value is None or pd.isna(value):
        return None
    return value.isoformat()


def first_full_listing_day(listing_time: pd.Timestamp) -> pd.Timestamp:
    """Return the first UTC daily candle that is not a partial listing day."""
    day = listing_time.normalize()
    return day if listing_time == day else day + pd.Timedelta(days=1)


def first_contiguous_window_end(
    timestamps: pd.DatetimeIndex,
    *,
    points: int,
    interval: pd.Timedelta,
) -> pd.Timestamp | None:
    """Return the end of the first exact fixed-interval window."""
    if len(timestamps) < points:
        return None
    interval_ns = int(interval.value)
    values = timestamps.to_numpy(dtype="datetime64[ns]").astype("int64")
    run_length = 1
    for index in range(1, len(values)):
        run_length = run_length + 1 if values[index] - values[index - 1] == interval_ns else 1
        if run_length >= points:
            return timestamps[index]
    return None


def serializable_scalar(value: Any) -> Any:
    if value is None or pd.isna(value):
        return None
    if hasattr(value, "item"):
        return value.item()
    return value


def schema_of(frame: pd.DataFrame) -> list[dict[str, Any]]:
    schema = []
    for column in frame.columns:
        series = frame[column]
        schema.append(
            {
                "name": str(column),
                "dtype": str(series.dtype),
                "nullable": bool(series.isna().any()),
            }
        )
    return schema


def gap_rows(
    pair: str,
    spec: DatasetSpec,
    unique_times: pd.DatetimeIndex,
    funding_max_gap_seconds: int,
) -> tuple[list[dict[str, Any]], dict[str, Any]]:
    if len(unique_times) < 2:
        return [], {
            "off_grid_timestamps": 0,
            "internal_gap_count": 0,
            "missing_grid_points": 0,
            "largest_delta_seconds": None,
            "interval_distribution_seconds": {},
        }

    # pandas 3 preserves Arrow's timestamp unit (often ms/us), so ``asi8`` is
    # no longer guaranteed to mean nanoseconds. Normalize explicitly.
    nanoseconds = unique_times.to_numpy(dtype="datetime64[ns]").astype("int64")
    deltas = (nanoseconds[1:] - nanoseconds[:-1]) // 1_000_000_000
    interval_counts = Counter(int(value) for value in deltas)
    largest = int(deltas.max())

    if spec.grid_kind == "fixed":
        assert spec.expected_seconds is not None
        expected = spec.expected_seconds
        off_grid = int(((nanoseconds // 1_000_000_000) % expected != 0).sum())
        threshold = expected
        classification = "internal_grid_gap"
    else:
        # Funding events can legitimately move between 1h/2h/4h/8h schedules.
        # More than the configured maximum is a candidate gap, not proof that
        # an exchange event is missing.
        expected = funding_max_gap_seconds
        off_grid = int(((nanoseconds // 1_000_000_000) % 3600 != 0).sum())
        threshold = funding_max_gap_seconds
        classification = "funding_event_gap_candidate"

    gaps: list[dict[str, Any]] = []
    missing_total = 0
    for index, delta in enumerate(deltas):
        delta_int = int(delta)
        if delta_int <= threshold:
            continue
        missing = max(0, (delta_int - 1) // expected)
        missing_total += missing
        gaps.append(
            {
                "pair": pair,
                "dataset": spec.key,
                "classification": classification,
                "previous_time": unique_times[index].isoformat(),
                "next_time": unique_times[index + 1].isoformat(),
                "delta_seconds": delta_int,
                "expected_seconds": expected,
                "missing_grid_points": missing,
            }
        )

    return gaps, {
        "off_grid_timestamps": off_grid,
        "internal_gap_count": len(gaps),
        "missing_grid_points": missing_total,
        "largest_delta_seconds": largest,
        "interval_distribution_seconds": {
            str(seconds): count for seconds, count in sorted(interval_counts.items())
        },
    }


def analyze_frame(
    frame: pd.DataFrame,
    *,
    pair: str,
    spec: DatasetSpec,
    funding_max_gap_seconds: int,
    listing_time: pd.Timestamp | None,
) -> tuple[dict[str, Any], list[dict[str, Any]]]:
    if "date" not in frame.columns:
        raise ValueError("required date column is missing")

    raw_dates = frame["date"]
    dates = pd.to_datetime(raw_dates, utc=True, errors="coerce")
    valid_dates = dates.dropna()
    sorted_unique = pd.DatetimeIndex(valid_dates.drop_duplicates().sort_values())
    duplicate_count = int(valid_dates.duplicated().sum())
    invalid_count = int(dates.isna().sum())
    nan_by_column = {str(column): int(frame[column].isna().sum()) for column in frame.columns}
    gap_detail, gap_summary = gap_rows(
        pair, spec, sorted_unique, funding_max_gap_seconds
    )

    min_time = sorted_unique.min() if len(sorted_unique) else None
    max_time = sorted_unique.max() if len(sorted_unique) else None
    timezone_name = str(sorted_unique.tz) if len(sorted_unique) else None
    pre_listing = (
        bool(min_time < listing_time) if min_time is not None and listing_time is not None else None
    )
    xsmom_fields: dict[str, Any] = {
        "xsmom_28d_first_feature_asof": None,
        "xsmom_28d_first_eligible_at": None,
        "xsmom_28d_latest_endpoint_usable": None,
    }
    if spec.key == "futures_1d":
        # A 28-day return needs close[t] and close[t-28].  The feature at the
        # t daily close becomes causally eligible only after that candle has
        # completed, one daily interval later. When authoritative listing
        # metadata is available, discard every pre-listing row and an intraday
        # listing day's partial 1d candle before looking for a complete window.
        window_points = 29
        eligible_daily_times = sorted_unique
        if listing_time is not None:
            eligible_daily_times = eligible_daily_times[
                eligible_daily_times >= first_full_listing_day(listing_time)
            ]
        if len(eligible_daily_times) >= window_points:
            first_asof = first_contiguous_window_end(
                eligible_daily_times,
                points=window_points,
                interval=pd.Timedelta(days=1),
            )
            latest_window = eligible_daily_times[-window_points:]
            latest_ns = latest_window.to_numpy(dtype="datetime64[ns]").astype("int64")
            expected_ns = 24 * 60 * 60 * 1_000_000_000
            endpoint_usable = bool(((latest_ns[1:] - latest_ns[:-1]) == expected_ns).all())
            xsmom_fields = {
                "xsmom_28d_first_feature_asof": iso_or_none(first_asof),
                "xsmom_28d_first_eligible_at": (
                    (first_asof + pd.Timedelta(days=1)).isoformat()
                    if first_asof is not None
                    else None
                ),
                "xsmom_28d_latest_endpoint_usable": endpoint_usable,
            }

    funding_feasibility = None
    if spec.key == "funding_1h":
        if len(sorted_unique) < 2:
            funding_feasibility = "insufficient_observations"
        elif gap_summary["internal_gap_count"]:
            funding_feasibility = "usable_with_gap_candidates"
        else:
            funding_feasibility = "usable_event_coverage"

    if frame.empty:
        classification = "empty_file"
    elif invalid_count == len(frame):
        classification = "invalid_time_column"
    elif gap_summary["internal_gap_count"]:
        classification = "has_internal_gaps"
    else:
        classification = "continuous_within_observed_boundaries"

    expected_seconds = (
        spec.expected_seconds
        if spec.grid_kind == "fixed"
        else funding_max_gap_seconds
    )
    record = {
        "status": "ok" if not frame.empty and invalid_count < len(frame) else "invalid",
        "classification": classification,
        "rows": int(len(frame)),
        "schema": schema_of(frame),
        "min_time": iso_or_none(min_time),
        "max_time": iso_or_none(max_time),
        "timezone": timezone_name,
        "is_monotonic": bool(valid_dates.is_monotonic_increasing),
        "duplicate_timestamps": duplicate_count,
        "invalid_timestamps": invalid_count,
        "nan_cells": int(sum(nan_by_column.values())),
        "nan_by_column": nan_by_column,
        "columns_with_nan": sum(1 for count in nan_by_column.values() if count),
        "expected_grid": (
            f"fixed_{spec.expected_seconds}s"
            if spec.grid_kind == "fixed"
            else f"funding_event_max_gap_{funding_max_gap_seconds}s"
        ),
        "expected_seconds": expected_seconds,
        "boundary_classification": {
            "start": "observed_start_not_inferred_as_gap",
            "end": "observed_end_not_inferred_as_gap",
        },
        "listing_time": iso_or_none(listing_time),
        "pre_listing": pre_listing,
        **xsmom_fields,
        "funding_coverage_feasibility": funding_feasibility,
        **gap_summary,
    }
    return record, gap_detail


def audit_file(
    *,
    data_dir: Path,
    pair: str,
    spec: DatasetSpec,
    funding_max_gap_seconds: int,
    listing_time: pd.Timestamp | None,
) -> tuple[dict[str, Any], list[dict[str, Any]]]:
    path = pair_file(data_dir, pair, spec.suffix)
    base = {
        "pair": pair,
        "dataset": spec.key,
        "path": str(path),
        "exists": path.exists(),
    }
    if not path.exists():
        return {
            **base,
            "status": "missing",
            "classification": "missing_file",
            "rows": 0,
            "bytes": 0,
            "sha256": None,
            "schema": [],
            "min_time": None,
            "max_time": None,
            "timezone": None,
            "is_monotonic": None,
            "duplicate_timestamps": None,
            "invalid_timestamps": None,
            "nan_cells": None,
            "nan_by_column": {},
            "columns_with_nan": None,
            "expected_grid": (
                f"fixed_{spec.expected_seconds}s"
                if spec.grid_kind == "fixed"
                else f"funding_event_max_gap_{funding_max_gap_seconds}s"
            ),
            "expected_seconds": spec.expected_seconds or funding_max_gap_seconds,
            "off_grid_timestamps": None,
            "internal_gap_count": None,
            "missing_grid_points": None,
            "largest_delta_seconds": None,
            "interval_distribution_seconds": {},
            "listing_time": iso_or_none(listing_time),
            "pre_listing": None,
            "xsmom_28d_first_feature_asof": None,
            "xsmom_28d_first_eligible_at": None,
            "xsmom_28d_latest_endpoint_usable": None,
            "funding_coverage_feasibility": (
                "missing_file" if spec.key == "funding_1h" else None
            ),
        }, []

    file_meta = {
        "bytes": path.stat().st_size,
        "sha256": sha256_file(path),
    }
    try:
        frame = pd.read_feather(path)
        analysis, gaps = analyze_frame(
            frame,
            pair=pair,
            spec=spec,
            funding_max_gap_seconds=funding_max_gap_seconds,
            listing_time=listing_time,
        )
        return {**base, **file_meta, **analysis}, gaps
    except Exception as exc:
        return {
            **base,
            **file_meta,
            "status": "error",
            "classification": "read_error",
            "error": f"{type(exc).__name__}: {exc}",
            "rows": None,
            "schema": [],
            "min_time": None,
            "max_time": None,
            "timezone": None,
            "is_monotonic": None,
            "duplicate_timestamps": None,
            "invalid_timestamps": None,
            "nan_cells": None,
            "nan_by_column": {},
            "columns_with_nan": None,
            "expected_grid": (
                f"fixed_{spec.expected_seconds}s"
                if spec.grid_kind == "fixed"
                else f"funding_event_max_gap_{funding_max_gap_seconds}s"
            ),
            "expected_seconds": spec.expected_seconds or funding_max_gap_seconds,
            "off_grid_timestamps": None,
            "internal_gap_count": None,
            "missing_grid_points": None,
            "largest_delta_seconds": None,
            "interval_distribution_seconds": {},
            "listing_time": iso_or_none(listing_time),
            "pre_listing": None,
            "xsmom_28d_first_feature_asof": None,
            "xsmom_28d_first_eligible_at": None,
            "xsmom_28d_latest_endpoint_usable": None,
            "funding_coverage_feasibility": (
                "read_error" if spec.key == "funding_1h" else None
            ),
        }, []


def csv_value(value: Any) -> Any:
    if isinstance(value, (dict, list)):
        return json.dumps(value, ensure_ascii=False, sort_keys=True)
    return serializable_scalar(value)


def write_csv(path: Path, fields: tuple[str, ...], rows: list[dict[str, Any]]) -> None:
    with path.open("w", encoding="utf-8", newline="") as handle:
        writer = csv.DictWriter(handle, fieldnames=fields, extrasaction="ignore")
        writer.writeheader()
        for row in rows:
            writer.writerow({field: csv_value(row.get(field)) for field in fields})


def readme_text(manifest: dict[str, Any]) -> str:
    summary = manifest["summary"]
    return f"""# Systematic trading research-data audit

Generated: `{manifest["generated_at"]}`

Source role: `{manifest["source_role"]}`

Data directory: `{manifest["data_dir"]}`

This is a read-only coverage and integrity audit. It does **not** prove which
strategy or dataset is deployed. A result with source role `research_cache` is
provisional; deployment/runtime truth requires running the same script against
the production container's mounted data and preserving that exported manifest.

## Outputs

- `data_manifest.json`: full provenance, schema, hashes, NaN and interval data.
- `data_coverage.csv`: one flattened row per pair and dataset.
- `gap_details.csv`: fixed-grid internal gaps and conservative funding-event
  gap candidates.

## Summary

- Pairs: {summary["pair_count"]}
- Expected files: {summary["expected_file_count"]}
- Read successfully: {summary["ok_file_count"]}
- Missing files: {summary["missing_file_count"]}
- Read/validation errors: {summary["error_file_count"]}
- Files with internal gap candidates: {summary["files_with_gaps"]}
- Detailed gap rows: {summary["gap_count"]}

## Interpretation rules

- The first and last observed timestamps are coverage boundaries. They are not
  counted as internal gaps and are not automatically treated as listing dates.
- `pre_listing` is only populated when authoritative metadata is supplied with
  `--listing-dates`; the script never guesses a listing date from the first bar.
- 5m futures, 1d futures and 1h mark data use fixed expected grids.
- Funding is an event series and Binance can change the settlement interval.
  Only deltas greater than the configured {manifest["funding_max_gap_hours"]:g}h
  maximum are emitted as `funding_event_gap_candidate`; this is conservative
  and is not proof that an event is missing.
- Daily rows report the first causally eligible 28-day XSMOM endpoint and
  whether the latest 29 eligible daily observations are contiguous. When
  authoritative listing metadata is supplied, pre-listing rows and an intraday
  listing day's partial candle are excluded.
- Funding rows report `funding_coverage_feasibility`; this only assesses event
  coverage and does not equate settled historical funding with the live
  exchange's predicted/current funding response.
- SHA256 covers the source Feather file. Source files are never modified.
"""


def build_parser() -> argparse.ArgumentParser:
    parser = argparse.ArgumentParser(description=__doc__)
    parser.add_argument(
        "--data-dir",
        type=Path,
        required=True,
        help="directory containing Binance futures Feather files",
    )
    parser.add_argument(
        "--config",
        type=Path,
        help="Freqtrade JSON config; reads exchange.pair_whitelist",
    )
    parser.add_argument(
        "--pairs",
        action="append",
        default=[],
        help="comma-separated pair list; repeatable and takes precedence over --config",
    )
    parser.add_argument("--output-dir", type=Path, required=True)
    parser.add_argument(
        "--listing-dates",
        type=Path,
        help="optional JSON mapping pair to authoritative listing timestamp",
    )
    parser.add_argument(
        "--funding-max-gap-hours",
        type=float,
        default=8.0,
        help="flag funding event deltas larger than this (default: 8)",
    )
    parser.add_argument(
        "--source-role",
        choices=("research_cache", "production_runtime_export", "other"),
        default="research_cache",
        help="provenance label; research_cache must not be reported as deployment truth",
    )
    parser.add_argument("--audit-label", default="systematic-trading-validation-2026-07-29")
    return parser


def main() -> int:
    args = build_parser().parse_args()
    if args.funding_max_gap_hours <= 0:
        raise SystemExit("--funding-max-gap-hours must be positive")
    if not args.data_dir.is_dir():
        raise SystemExit(f"data directory does not exist: {args.data_dir}")
    if args.config is not None and not args.config.is_file():
        raise SystemExit(f"config does not exist: {args.config}")

    pairs = load_pairs(args.config, args.pairs)
    if not pairs:
        raise SystemExit("pair list is empty")
    listing_dates = load_listing_dates(args.listing_dates)
    funding_max_gap_seconds = int(args.funding_max_gap_hours * 3600)

    records: list[dict[str, Any]] = []
    gaps: list[dict[str, Any]] = []
    for pair in pairs:
        for spec in DATASETS:
            record, file_gaps = audit_file(
                data_dir=args.data_dir,
                pair=pair,
                spec=spec,
                funding_max_gap_seconds=funding_max_gap_seconds,
                listing_time=listing_dates.get(pair),
            )
            records.append(record)
            gaps.extend(file_gaps)

    summary = {
        "pair_count": len(pairs),
        "expected_file_count": len(records),
        "ok_file_count": sum(record["status"] == "ok" for record in records),
        "missing_file_count": sum(record["status"] == "missing" for record in records),
        "error_file_count": sum(record["status"] in {"error", "invalid"} for record in records),
        "files_with_gaps": sum(bool(record.get("internal_gap_count")) for record in records),
        "gap_count": len(gaps),
        "duplicate_timestamp_count": sum(
            record.get("duplicate_timestamps") or 0 for record in records
        ),
        "nan_cell_count": sum(record.get("nan_cells") or 0 for record in records),
    }
    script_path = Path(__file__).resolve()
    manifest = {
        "audit_label": args.audit_label,
        "script_version": SCRIPT_VERSION,
        "script_path": str(script_path),
        "script_sha256": sha256_file(script_path),
        "generated_at": utc_now(),
        "generated_host": socket.gethostname(),
        "source_role": args.source_role,
        "data_dir": str(args.data_dir.resolve()),
        "config": str(args.config.resolve()) if args.config else None,
        "config_sha256": sha256_file(args.config) if args.config else None,
        "listing_dates_source": (
            str(args.listing_dates.resolve()) if args.listing_dates else None
        ),
        "listing_dates_sha256": (
            sha256_file(args.listing_dates) if args.listing_dates else None
        ),
        "funding_max_gap_hours": args.funding_max_gap_hours,
        "pairs": pairs,
        "datasets": [
            {
                "key": spec.key,
                "suffix": spec.suffix,
                "grid_kind": spec.grid_kind,
                "expected_seconds": spec.expected_seconds,
            }
            for spec in DATASETS
        ],
        "environment": {
            "python": sys.version,
            "platform": platform.platform(),
            "pandas": pd.__version__,
        },
        "summary": summary,
        "files": records,
    }

    args.output_dir.mkdir(parents=True, exist_ok=True)
    manifest_path = args.output_dir / "data_manifest.json"
    coverage_path = args.output_dir / "data_coverage.csv"
    gaps_path = args.output_dir / "gap_details.csv"
    readme_path = args.output_dir / "README.md"
    manifest_path.write_text(
        json.dumps(manifest, indent=2, ensure_ascii=False, sort_keys=True) + "\n",
        encoding="utf-8",
    )
    write_csv(coverage_path, COVERAGE_FIELDS, records)
    write_csv(gaps_path, GAP_FIELDS, gaps)
    readme_path.write_text(readme_text(manifest), encoding="utf-8")

    print(json.dumps(summary, ensure_ascii=False, sort_keys=True))
    print(f"wrote {manifest_path}")
    print(f"wrote {coverage_path}")
    print(f"wrote {gaps_path}")
    print(f"wrote {readme_path}")
    return 1 if summary["missing_file_count"] or summary["error_file_count"] else 0


if __name__ == "__main__":
    raise SystemExit(main())
