"""
monitor.py

Main orchestrator. Runs the full daily lifecycle:

    START
     -> premarket_scanner.py (09:29, one full-universe scan) -> top
        candidates (premarket_20)
     -> start/maintain SIP stream for premarket_20
     -> opening confirmation -> entry (evaluated across premarket_20)
     -> continuous intraday health scoring of premarket_20 (every
        intraday_health.eval_interval_seconds)
     -> full intraday rescan + premarket_20 rotation (every
        intraday_health.full_rescan_interval_minutes, until the
        no-new-entries cutoff)
     -> position management + trailing stops
     -> on exit -> top_stocks.py -> find replacement (from premarket_20)
        -> confirm -> enter
     -> end-of-day liquidation
     -> shutdown

[UPDATED 2026-09-01] Per explicit request, the single scan moved from
09:00 to schedule.premarket_scan_time (09:29) and the 09:00-09:25
development-monitoring window plus the 09:25 final-scoring review were
eliminated entirely -- there's no time left for either between a 09:29
scan and the 09:30 open. final_10 is now populated directly from
premarket_20's top candidates.total_score ranking right after the scan
(see _run_premarket_scan()), with no development-trend adjustment
applied (there's no multi-snapshot history to compute one from anymore).

[FEATURE 2026-08-17] premarket_20 is now the primary intraday candidate
universe (previously this was final_10, which restricted the bot to
only 10 of the original 20 candidates for the entire day). final_10 is
still written for the dashboard's own tab and for backward
compatibility, but no longer gates entries or replacement search.

Run directly:
    python monitor.py

Designed to be started once per day via cron shortly before 09:00 ET,
and to exit cleanly after end-of-day liquidation. A singleton file
lock prevents two instances from ever running against the same account
simultaneously.
"""

import sys
import time
import signal
import atexit
from datetime import datetime, timezone, timedelta

from config_loader import get_config, get_env, validate_config
from logger_setup import get_logger
import market_time
import data_store
import premarket_scanner
import top_stocks
import intraday_health
from entry_engine import evaluate_entry
from prediction_pipeline import evaluate_entry_pipeline
from position_manager import PositionManager
from stream import StreamManager
from alpaca_client import get_client

log = get_logger("monitor")


class _EntryAttemptTracker:
    """
    [FEATURE 2026-09-01] Fixes the shortlist-starvation bug found in the
    2026-09-01 session review -- see config.json's
    intraday_health._note_confirmation_failure_cooldown for the full
    real-world case (F monopolizing a shortlist slot for ~24 minutes
    while CRML, health-eligible and competitively scored, never got a
    single evaluate_entry() call before rotating back out).

    _scan_for_entries() ranks premarket_20 by health_score and only
    evaluates the top entry_shortlist_size symbols each poll cycle. A
    symbol that ranks highly but keeps failing entry_engine's
    confirmation on a condition that doesn't change bar-to-bar (e.g. a
    hard momentum disqualifier) can occupy one of those slots forever,
    since nothing here ever demoted it. This class tracks consecutive
    not-confirmed results per symbol and temporarily benches one after
    max_failures in a row, freeing its shortlist slot for the next-best
    candidate for cooldown_seconds.

    Deliberately does NOT change entry_engine's actual confirmation
    criteria in any way -- a benched symbol that comes off the bench is
    evaluated against the exact same rules as before. This only changes
    WHICH candidates get a turn, never what it takes to pass.

    Pure in-memory, per-session state (like self.premarket_20 itself) --
    intentionally not persisted, since a fresh session should start
    every symbol unbenched.
    """

    def __init__(self, max_failures: int, cooldown_seconds: float):
        self.max_failures = max_failures
        self.cooldown_seconds = cooldown_seconds
        self._consecutive_failures = {}
        self._benched_until = {}

    def is_benched(self, symbol: str, now: datetime) -> bool:
        until = self._benched_until.get(symbol)
        return until is not None and now < until

    def record_result(self, symbol: str, confirmed: bool, now: datetime):
        if confirmed:
            self._consecutive_failures.pop(symbol, None)
            self._benched_until.pop(symbol, None)
            return

        count = self._consecutive_failures.get(symbol, 0) + 1
        if count >= self.max_failures:
            self._benched_until[symbol] = now + timedelta(seconds=self.cooldown_seconds)
            self._consecutive_failures[symbol] = 0
            log.info(f"[ENTRY] {symbol} benched from the entry shortlist for "
                     f"{self.cooldown_seconds:.0f}s after {count} consecutive "
                     f"not-confirmed attempts -- giving other candidates a turn")
        else:
            self._consecutive_failures[symbol] = count


