"""
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"],
        "fast_engine": BASE_DIR / cfg["fast_engine_dir"],
        "state": BASE_DIR / cfg["state_dir"],
        "breakout_scanner": BASE_DIR / cfg["breakout_scanner_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_breakout_scanner_candidates(candidates: list):
    """[FEATURE 2026-09-16] breakout_scanner.py's own ranked candidate
    list -- deliberately a SEPARATE file from write_premarket_candidates
    above (different directory, different name) so it can never be
    confused with or overwrite sip_bot's own live candidate output.
    Date-stamped filename, same convention as write_premarket_candidates,
    so a day's breakout_scanner run and sip_bot run can be diffed later."""
    path = _dirs()["breakout_scanner"] / f"{_today_str()}_scanner.json"
    _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)


# ---------------------------------------------------------------------------
# Fast prediction engine decisions
# ---------------------------------------------------------------------------

def append_fast_engine_decision(record: dict):
    """
    [FEATURE 2026-09-11, extended 2026-09-11 v3] One JSON line per
    ranked candidate, every poll cycle, for EVERY health-eligible
    premarket_20 symbol -- not just the top-5 shortlist that gets a real
    fast_entry_gate decision. `shortlisted`/`confidence_rank` distinguish
    "ranked but not evaluated this cycle" (`shortlisted: False`,
    `state`/`should_enter`: None) from "made the top-5, got a real
    BUY/WAIT/REJECT" (`shortlisted: True`). Per explicit instruction
    ("log all bot decisions so we can review and adjust what is not
    working", later extended to "make sure we have all the indicator
    values logged during the life cycle of the symbol"): the whole point
    is being able to load this file into pandas afterward and see
    exactly what the engine saw for every symbol at every poll, whether
    or not it ever became a real trade (which append_trade_record above
    already covers) or even made the shortlist. Each record also carries
    a full raw-indicator snapshot (fast_pipeline.snapshot_dict() --
    every stream_features.StreamFeatures/FastPredictionReading field,
    not just the classified labels) alongside the classified summary
    fields kept here for quick filtering without unpacking the nested
    snapshot. Same per-day .jsonl pattern as premarket/intraday
    snapshots -- pandas-readable, no database dependency. Volume note:
    at ~20-30 eligible symbols x poll_interval_seconds_intraday=5s, this
    is meaningfully larger than the old shortlisted-only version (was
    ~5 records/cycle, now ~20-30/cycle) -- expect a materially bigger
    file per trading day.
    """
    path = _dirs()["fast_engine"] / f"{_today_str()}.jsonl"
    record = {"timestamp": datetime.utcnow().isoformat(), **record}
    _append_jsonl(path, record)


def append_momentum_fade_snapshot(record: dict):
    """
    [FEATURE 2026-09-11] One JSON line per open-position poll where
    position_manager._check_momentum_fade_exit() successfully computed
    features -- added the same day that check was introduced, after
    discovering there was no way to validate its thresholds
    (velocity_stall_threshold_pct_per_sec / pressure_threshold /
    min_ticks_sub, see config.json's stop.momentum_fade_exit) against
    real data: a full replay of 2026-09-11's 16 trades against the new
    check turned out to be impossible, because append_fast_engine_
    decision() above only dumps CANDIDATE-scan snapshots, which stop the
    moment a symbol becomes an open position (fast_pipeline.py has no
    reason to keep ranking something already held) -- so no tick-level
    velocity_sub/pressure.score history exists for any open-position hold
    window from that day. This closes that gap going forward: every poll
    where the check actually evaluates a reading (not the early-return
    guard cases -- disabled, missing bars_sub/quote, insufficient_data --
    which aren't informative for tuning) gets one row here, whether or
    not it fired, so the thresholds can be tuned against real fills
    instead of Yahoo-1-minute-bar approximations next time. Separate file
    from fast_engine's per-day dump -- different schema/purpose (a single
    open position's fade check vs. ranking all eligible candidates) even
    though it lives in the same directory.
    """
    path = _dirs()["fast_engine"] / f"momentum_fade_{_today_str()}.jsonl"
    record = {"timestamp": datetime.utcnow().isoformat(), **record}
    _append_jsonl(path, record)


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()
