"""
data_store.py

All persistence lives here: premarket snapshots, intraday snapshots,
trade records, and bot state (positions, watchlists, cooldowns).

Design notes:
- Uses simple JSON-lines files per day for premarket/intraday data so
  they can be analyzed later with pandas without a database dependency.
- State files (positions.json, watchlist.json, cooldowns.json) are
  written atomically (write to tmp file, then os.replace) to avoid
  corruption if the process is killed mid-write.
- A single filelock-style guard (using a lock file + fcntl) protects
  state writes in case of accidental concurrent processes, matching
  a lesson learned in prior systems where duplicate processes raced
  on shared state.
"""

import json
import os
import fcntl
import time
from datetime import datetime, date
from pathlib import Path
from contextlib import contextmanager

from config_loader import get_config
from logger_setup import get_logger

log = get_logger("data_store")

BASE_DIR = Path(__file__).resolve().parent


def _dirs():
    cfg = get_config()["data_storage"]
    return {
        "premarket": BASE_DIR / cfg["premarket_dir"],
        "intraday": BASE_DIR / cfg["intraday_dir"],
        "trades": BASE_DIR / cfg["trades_dir"],
        "state": BASE_DIR / cfg["state_dir"],
    }


def ensure_dirs():
    for d in _dirs().values():
        d.mkdir(parents=True, exist_ok=True)


def _today_str():
    return date.today().isoformat()


@contextmanager
def _file_lock(lock_path: Path):
    lock_path.parent.mkdir(parents=True, exist_ok=True)
    with open(lock_path, "w") as lf:
        fcntl.flock(lf, fcntl.LOCK_EX)
        try:
            yield
        finally:
            fcntl.flock(lf, fcntl.LOCK_UN)


def _atomic_write_json(path: Path, data):
    path.parent.mkdir(parents=True, exist_ok=True)
    tmp = path.with_suffix(path.suffix + ".tmp")
    with open(tmp, "w") as f:
        json.dump(data, f, indent=2, default=str)
    os.replace(tmp, path)


def _read_json(path: Path, default):
    if not path.exists():
        return default
    try:
        with open(path, "r") as f:
            return json.load(f)
    except (json.JSONDecodeError, OSError) as e:
        log.warning(f"Failed to read {path}: {e}. Returning default.")
        return default


def _append_jsonl(path: Path, record: dict):
    path.parent.mkdir(parents=True, exist_ok=True)
    with open(path, "a") as f:
        f.write(json.dumps(record, default=str) + "\n")


# ---------------------------------------------------------------------------
# Premarket data
# ---------------------------------------------------------------------------

def append_premarket_snapshot(symbol: str, record: dict):
    path = _dirs()["premarket"] / f"{_today_str()}.jsonl"
    record = {"timestamp": datetime.utcnow().isoformat(), "symbol": symbol, **record}
    _append_jsonl(path, record)


def write_premarket_candidates(candidates: list, filename="candidates_20.json"):
    path = _dirs()["premarket"] / f"{_today_str()}_{filename}"
    _atomic_write_json(path, candidates)


def write_top_stocks(top_list: list):
    """Writes the canonical top_stocks.json referenced by the whole system."""
    path = _dirs()["state"] / "top_stocks.json"
    payload = {"updated_at": datetime.utcnow().isoformat(), "stocks": top_list}
    _atomic_write_json(path, payload)


def read_top_stocks() -> list:
    path = _dirs()["state"] / "top_stocks.json"
    data = _read_json(path, {"stocks": []})
    return data.get("stocks", [])


# ---------------------------------------------------------------------------
# Intraday data
# ---------------------------------------------------------------------------

def append_intraday_snapshot(symbol: str, record: dict):
    path = _dirs()["intraday"] / f"{_today_str()}.jsonl"
    record = {"timestamp": datetime.utcnow().isoformat(), "symbol": symbol, **record}
    _append_jsonl(path, record)


def write_intraday_json(symbol: str, record: dict):
    """Mirrors the latest intraday state for a symbol into intraday.json
    (a single mutable snapshot file, distinct from the historical .jsonl log),
    matching the project's requested intraday.json state file pattern."""
    path = _dirs()["state"] / "intraday.json"
    with _file_lock(_dirs()["state"] / "intraday.json.lock"):
        all_data = _read_json(path, {})
        all_data[symbol] = {"updated_at": datetime.utcnow().isoformat(), **record}
        _atomic_write_json(path, all_data)


# ---------------------------------------------------------------------------
# Trades
# ---------------------------------------------------------------------------