class SessionOrchestrator:
    def __init__(self):
        self.cfg = get_config()
        validate_config(self.cfg)
        get_env().validate()

        self.simulation = self.cfg["mode"]["execution_mode"] == "simulation"
        self.client = get_client()
        self.position_mgr = PositionManager(simulation=self.simulation)
        self.stream = StreamManager()
        self._entry_tracker = _EntryAttemptTracker(
            max_failures=self.cfg["intraday_health"]["max_consecutive_confirmation_failures"],
            cooldown_seconds=self.cfg["intraday_health"]["confirmation_failure_cooldown_seconds"],
        )

        # [BUGFIX 2026-08-31] see intraday_health.prune_stale_health_state()
        # -- same "no today-only filter" gap already fixed for positions.json
        # via PositionManager._prune_stale_closed_positions().
        pruned_health_state = intraday_health.prune_stale_health_state(data_store.load_health_state())
        data_store.save_health_state(pruned_health_state)

        self.premarket_20 = []          # list of scored dicts -- the primary intraday universe
        self.final_10 = []              # informational/dashboard only as of 2026-08-17
        self._prefiltered_universe = []  # cached snapshot-filtered symbol list, reused by 30-min rescans
        self._opening_volume_baseline = {}  # symbol -> baseline volume, kept in sync as premarket_20 changes
        self._last_health_eval = None
        self._last_full_rescan = None
        self._shutdown = False

        # [FEATURE 2026-09-05] Caller-held per-symbol state for
        # prediction_pipeline.evaluate_entry_pipeline() (prior_features/
        # prior_prediction/confirmation_state) -- only populated/consumed
        # when entry.decision_engine == "prediction_pipeline"; harmless
        # empty dict otherwise. See config.json's entry._note_decision_engine.
        self._pipeline_state = {}

        signal.signal(signal.SIGTERM, self._handle_signal)
        signal.signal(signal.SIGINT, self._handle_signal)
        atexit.register(self._cleanup)

    # ------------------------------------------------------------------
    def _handle_signal(self, signum, frame):
        log.info(f"[SHUTDOWN] Received signal {signum}, shutting down gracefully")
        self._shutdown = True

    def _cleanup(self):
        try:
            self.stream.stop()
        except Exception:
            pass
        log.info("[SHUTDOWN] Cleanup complete")

    # ------------------------------------------------------------------
    def run(self):
        log.info(f"[START] monitor.py starting | mode={self.cfg['mode']['execution_mode']}")

        if not market_time.is_weekday():
            log.info("[START] Not a weekday, exiting")
            return

        self._wait_for_premarket_scan()
        if self._shutdown:
            return

        self._run_premarket_scan()
        self._start_streaming()
        self._wait_for_market_open()
        self._main_trading_loop()
        self._end_of_day_liquidation()

        log.info("[SHUTDOWN] Session complete")

    # ------------------------------------------------------------------
    def _wait_for_premarket_scan(self):
        while not market_time.is_premarket_scan_time() and not self._shutdown:
            time.sleep(5)

    def _run_premarket_scan(self):
        # [FEATURE 2026-08-17] Compute the snapshot-filtered universe
        # ourselves and hand it into scan() so we can cache it for reuse
        # by every 30-minute intraday rescan for the rest of the day --
        # avoids re-pulling the full tradable-asset list + a fresh bulk
        # snapshot call every single rescan.
        universe = premarket_scanner.get_universe_symbols(self.client)
        log.info(f"[PREMARKET] Universe size after asset filtering: {len(universe)}")
        self._prefiltered_universe = premarket_scanner.prefilter_by_snapshot(self.client, universe)
        log.info(f"[PREMARKET] Prefiltered to {len(self._prefiltered_universe)} symbols in price/volume band")

        self.premarket_20 = premarket_scanner.scan(
            self.cfg["candidates"]["premarket_candidate_count"],
            prefiltered=self._prefiltered_universe,
        )

        # [UPDATED 2026-09-01] final_10 used to come from a 09:25 final-
        # scoring pass that applied a development-trend adjustment built
        # from 09:00-09:25 snapshot history. That window no longer exists
        # (scan now runs once at schedule.premarket_scan_time, right
        # before the open) -- so final_10 is just premarket_20's top
        # final_candidate_count by the scan's own total_score, computed
        # right here instead of waiting for a step that no longer runs.
        # premarket_20 is already sorted descending by total_score
        # (premarket_scanner.scan() does this), so this is a plain slice.
        self.final_10 = self.premarket_20[: self.cfg["candidates"]["final_candidate_count"]]

        data_store.save_watchlist({
            "premarket_20": [{"symbol": r["symbol"], "score": r["total_score"]} for r in self.premarket_20],
            "final_10": [{"symbol": r["symbol"], "score": r["total_score"]} for r in self.final_10],
        })
        data_store.write_top_stocks(self.premarket_20)
        self._sync_opening_volume_baseline(self.premarket_20)

    def _sync_opening_volume_baseline(self, candidates: list):
        """Keeps self._opening_volume_baseline current as premarket_20
        gains symbols (initial scan, slot-replacement additions, and
        full-rescan rotations) -- entry_engine's opening-volume-expansion
        check needs a real baseline per symbol, not the accidental "no
        baseline -> ratio always looks huge" behavior of a missing key."""
        for r in candidates:
            metrics = r.get("metrics", {})
            self._opening_volume_baseline[r["symbol"]] = max(metrics.get("premarket_volume", 1), 1)

    # ------------------------------------------------------------------
    def _start_streaming(self):
        # [FEATURE 2026-08-17] Stream the full premarket_20, not just
        # final_10 -- entry scanning, replacement search, and continuous
        # health evaluation all need live bars for all 20 candidates now.
        symbols = [r["symbol"] for r in self.premarket_20]
        if not symbols:
            log.warning("[OPEN] No candidates to stream; skipping stream start")
            return
        self.stream.start(symbols)
        if self.stream.is_healthy():
            log.info("[OPEN] SIP stream active")
        else:
            log.error("[OPEN] SIP stream failed to become healthy; will rely on REST fallback")

    def _wait_for_market_open(self):
        while not market_time.is_past_market_open() and not self._shutdown:
            time.sleep(2)

    # ------------------------------------------------------------------
    def _main_trading_loop(self):
        interval = self.cfg["schedule"]["poll_interval_seconds_intraday"]
        health_cfg = self.cfg["intraday_health"]
        now = datetime.now(timezone.utc)
        self._last_health_eval = now
        self._last_full_rescan = now

        while not self._shutdown and not market_time.is_force_liquidate_time():
            # 0. Advance any exit already in flight (poll its order status
            #    and finalize on fill) before doing anything else.
            #    [BUGFIX 2026-08-18] -- this must run every cycle so a
            #    "closing" position never sits unresolved while stale.
            self.position_mgr.poll_pending_exits()

            # 1. Manage existing open positions (stops)
            self._update_open_positions()

            # 2. Look for new entries among premarket_20 symbols with free slots
            if self.position_mgr.has_available_slot() and not market_time.is_new_entries_cutoff():
                self._scan_for_entries()

            # 3. Reconcile broker state periodically (cheap safety net)
            self.position_mgr.reconcile_with_broker()

            # 4. Continuous intraday health evaluation of premarket_20
            #    [FEATURE 2026-08-17]
            now = datetime.now(timezone.utc)
            if health_cfg.get("enabled", True):
                elapsed = (now - self._last_health_eval).total_seconds()
                if elapsed >= health_cfg["eval_interval_seconds"]:
                    self._update_intraday_health()
                    self._last_health_eval = now

            # 5. Full intraday rescan + premarket_20 rotation, every
            #    full_rescan_interval_minutes, but only while new entries
            #    are still allowed -- rebuilding the candidate pool after
            #    the cutoff would have nothing to act on.
            if not market_time.is_new_entries_cutoff():
                elapsed = (now - self._last_full_rescan).total_seconds()
                if elapsed >= health_cfg["full_rescan_interval_minutes"] * 60:
                    self._run_intraday_full_rescan()
                    self._last_full_rescan = now

            time.sleep(interval)

    def _update_open_positions(self):
        for symbol in self.position_mgr.get_open_symbols():
            bars = self.stream.get_bars(symbol)
            if not bars:
                continue
            current_price = bars[-1]["c"]
            exit_signal = self.position_mgr.update_position(symbol, current_price, bars)
            if exit_signal:
                self.position_mgr.exit_position(symbol, current_price, exit_signal)
                top_stocks.register_cooldown(symbol)
                self._handle_slot_freed(symbol)

    def _handle_slot_freed(self, freed_symbol: str):
        log.info(f"[SLOT] 1 position slot available")
        open_symbols = self.position_mgr.get_open_symbols()
        # find_replacement() now defaults to searching premarket_20 (see
        # top_stocks.py) and ranks by continuous health score.
        candidate = top_stocks.find_replacement(exclude_symbols=open_symbols)
        if not candidate:
            return
        # Ensure the stream is subscribed to the new candidate before we
        # try to read live bars for it.
        self.stream.add_symbols([candidate["symbol"]])

        # [BUGFIX 2026-08-19] This used to be "insert only if absent," which
        # meant that if `candidate["symbol"]` already had a (possibly stale)
        # entry sitting in self.premarket_20 -- e.g. a symbol that was in
        # the pool earlier, dropped health, and is now being freshly
        # re-selected as a replacement by top_stocks.find_replacement() --
        # the fresh pm_high/resistance and score data top_stocks just
        # computed would be silently discarded and the old dict left in
        # place. Upsert instead: replace any existing entry for this
        # symbol with the fresh one so resistance/pm_high always reflects
        # what top_stocks.find_replacement() just calculated.
        self.premarket_20 = [r for r in self.premarket_20 if r["symbol"] != candidate["symbol"]]
        self.premarket_20.append(candidate)
        self._sync_opening_volume_baseline([candidate])

    def _scan_for_entries(self):
        # [FEATURE 2026-08-17] Ranked shortlist, not first-come-first-served
        # list order. Every poll cycle, all not-yet-open, health-eligible
        # premarket_20 symbols are ranked by their current health_score
        # (tie-broken by the premarket-style total_score), and only the
        # top intraday_health.entry_shortlist_size (default 5) are actually
        # evaluated through entry_engine.evaluate_entry() this cycle. This
        # is intentionally independent from trading.max_positions (how
        # many can be OPEN at once) -- the shortlist is "who's currently
        # worth trying," the max_positions cap is "how many slots exist."
        # A symbol outside today's shortlist simply isn't attempted this
        # cycle; it's re-ranked fresh on the next one as health scores
        # change, so nothing is permanently excluded the way REMOVED is.
        #
        # IMPORTANT: this only affects whether a NEW entry can be taken.
        # It has no effect on any already-open position, which is managed
        # exclusively by position_manager.py / risk_manager.py.
        health_state = data_store.load_health_state()
        shortlist_size = self.cfg["intraday_health"]["entry_shortlist_size"]
        now = datetime.now(timezone.utc)

        eligible = []
        for r in self.premarket_20:
            symbol = r["symbol"]
            if self.position_mgr.is_symbol_open(symbol):
                continue

            entry = health_state.get(symbol, {})
            confirmed_state = entry.get("state", "WATCH")
            if not intraday_health.is_eligible_for_entry(confirmed_state):
                log.debug(f"[ENTRY] Skipping {symbol}: health={confirmed_state}")
                continue

            if self.stream.is_symbol_stale(symbol):
                log.debug(f"[ENTRY] Skipping {symbol}: stale stream data")
                continue

            # [FEATURE 2026-09-01] See _EntryAttemptTracker's docstring --
            # a symbol repeatedly failing confirmation is temporarily
            # excluded from the shortlist so it can't monopolize a slot
            # forever and starve lower-ranked, potentially viable
            # candidates of ever getting a turn.
            if self._entry_tracker.is_benched(symbol, now):
                log.debug(f"[ENTRY] Skipping {symbol}: benched after repeated "
                          f"confirmation failures")
                continue

            health_score = entry.get("health_score", 0.0)
            eligible.append((health_score, r.get("total_score", 0.0), r))

        eligible.sort(key=lambda t: (t[0], t[1]), reverse=True)
        shortlist = eligible[:shortlist_size]

        if shortlist:
            log.debug(f"[ENTRY] Shortlist this cycle: " +
                      ", ".join(f"{r['symbol']}(health={hs:.1f})" for hs, _, r in shortlist))

        avg_vol_baseline = self.cfg["universe"]["min_avg_daily_volume"]

        for health_score, total_score, r in shortlist:
            symbol = r["symbol"]
            if not self.position_mgr.has_available_slot():
                return

            bars = self.stream.get_bars(symbol)
            quote = self.stream.get_quote(symbol)
            baseline = self._opening_volume_baseline.get(symbol, 1)

            if not bars or len(bars) < 3:
                log.debug(f"[ENTRY] Skipping {symbol}: insufficient live bars "
                          f"for a fresh health check")
                continue

            # [BUGFIX 2026-08-18] The eligibility check above used
            # state/intraday_health.json -- fine for cheaply ranking the
            # shortlist, but only as current as the last periodic
            # _update_intraday_health() cycle (up to eval_interval_seconds
            # old, more if a cycle lagged). Confirmed live: two RCAT
            # entries fired on a cached WATCH reading that was ~10
            # minutes stale by the time of entry. Recompute fresh, right
            # now, against the exact same bars this entry is about to be
            # decided on, and let entry_engine gate on THAT -- not the
            # cached state -- as one of its own required checks.
            fresh_reading = intraday_health.compute_health(symbol, bars, avg_vol_baseline)
            cached_state = health_state.get(symbol, {}).get("state", "?")
            if fresh_reading.raw_state != cached_state:
                log.debug(f"[ENTRY] {symbol}: cached health={cached_state} but fresh "
                          f"read={fresh_reading.raw_state} (score={fresh_reading.health_score:.1f}) "
                          f"-- using the fresh read for the entry decision")

            # [FEATURE 2026-08-27] Score-based re-entry gate -- see
            # entry_engine.evaluate_entry()'s is_reentry docstring.
            is_reentry = self.position_mgr.has_closed_position_today(symbol)

            if self.cfg["entry"].get("decision_engine") == "prediction_pipeline":
                # [FEATURE 2026-09-05] See config.json's entry._note_decision_engine.
                # Only the entry TRIGGER changes -- fresh_reading/is_reentry above,
                # and everything after should_enter below (sizing, stops, exits),
                # are identical regardless of which path decided should_enter.
                bars_sub = self.stream.get_bars_sub(symbol)
                sub_bucket_seconds = self.cfg["streaming"].get("sub_minute_bucket_seconds", 30)
                pipeline_state = self._pipeline_state.setdefault(symbol, {})
                decision = evaluate_entry_pipeline(
                    symbol, bars, bars_sub, quote, baseline,
                    session_elapsed_minutes=market_time.minutes_since_open(),
                    state=pipeline_state, sub_bucket_seconds=sub_bucket_seconds)
                # [BUGFIX 2026-09-08] evaluate_entry_pipeline() (unlike
                # entry_engine.evaluate_entry(), which has always logged
                # this internally -- see its own 2026-08-28 bugfix note)
                # never had its reasons_for/reasons_against written
                # anywhere. Confirmed on 2026-09-08's live paper session:
                # 0 trades, 260 benched-after-5-failures events, and NO
                # way to see why any of them failed -- entry_score.py
                # computes a specific reason (e.g. "prediction_score 42.3
                # below minimum 60") every single call, it was just being
                # thrown away here. Mirrors entry_engine.py's own
                # logging.log_rejections gate exactly, so both decision
                # engines behave the same way under that flag.
                if decision.should_enter:
                    log.info(f"[PIPELINE] {symbol} confirmed ({decision.confirmation_score:.1f}): "
                             f"{'; '.join(decision.reasons_for) if decision.reasons_for else 'no specific reason recorded'}")
                else:
                    log_rejections = self.cfg.get("logging", {}).get("log_rejections", False)
                    message = (f"[PIPELINE] {symbol} not confirmed ({decision.confirmation_score:.1f}): "
                               f"{'; '.join(decision.reasons_against) if decision.reasons_against else 'no specific reason recorded'}")
                    if log_rejections:
                        log.info(message)
                    else:
                        log.debug(message)
            else:
                decision = evaluate_entry(symbol, r, bars, quote, baseline,
                                           health_reading=fresh_reading, is_reentry=is_reentry)
            self._entry_tracker.record_result(symbol, decision.should_enter, datetime.now(timezone.utc))
            if decision.should_enter:
                entry_price = bars[-1]["c"]
                account = self.client.get_account() if not self.simulation else None
                sim_equity = self.cfg["trading"]["simulated_equity_default"]
                equity = float(account.equity) if account else sim_equity
                reason = "; ".join(decision.reasons_for)
                self.position_mgr.enter_position(symbol, entry_price, equity, bars, reason)

    # ------------------------------------------------------------------
    def _update_intraday_health(self):
        """
        [FEATURE 2026-08-17] Evaluates every symbol currently in
        premarket_20 through intraday_health.evaluate_symbol(), using
        the SIP stream's already-buffered bars (no extra API calls).
        Persists the hysteresis-confirmed state to
        state/intraday_health.json, which _scan_for_entries() and
        top_stocks.find_replacement() both read.

        This function NEVER touches position_manager.py or triggers an
        exit -- an existing open position is left alone regardless of
        what health state its symbol reports here.
        """
        avg_vol_baseline = self.cfg["universe"]["min_avg_daily_volume"]
        health_state = data_store.load_health_state()

        for r in self.premarket_20:
            symbol = r["symbol"]
            bars = self.stream.get_bars(symbol)
            if not bars or len(bars) < 3:
                continue
            reading, new_entry = intraday_health.evaluate_symbol(
                symbol, bars, avg_vol_baseline, health_state)
            health_state[symbol] = new_entry

        data_store.save_health_state(health_state)

    def _run_intraday_full_rescan(self):
        """
        [FEATURE 2026-08-17] Every intraday_health.full_rescan_interval_
        minutes (default 30), re-scores the ENTIRE cached prefiltered
        universe (not just the current premarket_20) against live,
        current bars -- using the same scoring function and the same
        health-eligibility filter as the rest of the intraday system --
        and rotates premarket_20 to the resulting top N. This is what
        lets the bot discover a genuinely new candidate that wasn't
        part of the original 09:00 list, not just rotate among the
        original 20 all day.

        Reuses self._prefiltered_universe (captured once at 09:00)
        rather than re-pulling the full tradable-asset list and a fresh
        bulk snapshot every 30 minutes.

        Any symbol with an open position is always kept in the pool
        regardless of whether the fresh rescan re-selects it -- dropping
        it here would only affect future entry/replacement eligibility
        bookkeeping, never the open position itself, but keeping it
        visible avoids losing track of it in logs/dashboard.
        """
        if not self._prefiltered_universe:
            log.warning("[INTRADAY-RESCAN] No cached prefiltered universe available; skipping")
            return

        log.info(f"[INTRADAY-RESCAN] Starting full rescan of "
                 f"{len(self._prefiltered_universe)} prefiltered symbols")
        cfg = self.cfg["intraday_health"]
        avg_vol_baseline = self.cfg["universe"]["min_avg_daily_volume"]
        health_state = data_store.load_health_state()

        rescored = premarket_scanner.scan(
            candidate_count=len(self._prefiltered_universe),
            prefiltered=self._prefiltered_universe,
            lookback_hours=cfg["full_rescan_lookback_hours"],
        )

        eligible = []
        for r in rescored:
            symbol = r["symbol"]
            bars = self.stream.get_bars(symbol) if symbol in [p["symbol"] for p in self.premarket_20] else None
            if not bars or len(bars) < 3:
                # Not currently streamed (a genuinely new candidate) --
                # evaluate health directly off the rescan's own bars via
                # a fresh short pull is unnecessary here since scan()
                # already required >=3 bars to produce a scored result;
                # fall back to a neutral WATCH read so a brand-new name
                # isn't unfairly excluded on its very first appearance.
                r["health_score"] = r["total_score"]
                r["health_state"] = "WATCH"
                eligible.append(r)
                continue
            reading, new_entry = intraday_health.evaluate_symbol(
                symbol, bars, avg_vol_baseline, health_state)
            health_state[symbol] = new_entry
            r["health_score"] = reading.health_score
            r["health_state"] = reading.confirmed_state
            if intraday_health.is_eligible_for_entry(reading.confirmed_state):
                eligible.append(r)

        data_store.save_health_state(health_state)

        if not eligible:
            log.warning("[INTRADAY-RESCAN] No health-eligible candidates found; "
                        "keeping the existing premarket_20 unchanged")
            return

        eligible.sort(key=lambda r: (r["health_score"], r["total_score"]), reverse=True)
        new_pool = eligible[: cfg["full_rescan_pool_size"]]

        # Never drop a symbol with an open position from the tracked pool.
        #
        # [BUGFIX 2026-08-19] This used to re-insert the OLD dict from
        # self.premarket_20 for any open-position symbol that fell outside
        # the top full_rescan_pool_size cut. That dict's pm_high/resistance
        # was whatever was computed the last time this symbol made the cut
        # -- possibly hours earlier -- and reusing the object (rather than
        # the fresh one this very rescan just computed in `rescored`) meant
        # resistance/pm_high could freeze indefinitely: since self.premarket_20
        # is reassigned to new_pool at the end of this function, the same
        # stale object gets carried forward into the NEXT rescan too, so a
        # symbol that keeps having an open position or keeps narrowly
        # missing the health cut never gets its resistance level refreshed
        # again for the rest of the session. Confirmed in production logs:
        # CLSK entered 5 times between 11:27 and 13:24 citing the identical
        # "breakout of resistance $11.70" on every single entry, despite
        # trading as high as $12.20 in between -- the resistance value was
        # never actually being recalculated after the first time CLSK
        # dropped out of the top-20 health cut with a position open.
        #
        # Fix: `rescored` already contains a FRESH score (fresh pm_high
        # included) for every symbol in the prefiltered universe, since
        # premarket_scanner.scan() just rescored all of them this pass --
        # it's only `eligible[:pool_size]` that truncates. So for any
        # open-position symbol that got cut, look up its fresh entry from
        # `rescored` first and use that. Only fall back to the old
        # self.premarket_20 entry in the (very unlikely) case the symbol
        # is missing from `rescored` entirely, e.g. halted/delisted intraday.
        # `rescored` entries already carry fresh health_score/health_state
        # from the loop above (set for every symbol, eligible or not), so
        # the fresh dict pulled from here needs no further patching --
        # it's a complete, current record, just one that didn't make the
        # top-N cut on health/total_score.
        rescored_by_symbol = {r["symbol"]: r for r in rescored}
        open_symbols = set(self.position_mgr.get_open_symbols())
        pool_symbols = {r["symbol"] for r in new_pool}
        for r in self.premarket_20:
            symbol = r["symbol"]
            if symbol in open_symbols and symbol not in pool_symbols:
                fresh = rescored_by_symbol.get(symbol)
                if fresh is not None:
                    new_pool.append(fresh)
                    pool_symbols.add(symbol)
                    log.debug(f"[INTRADAY-RESCAN] {symbol} kept in pool (open position) "
                              f"with refreshed pm_high=${fresh.get('pm_high', 0):.2f}")
                else:
                    log.warning(f"[INTRADAY-RESCAN] {symbol} kept in pool (open position) but "
                                f"missing from this rescan's universe -- reusing stale "
                                f"pm_high=${r.get('pm_high', 0):.2f}; resistance will not "
                                f"reflect current price action until this symbol reappears "
                                f"in the prefiltered universe")
                    new_pool.append(r)
                    pool_symbols.add(symbol)

        self.premarket_20 = new_pool
        self._sync_opening_volume_baseline(new_pool)
        self.stream.add_symbols([r["symbol"] for r in new_pool])

        # [FEATURE 2026-08-27] final_10 (sized by candidates.final_candidate_count,
        # currently 15 despite the legacy key name) used to be written ONCE at
        # 09:25 ("informational only -- premarket_20 remains the active
        # intraday universe") and never touched again -- meaning the
        # dashboard's "FINAL 10" panel showed a frozen premarket snapshot
        # all day while "PREMARKET 20" was the only genuinely live one.
        # Per explicit request ("I want final 15 to be updated as well, so
        # we know the symbols the bot is considering"), this now refreshes
        # final_10 on every rescan too, reusing new_pool -- which is
        # already sorted by (health_score, total_score) above -- so no
        # extra ranking work is needed, just taking its top N. Same key
        # name and same {symbol, score, health} shape premarket_20 already
        # uses, so data.php/dashboard.js need no changes to pick this up.
        final_count = self.cfg.get("candidates", {}).get("final_candidate_count", 15)
        final_slice = new_pool[:final_count]

        watchlist = data_store.load_watchlist()
        watchlist["premarket_20"] = [
            {"symbol": r["symbol"], "score": r["total_score"], "health": r.get("health_state")}
            for r in new_pool
        ]
        watchlist["final_10"] = [
            {"symbol": r["symbol"], "score": r["total_score"], "health": r.get("health_state")}
            for r in final_slice
        ]
        data_store.save_watchlist(watchlist)
        data_store.write_top_stocks(new_pool)

        for r in new_pool:
            log.info(f"[INTRADAY-RESCAN] {r['symbol']} score={r['total_score']} "
                     f"health={r.get('health_state')}")
        log.info(f"[INTRADAY-RESCAN] Rotated premarket_20 -> {len(new_pool)} candidates, "
                 f"final_10 -> top {len(final_slice)}")

    # ------------------------------------------------------------------
    def _end_of_day_liquidation(self):
        """
        [BUGFIX 2026-08-24] Previously this called liquidate_all() ONCE,
        stopped the stream, and logged the summary immediately -- with
        no retry window at all. Any position whose close attempt failed
        on that single pass (e.g. the "position not found" entry-
        settlement race the same day's fix addresses in
        position_manager.py, or a wash-trade rejection needing a poll
        cycle to resolve) was simply left "open"/"closing" in the local
        ledger, EXCLUDED from that day's trades/wins/losses/total_P/L
        summary, and only discovered -- and only actually logged as a
        completed trade -- by the NEXT day's startup reconciliation.
        Confirmed in real logs: 2026-08-21's 5 stuck positions all
        surfaced as EXIT_REASON=END_OF_DAY entries in the 2026-08-24
        log, at market open, a full trading day later than they should
        have resolved.
        Fix: keep calling liquidate_all() (idempotent -- see
        exit_position()) and poll_pending_exits() every poll interval
        for up to eod_liquidation_max_wait_seconds, so anything that
        just needs a few more cycles to settle actually finishes THIS
        session and gets counted THIS day's summary. Only genuinely
        stuck stragglers (which should now be rare) still fall through
        to next-day reconciliation.
        """
        log.info("[EOD] Force-liquidating all open positions")
        self.position_mgr.liquidate_all(reason="END_OF_DAY")

        max_wait = self.cfg.get("schedule", {}).get("eod_liquidation_max_wait_seconds", 120)
        poll_interval = self.cfg["schedule"]["poll_interval_seconds_intraday"]
        deadline = time.time() + max_wait

        while time.time() < deadline:
            self.position_mgr.poll_pending_exits()
            still_open = self.position_mgr.get_open_symbols() + self.position_mgr.get_closing_symbols()
            if not still_open:
                break
            time.sleep(poll_interval)
            # Re-submit closes for anything still genuinely "open"
            # (e.g. a "position not found" backoff window just expired,
            # or a close never got submitted in the first place) --
            # exit_position() is idempotent per symbol, so this is safe
            # to call repeatedly.
            self.position_mgr.liquidate_all(reason="END_OF_DAY")
        else:
            still_open = self.position_mgr.get_open_symbols() + self.position_mgr.get_closing_symbols()
            if still_open:
                log.warning(f"[EOD] {len(still_open)} position(s) still not resolved after "
                            f"{max_wait}s of retries -- will fall through to next session's "
                            f"startup reconciliation: {sorted(still_open)}")

        self.stream.stop()

        # [BUGFIX 2026-08-17] Build the summary from the day's complete
        # trade-by-trade JSONL log (data_store.load_today_trades()),
        # not from self.position_mgr.positions.values(). That dict is
        # keyed by symbol, so a symbol traded more than once in a day
        # (e.g. stopped out and re-entered) had every earlier round trip
        # silently overwritten in memory, undercounting both the trade
        # count and total P/L in this summary. liquidate_all() above has
        # already appended any final EOD exits to the JSONL log by the
        # time we read it here, so this includes the full day.
        trades = data_store.load_today_trades()
        data_store.write_trades_summary(trades)

        total_pl = sum(t.get("current_pl", 0) for t in trades if t.get("status") == "closed")
        wins = sum(1 for t in trades if t.get("status") == "closed" and t.get("current_pl", 0) > 0)
        losses = sum(1 for t in trades if t.get("status") == "closed" and t.get("current_pl", 0) <= 0)
        log.info(f"[EOD] Session summary: trades={len(trades)} wins={wins} losses={losses} "
                 f"total_P/L=${total_pl:.2f}")


def main():
    try:
        with data_store.acquire_singleton_lock():
            data_store.ensure_dirs()
            orchestrator = SessionOrchestrator()
            orchestrator.run()
    except RuntimeError as e:
        log.error(str(e))
        sys.exit(1)


if __name__ == "__main__":
    main()