"""
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
import threading
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 breakout_scanner
import top_stocks
import intraday_health
from fast_pipeline import compute_fast_ranking, evaluate_fast_entry_only, snapshot_dict
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
        # [Phase 5 -- 2026-09-08] symbol -> real prior-day volume, populated
        # once at the 09:29 scan (premarket_scanner.prefilter_by_snapshot's
        # baseline_out) and reused everywhere RVOL is computed for the rest
        # of the day, replacing the flat universe.min_avg_daily_volume
        # constant that was previously used for every symbol regardless of
        # its own normal volume. See premarket_scanner.py's docstring for
        # the full bug this fixes.
        self._volume_baseline = {}
        # [FEATURE 2026-09-11] symbol -> previous session's high, same
        # populate-once-reuse-all-day pattern as _volume_baseline above.
        # Used by fast_prediction_engine.py's resistance context.
        self._prev_day_high = {}
        # [FEATURE 2026-09-16] symbol -> previous session's close, same
        # populate-once-reuse pattern as _prev_day_high above. Only
        # consumed by breakout_scanner.py's gap_pct (see _run_premarket_
        # scan() below) -- sip_bot's own pipeline has no use for this.
        self._prev_close = {}
        # [FEATURE 2026-09-15] symbol -> trailing-30-day avg daily ATR%,
        # same populate-once-reuse-all-day pattern as _volume_baseline
        # above. Feeds scorer.py's structural volatility floor (catches
        # BDCs/pre-merger SPACs like WHF/APXT that pass every other
        # filter) -- see premarket_scanner.compute_volatility_baseline().
        self._volatility_baseline = {}
        self._last_health_eval = None
        self._last_full_rescan = None
        self._shutdown = False

        # [FEATURE 2026-09-11] Caller-held per-symbol state for
        # fast_entry_gate.evaluate_fast_entry()'s persistence timer.
        self._fast_engine_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._run_breakout_scanner_shadow()
        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, baseline_out=self._volume_baseline,
            prev_day_high_out=self._prev_day_high, prev_close_out=self._prev_close)
        log.info(f"[PREMARKET] Prefiltered to {len(self._prefiltered_universe)} symbols in price/volume band")

        self._volatility_baseline = premarket_scanner.compute_volatility_baseline(
            self.client, self._prefiltered_universe)

        self.premarket_20 = premarket_scanner.scan(
            self.cfg["candidates"]["premarket_candidate_count"],
            prefiltered=self._prefiltered_universe,
            volume_baselines=self._get_volume_baselines(),
            prev_day_highs=self._prev_day_high,
            volatility_baselines=self._get_volatility_baselines(),
        )

        # [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 _run_breakout_scanner_shadow(self):
        """
        [FEATURE 2026-09-16] Runs breakout_scanner.py's candidate-scoring
        model (docs/breakout_bot_design_conversation.pdf's Part II
        design) once, right after the real premarket scan above, against
        the SAME cached prefiltered universe/baselines -- purely for
        side-by-side comparison (writes data/breakout_scanner/{date}_
        scanner.json, a completely separate file from sip_bot's own
        candidate output). Never touches self.premarket_20/self.final_10
        or anything else the live trading pipeline reads, and unlike
        _run_premarket_scan() above, nothing calls this again during the
        day -- it's a one-shot morning comparison, not part of the
        30-minute intraday rescan cadence.

        Runs on a background thread rather than inline: measured
        2026-09-16 at ~17s for a 307-symbol prefiltered pool (scan +
        compute_range_20d_high's extra bulk daily-bars call) -- scaled to
        a real ~540-symbol 9:29 premarket pool that's roughly another
        30s stacked sequentially on top of the real scan's own ~15-25s,
        which this project's own schedule.premarket_scan_time note
        already flags as having very little buffer before market_open_
        time. A background thread means this can take as long as it
        needs without ever pushing _start_streaming()/market open later.
        Reads (self.client / prefiltered/baseline data) are snapshotted
        onto plain copies before the thread starts, since self._volume_
        baseline/_prev_day_high keep getting mutated by the live pipeline
        as the day goes on (e.g. replacement candidates) -- the thread
        works off its own copies rather than racing those live reads/
        writes. self.client itself IS shared with the main thread for
        the rest of the session (concurrent read-only GET calls against
        alpaca-py's underlying requests.Session/urllib3 connection pool,
        which pools connections per-thread safely for this kind of
        simple concurrent usage) -- any failure here is caught and
        logged, never raised into the main thread.

        Wrapped in try/except inside the worker: a bug in this shadow
        scanner must never delay market open or crash the real session
        -- see breakout_scanner.py's own module docstring for the full
        design and why it's an intentionally separate, differently-
        scored implementation rather than a variant of premarket_
        scanner.py.
        """
        if not self.cfg.get("breakout_scanner", {}).get("enabled", True):
            return

        prefiltered_snapshot = list(self._prefiltered_universe)
        volume_baselines_snapshot = dict(self._get_volume_baselines())
        prev_day_highs_snapshot = dict(self._prev_day_high)
        prev_closes_snapshot = dict(self._prev_close)

        def _worker():
            try:
                range_20d_highs = breakout_scanner.compute_range_20d_high(
                    self.client, prefiltered_snapshot)
                breakout_scanner.scan(
                    self.client,
                    prefiltered=prefiltered_snapshot,
                    volume_baselines=volume_baselines_snapshot,
                    prev_day_highs=prev_day_highs_snapshot,
                    prev_closes=prev_closes_snapshot,
                    range_20d_highs=range_20d_highs,
                )
            except Exception:
                log.exception("[BREAKOUT-SCANNER] Shadow scan failed -- continuing sip_bot session unaffected")

        threading.Thread(target=_worker, name="breakout-scanner-shadow", daemon=True).start()

    def _get_volume_baselines(self) -> dict:
        """[Phase 5 -- 2026-09-08] Defensive accessor for self._volume_
        baseline -- some of this project's existing tests construct a
        SessionOrchestrator via __new__() (bypassing __init__ entirely) and
        set only the specific attributes each test needs, per this
        project's established lightweight-test-double pattern (see
        tests/test_resistance_freeze_bugfix.py). getattr with a default
        keeps this new attribute from breaking those tests, and behaves
        identically to a plain self._volume_baseline read for every real
        session (where __init__ always sets it)."""
        return getattr(self, "_volume_baseline", {})

    def _get_prev_day_highs(self) -> dict:
        """Same defensive-accessor pattern as _get_volume_baselines() above,
        for the same reason (test doubles built via __new__() that don't
        set every __init__ attribute) -- see that method's docstring."""
        return getattr(self, "_prev_day_high", {})

    def _get_volatility_baselines(self) -> dict:
        """Same defensive-accessor pattern as _get_volume_baselines() above,
        for the same reason (test doubles built via __new__() that don't
        set every __init__ attribute) -- see that method's docstring."""
        return getattr(self, "_volatility_baseline", {})

    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"]
            # [FEATURE 2026-09-11] bars_sub/quote feed the momentum-fade
            # exit check (see position_manager._check_momentum_fade_exit)
            # -- the same fast, sub-30s-bucket stream_features.py pipeline
            # fast_entry_gate.py already uses for entries, now reused on
            # the exit side instead of intraday_health's much slower
            # session-level slope.
            bars_sub = self.stream.get_bars_sub(symbol)
            quote = self.stream.get_quote(symbol)
            exit_signal = self.position_mgr.update_position(
                symbol, current_price, bars, bars_sub=bars_sub, quote=quote)
            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,
                                                 volume_baselines=self._get_volume_baselines(),
                                                 volatility_baselines=self._get_volatility_baselines())
        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-09-11 v2] Ranked shortlist, not first-come-first-served
        # list order -- but the ranking key is now fast_prediction_engine's
        # own Direction/Confidence reading, not intraday_health.health_score.
        # Per explicit instruction ("the confidence value should be the
        # value used for selecting the top_5"): every poll cycle, all
        # not-yet-open, health-eligible premarket_20 symbols first get a
        # fast_pipeline.compute_fast_ranking() call (features + fast_
        # prediction only, no persistence-gate work yet), are ranked by
        # prediction.confidence, and only the top
        # intraday_health.entry_shortlist_size (default 5) go on to a full
        # evaluate_fast_entry_only() persistence-gate call 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 gated this cycle;
        # it's re-ranked fresh on the next one as confidence changes, so
        # nothing is permanently excluded the way REMOVED is.
        #
        # health_score/health state is still a hard ELIGIBILITY filter
        # (WATCH/HEALTHY-only, staleness, benching) below -- only the
        # RANKING among eligible candidates changed.
        #
        # 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)
        sub_bucket_seconds = self.cfg["streaming"].get("sub_minute_bucket_seconds", 30)

        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

            eligible.append(r)

        # Rank ALL health-eligible candidates by fast_prediction confidence
        # -- not just today's previous shortlist -- so a symbol that was
        # ranked low (or unranked) last cycle can still compete for a slot
        # the moment its own confidence rises.
        ranked = []
        for r in eligible:
            symbol = r["symbol"]
            bars = self.stream.get_bars(symbol)
            if not bars or len(bars) < 3:
                log.debug(f"[ENTRY] Skipping {symbol}: insufficient live bars "
                          f"for a fresh ranking read")
                continue
            bars_sub = self.stream.get_bars_sub(symbol)
            quote = self.stream.get_quote(symbol)
            baseline = self._opening_volume_baseline.get(symbol, 1)
            real_baseline = self._get_volume_baselines().get(symbol, baseline)
            fast_state = self._fast_engine_state.setdefault(symbol, {})
            prediction, features = compute_fast_ranking(
                symbol, bars, bars_sub, quote, real_baseline, state=fast_state,
                sub_bucket_seconds=sub_bucket_seconds, as_of=now, premarket_result=r)
            ranked.append((prediction.confidence, symbol, r, bars, quote, fast_state, prediction, features))

        ranked.sort(key=lambda t: t[0], reverse=True)
        shortlist = ranked[:shortlist_size]
        shortlisted_symbols = {sym for _, sym, *_ in shortlist}

        if shortlist:
            log.debug(f"[ENTRY] Shortlist this cycle: " +
                      ", ".join(f"{sym}(confidence={conf:.1f})" for conf, sym, *_ in shortlist))

        # [FEATURE 2026-09-11 v3] Full raw-indicator snapshot for EVERY
        # ranked candidate this cycle, not just the top-5 that go on to a
        # real gate decision -- per explicit instruction ("make sure we
        # have all the indicator values logged during the life cycle of
        # the symbol"). Without this, a symbol sitting in premarket_20 but
        # never cracking the top-5 (e.g. HGTY-style slow builders, or
        # anything that loses a close ranking race) would leave NO record
        # of what the engine saw it doing while it waited -- only
        # shortlisted cycles were ever logged before. Non-shortlisted
        # candidates get a snapshot with gate fields left None (never
        # evaluated this cycle, not "rejected"); the shortlist loop below
        # logs the same shape WITH the real gate decision merged in, so
        # both cases land in the same file/schema for one clean
        # pandas.read_json(lines=True) load afterward.
        for rank, (confidence, symbol, r, bars, quote, fast_state, prediction, features) in enumerate(ranked, start=1):
            if symbol in shortlisted_symbols:
                continue  # logged with its real gate decision below instead
            data_store.append_fast_engine_decision({
                "symbol": symbol, "shortlisted": False, "confidence_rank": rank,
                "confidence": confidence, "state": None, "should_enter": None,
                "direction": prediction.direction, "momentum": prediction.momentum,
                "volume_state": prediction.volume_state, "vwap_state": prediction.vwap_state,
                "ema9_state": prediction.ema9_state, "resistance": prediction.resistance,
                "extension": prediction.extension, "persistence_seconds_elapsed": None,
                "reasons_for": [], "reasons_against": [],
                **snapshot_dict(features, prediction),
            })

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

        for rank, (confidence, symbol, r, bars, quote, fast_state, prediction, features) in enumerate(shortlist, start=1):
            if not self.position_mgr.has_available_slot():
                return

            # [BUGFIX 2026-08-18] The eligibility check above used
            # state/intraday_health.json -- fine for cheaply filtering
            # eligibility, 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, purely for this log/visibility check.
            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-09-11] fast_pipeline's fast_prediction_engine ->
            # fast_entry_gate chain is now the ONLY entry decision path --
            # see fast_prediction_engine.py's module docstring for the full
            # design history (2026-09-10's zero-trade session) and this
            # session's cleanup notes for why the earlier legacy/
            # prediction_pipeline alternatives were removed rather than
            # kept selectable. prediction/features were already computed
            # during ranking above -- evaluate_fast_entry_only() just
            # finishes the persistence-gate step, no recompute.
            decision = evaluate_fast_entry_only(symbol, prediction, features, state=fast_state, as_of=now)

            # [FEATURE 2026-09-11] Log EVERY decision (BUY/WAIT/REJECT
            # alike), always at INFO -- not gated behind logging.log_rejections.
            # Per explicit instruction ("log all bot decisions so we can
            # review and adjust what is not working"): under-logging here
            # would defeat the whole point of an evaluation run.
            log.info(f"[FAST] {symbol} {decision.state} (confidence={decision.confidence}, "
                     f"direction={decision.direction}, momentum={decision.momentum}, "
                     f"volume={decision.volume_state}, vwap={decision.vwap_state}, "
                     f"ema9={decision.ema9_state}, resistance={decision.resistance}, "
                     f"extension={decision.extension}, "
                     f"held={decision.persistence_seconds_elapsed:.1f}s): "
                     f"{'; '.join(decision.reasons_for or decision.reasons_against) or 'no specific reason recorded'}")
            data_store.append_fast_engine_decision({
                "symbol": symbol, "shortlisted": True, "confidence_rank": rank,
                "state": decision.state, "should_enter": decision.should_enter,
                "confidence": decision.confidence, "direction": decision.direction,
                "momentum": decision.momentum, "volume_state": decision.volume_state,
                "vwap_state": decision.vwap_state, "ema9_state": decision.ema9_state,
                "resistance": decision.resistance, "extension": decision.extension,
                "persistence_seconds_elapsed": decision.persistence_seconds_elapsed,
                "reasons_for": decision.reasons_for, "reasons_against": decision.reasons_against,
                **snapshot_dict(features, prediction),
            })
            # [BUGFIX 2026-09-15] Don't count insufficient_data WAITs against
            # the consecutive-failure streak -- that's warm-up (fewer than
            # fast_prediction_engine's min_bars one-minute bars), not a real
            # evaluated-and-rejected result, and it resolves on its own in a
            # few minutes regardless of benching. Counting it here let
            # symbols get benched for confirmation_failure_cooldown_seconds
            # purely for not having enough bars yet, since
            # poll_interval_seconds_intraday (5s) x max_consecutive_
            # confirmation_failures (5) = 25s -- far faster than the warm-up
            # window itself -- so a symbol could cycle bench/unbench through
            # its entire warm-up without ever getting a real evaluation.
            if not getattr(prediction, "insufficient_data", False):
                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 otherwise never touches position_manager.py or
        triggers an exit -- an existing open position's trailing stop is
        left alone regardless of what health state its symbol reports
        here. [FEATURE 2026-09-11] One narrow exception: if the symbol
        has an open position AND intraday_health.
        should_force_exit_on_deterioration() confirms its health has
        read STALE-or-worse for exit_confirm_reads consecutive
        evaluations in a row, force-close it here -- see
        intraday_health.py's module docstring for why.
        """
        fallback_baseline = self.cfg["universe"]["min_avg_daily_volume"]
        health_state = data_store.load_health_state()
        open_symbols = set(self.position_mgr.get_open_symbols())

        for r in self.premarket_20:
            symbol = r["symbol"]
            bars = self.stream.get_bars(symbol)
            if not bars or len(bars) < 3:
                continue
            # [Phase 5 -- 2026-09-08] real per-symbol baseline, see
            # self._volume_baseline's own comment.
            avg_vol_baseline = self._get_volume_baselines().get(symbol, fallback_baseline)
            reading, new_entry = intraday_health.evaluate_symbol(
                symbol, bars, avg_vol_baseline, health_state)
            health_state[symbol] = new_entry

            if symbol in open_symbols and intraday_health.should_force_exit_on_deterioration(
                    new_entry, self.cfg["intraday_health"]):
                current_price = bars[-1]["c"]
                log.info(
                    f"[HEALTH-EXIT] {symbol} health confirmed deteriorating for "
                    f"{new_entry['consecutive_deteriorating']} consecutive reads "
                    f"(state={new_entry['state']}); forcing exit"
                )
                self.position_mgr.exit_position(symbol, current_price, "HEALTH_DETERIORATION")
                top_stocks.register_cooldown(symbol)
                self._handle_slot_freed(symbol)

        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"]
        fallback_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"],
            volume_baselines=self._get_volume_baselines(),
            prev_day_highs=self._get_prev_day_highs(),
            volatility_baselines=self._get_volatility_baselines(),
        )

        eligible = []
        for r in rescored:
            symbol = r["symbol"]
            avg_vol_baseline = self._get_volume_baselines().get(symbol, fallback_baseline)
            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()