def append_trade_record(trade: dict):
    path = _dirs()["trades"] / f"{_today_str()}_trades.jsonl"
    _append_jsonl(path, trade)


def write_trades_summary(trades: list):
    path = _dirs()["trades"] / f"{_today_str()}_summary.json"
    _atomic_write_json(path, trades)


def load_today_trades() -> list:
    """
    Reads every trade record appended today via append_trade_record() —
    the complete, per-trade history for the day.

    [BUGFIX 2026-08-17] This is deliberately separate from
    PositionManager.positions, which is keyed by symbol and therefore
    only retains the MOST RECENT trade for any symbol traded more than
    once in a day. monitor.py's end-of-day summary previously built
    itself from positions.values() and silently undercounted days with
    repeated same-symbol round trips (verified: a session with 19 actual
    closed trades logged a summary of only 10 — one per unique symbol).
    Use this function instead for anything that needs a true count of
    the day's trades or an accurate total P/L.
    """
    path = _dirs()["trades"] / f"{_today_str()}_trades.jsonl"
    if not path.exists():
        return []
    trades = []
    with open(path, "r") as f:
        for line in f:
            line = line.strip()
            if not line:
                continue
            try:
                trades.append(json.loads(line))
            except json.JSONDecodeError:
                log.warning(f"Skipping malformed trade record line in {path}")
                continue
    return trades


# ---------------------------------------------------------------------------
# State: positions, watchlist, cooldowns, session status
# ---------------------------------------------------------------------------

def _state_path(name: str) -> Path:
    return _dirs()["state"] / name


def load_positions() -> dict:
    raw = _read_json(_state_path("positions.json"), {})
    today = _today_str()
    filtered = {}
    for sym, p in raw.items():
        if p.get("status") in ("open", "closing"):
            filtered[sym] = p  # never drop live positions regardless of date
        elif p.get("entry_time", "").startswith(today):
            filtered[sym] = p  # closed today -- keep for has_closed_position_today()
        # else: closed on a prior day -- drop, it's stale
    return filtered


def save_positions(positions: dict):
    with _file_lock(_state_path("positions.json.lock")):
        _atomic_write_json(_state_path("positions.json"), positions)


def load_watchlist() -> dict:
    return _read_json(_state_path("watchlist.json"), {"premarket_20": [], "final_10": []})


def save_watchlist(watchlist: dict):
    _atomic_write_json(_state_path("watchlist.json"), watchlist)


def load_cooldowns() -> dict:
    return _read_json(_state_path("cooldowns.json"), {})


def save_cooldowns(cooldowns: dict):
    _atomic_write_json(_state_path("cooldowns.json"), cooldowns)


# [FEATURE 2026-08-17] Persisted hysteresis state for intraday_health.py's
# continuous premarket_20 health evaluation -- consecutive-good/bad read
# counters and each symbol's current confirmed HEALTHY/WATCH/STALE/
# UNHEALTHY/REMOVED state live here so they survive individual poll
# cycles (and, incidentally, a mid-day restart).
def load_health_state() -> dict:
    return _read_json(_state_path("intraday_health.json"), {})


def save_health_state(health_state: dict):
    with _file_lock(_state_path("intraday_health.json.lock")):
        _atomic_write_json(_state_path("intraday_health.json"), health_state)


def load_session_state() -> dict:
    return _read_json(_state_path("session_state.json"), {})


def save_session_state(state: dict):
    with _file_lock(_state_path("session_state.json.lock")):
        _atomic_write_json(_state_path("session_state.json"), state)


def acquire_singleton_lock() -> "SingletonLock":
    """Ensures only one monitor.py instance runs at a time — prevents the
    duplicate-process-racing-the-same-account failure class."""
    return SingletonLock(_state_path("monitor.pid.lock"))


class SingletonLock:
    def __init__(self, path: Path):
        self.path = path
        self._fh = None

    def __enter__(self):
        self.path.parent.mkdir(parents=True, exist_ok=True)
        self._fh = open(self.path, "w")
        try:
            fcntl.flock(self._fh, fcntl.LOCK_EX | fcntl.LOCK_NB)
        except BlockingIOError:
            raise RuntimeError(
                "Another monitor.py instance already holds the singleton lock. "
                "Refusing to start a second instance against the same account."
            )
        self._fh.write(str(os.getpid()))
        self._fh.flush()
        return self

    def __exit__(self, exc_type, exc_val, exc_tb):
        if self._fh:
            fcntl.flock(self._fh, fcntl.LOCK_UN)
            self._fh.close()
