o
    M¹ªj«Ë  ã                   @   s
  d Z ddlZddlZddlZddlZddlZddlmZmZmZ ddl	m
Z
mZmZ ddlmZ ddlZddlZddlZddlZddlZddlZddlmZmZmZ ddlmZ ddlmZ dd	lmZ ed
ƒZG dd„ dƒZ G dd„ dƒZ!dd„ Z"e#dkrƒe"ƒ  dS dS )ax  
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.
é    N)ÚdatetimeÚtimezoneÚ	timedelta)Ú
get_configÚget_envÚvalidate_config)Ú
get_logger)Úcompute_fast_rankingÚevaluate_fast_entry_onlyÚsnapshot_dict)ÚPositionManager)ÚStreamManager)Ú
get_clientÚmonitorc                   @   sN   e Zd ZdZdedefdd„Zdedede	fd	d
„Z
dede	defdd„ZdS )Ú_EntryAttemptTrackeraš  
    [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.
    Úmax_failuresÚcooldown_secondsc                 C   s   || _ || _i | _i | _d S ©N)r   r   Ú_consecutive_failuresÚ_benched_until)Úselfr   r   © r   ú
monitor.pyÚ__init__b   s   
z_EntryAttemptTracker.__init__ÚsymbolÚnowÚreturnc                 C   s   | j  |¡}|d uo||k S r   )r   Úget)r   r   r   Úuntilr   r   r   Ú
is_benchedh   s   z_EntryAttemptTracker.is_benchedÚ	confirmedc              	   C   s”   |r| j  |d ¡ | j |d ¡ d S | j  |d¡d }|| jkrC|t| jd� | j|< d| j |< t d|› d| jd›d|› d�¡ d S || j |< d S )	Nr   é   )Úsecondsú[ENTRY] z& benched from the entry shortlist for z.0fzs after zE consecutive not-confirmed attempts -- giving other candidates a turn)	r   Úpopr   r   r   r   r   ÚlogÚinfo)r   r   r    r   Úcountr   r   r   Úrecord_resultl   s   

ÿÿz"_EntryAttemptTracker.record_resultN)Ú__name__Ú
__module__Ú__qualname__Ú__doc__ÚintÚfloatr   Ústrr   Úboolr   r(   r   r   r   r   r   E   s
    r   c                   @   sÊ   e Zd Zdd„ Zdd„ Zdd„ Zdd„ Zd	d
„ Zdd„ Zdd„ Z	de
fdd„Zde
fdd„Zde
fdd„Zdefdd„Zdd„ Zdd„ Zdd„ Zdd „ Zd!efd"d#„Zd$d%„ Zd&d'„ Zd(d)„ Zd*d+„ Zd,S )-ÚSessionOrchestratorc                 C   s   t ƒ | _t| jƒ tƒ  ¡  | jd d dk| _tƒ | _t| jd�| _	t
ƒ | _t| jd d | jd d d�| _t t ¡ ¡}t |¡ g | _g | _g | _i | _i | _i | _i | _i | _d | _d | _d	| _i | _t tj | j!¡ t tj"| j!¡ t# $| j%¡ d S )
NÚmodeÚexecution_modeÚ
simulation)r4   Úintraday_healthÚ%max_consecutive_confirmation_failuresÚ%confirmation_failure_cooldown_seconds)r   r   F)&r   Úcfgr   r   Úvalidater4   r   Úclientr   Úposition_mgrr   Ústreamr   Ú_entry_trackerr5   Úprune_stale_health_stateÚ
data_storeÚload_health_stateÚsave_health_stateÚpremarket_20Úfinal_10Ú_prefiltered_universeÚ_opening_volume_baselineÚ_volume_baselineÚ_prev_day_highÚ_prev_closeÚ_volatility_baselineÚ_last_health_evalÚ_last_full_rescanÚ	_shutdownÚ_fast_engine_stateÚsignalÚSIGTERMÚ_handle_signalÚSIGINTÚatexitÚregisterÚ_cleanup)r   Úpruned_health_stater   r   r   r   ~   s8   

þ
zSessionOrchestrator.__init__c                 C   s   t  d|› d�¡ d| _d S )Nz[SHUTDOWN] Received signal z, shutting down gracefullyT)r%   r&   rL   )r   ÚsignumÚframer   r   r   rP   º   s   
z"SessionOrchestrator._handle_signalc                 C   s0   z| j  ¡  W n	 ty   Y nw t d¡ d S )Nz[SHUTDOWN] Cleanup complete)r<   ÚstopÚ	Exceptionr%   r&   ©r   r   r   r   rT   ¾   s   ÿzSessionOrchestrator._cleanupc                 C   s€   t  d| jd d › �¡ t ¡ st  d¡ d S |  ¡  | jr!d S |  ¡  |  ¡  |  	¡  |  
¡  |  ¡  |  ¡  t  d¡ d S )Nz#[START] monitor.py starting | mode=r2   r3   z[START] Not a weekday, exitingz[SHUTDOWN] Session complete)r%   r&   r8   Úmarket_timeÚ
is_weekdayÚ_wait_for_premarket_scanrL   Ú_run_premarket_scanÚ_run_breakout_scanner_shadowÚ_start_streamingÚ_wait_for_market_openÚ_main_trading_loopÚ_end_of_day_liquidationrZ   r   r   r   ÚrunÆ   s   
zSessionOrchestrator.runc                 C   ó6   t  ¡ s| jst d¡ t  ¡ s| jrd S d S d S d S )Né   )r[   Úis_premarket_scan_timerL   ÚtimeÚsleeprZ   r   r   r   r]   Û   ó   
ÿz,SessionOrchestrator._wait_for_premarket_scanc                 C   sò   t  | j¡}t dt|ƒ› �¡ t j| j|| j| j| j	d�| _
t dt| j
ƒ› d�¡ t  | j| j
¡| _t j| jd d | j
|  ¡ | j|  ¡ d�| _| jd | jd d … | _t d	d
„ | jD ƒdd
„ | jD ƒdœ¡ t | j¡ |  | j¡ d S )Nz1[PREMARKET] Universe size after asset filtering: )Úbaseline_outÚprev_day_high_outÚprev_close_outz[PREMARKET] Prefiltered to z symbols in price/volume bandÚ
candidatesÚpremarket_candidate_count)ÚprefilteredÚvolume_baselinesÚprev_day_highsÚvolatility_baselinesÚfinal_candidate_countc                 S   ó   g | ]}|d  |d dœ‘qS ©r   Útotal_score)r   Úscorer   ©Ú.0Úrr   r   r   Ú
<listcomp>  ó    z;SessionOrchestrator._run_premarket_scan.<locals>.<listcomp>c                 S   ru   rv   r   ry   r   r   r   r|     r}   )rB   rC   )Úpremarket_scannerÚget_universe_symbolsr:   r%   r&   ÚlenÚprefilter_by_snapshotrF   rG   rH   rD   Úcompute_volatility_baselinerI   Úscanr8   Ú_get_volume_baselinesÚ_get_volatility_baselinesrB   rC   r?   Úsave_watchlistÚwrite_top_stocksÚ_sync_opening_volume_baseline)r   Úuniverser   r   r   r^   ß   s0   
þÿûþz'SessionOrchestrator._run_premarket_scanc                    sp   ˆj  di ¡ dd¡sdS tˆjƒ‰ tˆ ¡ ƒ‰tˆjƒ‰tˆjƒ‰‡ ‡‡‡‡fdd„}tj	|ddd� 
¡  dS )	aÏ	  
        [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.
        Úbreakout_scannerÚenabledTNc                     sL   zt  ˆjˆ ¡} t jˆjˆ ˆˆˆ| d� W d S  ty%   t d¡ Y d S w )N)rp   rq   rr   Úprev_closesÚrange_20d_highszN[BREAKOUT-SCANNER] Shadow scan failed -- continuing sip_bot session unaffected)rŠ   Úcompute_range_20d_highr:   rƒ   rY   r%   Ú	exception)r�   ©Úprefiltered_snapshotÚprev_closes_snapshotÚprev_day_highs_snapshotr   Úvolume_baselines_snapshotr   r   Ú_worker;  s   ÿúÿzASessionOrchestrator._run_breakout_scanner_shadow.<locals>._workerzbreakout-scanner-shadow)ÚtargetÚnameÚdaemon)r8   r   ÚlistrD   Údictr„   rG   rH   Ú	threadingÚThreadÚstart)r   r•   r   r�   r   r_   	  s   *


z0SessionOrchestrator._run_breakout_scanner_shadowr   c                 C   ó   t | di ƒS )a^  [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).rF   ©ÚgetattrrZ   r   r   r   r„   L  s   
z)SessionOrchestrator._get_volume_baselinesc                 C   rž   )úÐ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.rG   rŸ   rZ   r   r   r   Ú_get_prev_day_highsX  ó   z'SessionOrchestrator._get_prev_day_highsc                 C   rž   )r¡   rI   rŸ   rZ   r   r   r   r…   ^  r£   z-SessionOrchestrator._get_volatility_baselinesrn   c                 C   s6   |D ]}|  di ¡}t|  dd¡dƒ| j|d < qdS )aY  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.ÚmetricsÚpremarket_volumer!   r   N)r   ÚmaxrE   )r   rn   r{   r¤   r   r   r   rˆ   d  s   þz1SessionOrchestrator._sync_opening_volume_baselinec                 C   sT   dd„ | j D ƒ}|st d¡ d S | j |¡ | j ¡ r#t d¡ d S t d¡ d S )Nc                 S   ó   g | ]}|d  ‘qS ©r   r   ry   r   r   r   r|   s  ó    z8SessionOrchestrator._start_streaming.<locals>.<listcomp>z5[OPEN] No candidates to stream; skipping stream startz[OPEN] SIP stream activezF[OPEN] SIP stream failed to become healthy; will rely on REST fallback)rB   r%   Úwarningr<   r�   Ú
is_healthyr&   Úerror)r   Úsymbolsr   r   r   r`   o  s   

z$SessionOrchestrator._start_streamingc                 C   re   )Né   )r[   Úis_past_market_openrL   rh   ri   rZ   r   r   r   ra   }  rj   z)SessionOrchestrator._wait_for_market_openc                 C   s  | j d d }| j d }t tj¡}|| _|| _| js‚t 	¡ s„| j
 ¡  |  ¡  | j
 ¡ r5t ¡ s5|  ¡  | j
 ¡  t tj¡}| dd¡rZ|| j  ¡ }||d krZ|  ¡  || _t ¡ st|| j  ¡ }||d d krt|  ¡  || _t |¡ | js†t 	¡ rd S d S d S d S )	NÚscheduleÚpoll_interval_seconds_intradayr5   r‹   TÚeval_interval_secondsÚfull_rescan_interval_minutesé<   )r8   r   r   r   ÚutcrJ   rK   rL   r[   Úis_force_liquidate_timer;   Úpoll_pending_exitsÚ_update_open_positionsÚhas_available_slotÚis_new_entries_cutoffÚ_scan_for_entriesÚreconcile_with_brokerr   Útotal_secondsÚ_update_intraday_healthÚ_run_intraday_full_rescanrh   ri   )r   ÚintervalÚ
health_cfgr   Úelapsedr   r   r   rb   ‚  s0   



Üz&SessionOrchestrator._main_trading_loopc                 C   sˆ   | j  ¡ D ]<}| j |¡}|sq|d d }| j |¡}| j |¡}| j j|||||d�}|rA| j  |||¡ t 	|¡ |  
|¡ qd S )NéÿÿÿÿÚc)Úbars_subÚquote)r;   Úget_open_symbolsr<   Úget_barsÚget_bars_subÚ	get_quoteÚupdate_positionÚexit_positionÚ
top_stocksÚregister_cooldownÚ_handle_slot_freed)r   r   ÚbarsÚcurrent_pricerÅ   rÆ   Úexit_signalr   r   r   r¸   ¯  s    
ÿ

€îz*SessionOrchestrator._update_open_positionsÚfreed_symbolc                    sx   t  d¡ | j ¡ }tj||  ¡ |  ¡ d�‰ ˆ sd S | j 	ˆ d g¡ ‡ fdd„| j
D ƒ| _
| j
 ˆ ¡ |  ˆ g¡ d S )Nz [SLOT] 1 position slot available)Úexclude_symbolsrq   rs   r   c                    s    g | ]}|d  ˆ d  kr|‘qS r¨   r   ry   ©Ú	candidater   r   r|   Ü  s     z:SessionOrchestrator._handle_slot_freed.<locals>.<listcomp>)r%   r&   r;   rÇ   rÍ   Úfind_replacementr„   r…   r<   Úadd_symbolsrB   Úappendrˆ   )r   rÓ   Úopen_symbolsr   rÕ   r   rÏ   Ä  s   

þz&SessionOrchestrator._handle_slot_freedc                  C   sþ  t  ¡ }| jd d }t tj¡}| jd  dd¡}g }| jD ]P}|d }| j	 
|¡r,q| |i ¡}| dd¡}	t |	¡sIt d	|› d
|	› �¡ q| j |¡rYt d	|› d�¡ q| j ||¡rjt d	|› d�¡ q| |¡ qg }
|D ]]}|d }| j |¡}|rˆt|ƒdk r’t d	|› d�¡ qt| j |¡}| j |¡}| j |d¡}|  ¡  ||¡}| j |i ¡}t|||||||||d�	\}}|
 |j|||||||f¡ qt|
jdd„ dd� |
d |… }dd„ |D ƒ}|rùt dd dd„ |D ƒ¡ ¡ t |
dd�D ]S\}\}}}}}}}}||v �rqÿt  !i d|“dd“d|“d|“dd “d d “d!|j"“d"|j#“d#|j$“d$|j%“d%|j&“d&|j'“d'|j(“d(d “d)g “d*g “t)||ƒ¥¡ qÿ| jd+ d, }t |dd�D �]\}\}}}}}}}}| j	 *¡ �sv d S t +|||¡}| |i ¡ dd-¡}|j,|k�r¢t d.|› d/|› d0|j,› d1|j-d2›d3�	¡ t.|||||d4�}t /d5|› d6|j0› d7|j› d8|j"› d9|j#› d:|j$› d;|j%› d<|j&› d=|j'› d>|j(› d?|j1d2›d@dA |j2�pâ|j3¡�pædB› �¡ t  !i d|“dd“d|“d|j0“d |j4“d|j“d!|j"“d"|j#“d#|j$“d$|j%“d%|j&“d&|j'“d'|j(“d(|j1“d)|j2“d*|j3“t)||ƒ¥¡ t5|dCdƒ�sE| j 6||j4t tj¡¡ |j4�r{|dD dE }| j7�sX| j8 9¡ nd }| jdF dG }|�rit:|j;ƒn|}dA |j2¡}| j	 <|||||¡ �q`d S )HNr5   Úentry_shortlist_sizeÚ	streamingÚsub_minute_bucket_secondsé   r   ÚstateÚWATCHz[ENTRY] Skipping z	: health=z: stale stream dataz.: benched after repeated confirmation failuresé   z1: insufficient live bars for a fresh ranking readr!   )rß   Úsub_bucket_secondsÚas_ofÚpremarket_resultc                 S   s   | d S )Nr   r   )Útr   r   r   Ú<lambda>1  s    z7SessionOrchestrator._scan_for_entries.<locals>.<lambda>T©ÚkeyÚreversec                 S   s   h | ]^}}}|’qS r   r   )rz   Ú_Úsymr   r   r   Ú	<setcomp>3  ó    z8SessionOrchestrator._scan_for_entries.<locals>.<setcomp>z[ENTRY] Shortlist this cycle: z, c                 s   s(   � | ]^}}}|› d |d›d�V  qdS )z(confidence=ú.1fú)Nr   )rz   Úconfrë   rê   r   r   r   Ú	<genexpr>7  s   €& z8SessionOrchestrator._scan_for_entries.<locals>.<genexpr>)r�   ÚshortlistedFÚconfidence_rankÚ
confidenceÚshould_enterÚ	directionÚmomentumÚvolume_stateÚ
vwap_stateÚ
ema9_stateÚ
resistanceÚ	extensionÚpersistence_seconds_elapsedÚreasons_forÚreasons_againstr‰   Úmin_avg_daily_volumeú?r#   z: cached health=z but fresh read=z (score=rî   z0) -- using the fresh read for the entry decision)rß   rã   z[FAST] ú z (confidence=z, direction=z, momentum=z	, volume=z, vwap=z, ema9=z, resistance=z, extension=z, held=zs): z; zno specific reason recordedÚinsufficient_datarÃ   rÄ   ÚtradingÚsimulated_equity_default)=r?   r@   r8   r   r   r   rµ   r   rB   r;   Úis_symbol_openr5   Úis_eligible_for_entryr%   Údebugr<   Úis_symbol_staler=   r   rÙ   rÈ   r€   rÉ   rÊ   rE   r„   rM   Ú
setdefaultr	   rô   ÚsortÚjoinÚ	enumerateÚappend_fast_engine_decisionrö   r÷   rø   rù   rú   rû   rü   r   r¹   Úcompute_healthÚ	raw_stateÚhealth_scorer
   r&   rß   rý   rþ   rÿ   rõ   r    r(   r4   r:   Úget_accountr.   ÚequityÚenter_position) r   Úhealth_stateÚshortlist_sizer   râ   Úeligibler{   r   ÚentryÚconfirmed_stateÚrankedrÐ   rÅ   rÆ   ÚbaselineÚreal_baselineÚ
fast_stateÚ
predictionÚfeaturesÚ	shortlistÚshortlisted_symbolsÚrankrô   Úavg_vol_baselineÚfresh_readingÚcached_stateÚdecisionÚentry_priceÚaccountÚ
sim_equityr  Úreasonr   r   r   r»   à  sB  


þÿ$
ÿÿÿþþþýýüüûûúúùùø&ÿ
ÿÿÿþþýýüûúÿÿÿþþýýüüûûúúùøø	÷€·z%SessionOrchestrator._scan_for_entriesc              
   C   s   | j d d }t ¡ }t| j ¡ ƒ}| jD ]c}|d }| j |¡}|r)t	|ƒdk r*q|  
¡  ||¡}t ||||¡\}}	|	||< ||v rxt |	| j d ¡rx|d d }
t d|› d	|	d
 › d|	d › d�¡ | j ||
d¡ t |¡ |  |¡ qt |¡ dS )a½  
        [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.
        r‰   r   r   rá   r5   rÃ   rÄ   z[HEALTH-EXIT] z$ health confirmed deteriorating for Úconsecutive_deterioratingz consecutive reads (state=rß   z); forcing exitÚHEALTH_DETERIORATIONN)r8   r?   r@   Úsetr;   rÇ   rB   r<   rÈ   r€   r„   r   r5   Úevaluate_symbolÚ"should_force_exit_on_deteriorationr%   r&   rÌ   rÍ   rÎ   rÏ   rA   )r   Úfallback_baseliner  rÚ   r{   r   rÐ   r#  ÚreadingÚ	new_entryrÑ   r   r   r   r¾   £  s:   
ÿ
ÿÿþÿ

€z+SessionOrchestrator._update_intraday_healthc                 C   s*  | j s
t d¡ dS t dt| j ƒ› d�¡ | jd }| jd d }t ¡ }tj	t| j ƒ| j |d |  
¡ |  ¡ |  ¡ d	�}g }|D ][}|d
 }|  
¡  ||¡}|dd„ | jD ƒv r_| j |¡nd}	|	rit|	ƒdk ry|d |d< d|d< | |¡ qAt ||	||¡\}
}|||< |
j|d< |
j|d< t |
j¡rœ| |¡ qAt |¡ |s«t d¡ dS |jdd„ dd� |d|d … }dd„ |D ƒ}t| j ¡ ƒ}dd„ |D ƒ}| jD ]P}|d
 }||v �r$||v�r$| |¡}|du�r	| |¡ | |¡ t d|› d| dd¡d ›�¡ qÔt d|› d!| dd¡d ›d"�¡ | |¡ | |¡ qÔ|| _|  |¡ | j d#d„ |D ƒ¡ | j d$i ¡ d%d&¡}|d|… }t  ¡ }d'd„ |D ƒ|d(< d)d„ |D ƒ|d*< t !|¡ t "|¡ |D ]}t d|d
 › d+|d › d,| d¡› �¡ �qkt d-t|ƒ› d.t|ƒ› �¡ dS )/a\  
        [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.
        zD[INTRADAY-RESCAN] No cached prefiltered universe available; skippingNz*[INTRADAY-RESCAN] Starting full rescan of z prefiltered symbolsr5   r‰   r   Úfull_rescan_lookback_hours)Úcandidate_countrp   Úlookback_hoursrq   rr   rs   r   c                 S   r§   r¨   r   )rz   Úpr   r   r   r|     r©   zASessionOrchestrator._run_intraday_full_rescan.<locals>.<listcomp>rá   rw   r  rà   r  zb[INTRADAY-RESCAN] No health-eligible candidates found; keeping the existing premarket_20 unchangedc                 S   s   | d | d fS )Nr  rw   r   )r{   r   r   r   ræ     s    z?SessionOrchestrator._run_intraday_full_rescan.<locals>.<lambda>Trç   Úfull_rescan_pool_sizec                 S   s   i | ]}|d  |“qS r¨   r   ry   r   r   r   Ú
<dictcomp>@  rí   zASessionOrchestrator._run_intraday_full_rescan.<locals>.<dictcomp>c                 S   s   h | ]}|d  ’qS r¨   r   ry   r   r   r   rì   B  r©   z@SessionOrchestrator._run_intraday_full_rescan.<locals>.<setcomp>z[INTRADAY-RESCAN] z6 kept in pool (open position) with refreshed pm_high=$Úpm_highr   ú.2fz` kept in pool (open position) but missing from this rescan's universe -- reusing stale pm_high=$zj; resistance will not reflect current price action until this symbol reappears in the prefiltered universec                 S   r§   r¨   r   ry   r   r   r   r|   W  r©   rn   rt   é   c                 S   ó&   g | ]}|d  |d |  d¡dœ‘qS ©r   rw   r  )r   rx   Úhealth©r   ry   r   r   r   r|   j  ó    ÿÿrB   c                 S   r<  r=  r?  ry   r   r   r   r|   n  r@  rC   z score=z health=z*[INTRADAY-RESCAN] Rotated premarket_20 -> z candidates, final_10 -> top )#rD   r%   rª   r&   r€   r8   r?   r@   r~   rƒ   r„   r¢   r…   r   rB   r<   rÈ   rÙ   r5   r.  r  r  r  rA   r  r-  r;   rÇ   Úaddr  rˆ   rØ   Úload_watchlistr†   r‡   )r   r8   r0  r  Úrescoredr  r{   r   r#  rÐ   r1  r2  Únew_poolÚrescored_by_symbolrÚ   Úpool_symbolsÚfreshÚfinal_countÚfinal_sliceÚ	watchlistr   r   r   r¿   Ô  s¤   

ÿ
ú	$
ÿ


€

#





ÿ
þ

€

þ
þ

ÿÿz-SessionOrchestrator._run_intraday_full_rescanc           	      C   sV  t  d¡ | jjdd� | j di ¡ dd¡}| jd d }t ¡ | }t ¡ |k rN| j ¡  | j ¡ | j 	¡  }|s<n0t 
|¡ | jjdd� t ¡ |k s*| j ¡ | j 	¡  }|rlt  dt|ƒ› d	|› d
t|ƒ› �¡ | j ¡  t ¡ }t |¡ tdd„ |D ƒƒ}tdd„ |D ƒƒ}tdd„ |D ƒƒ}t  dt|ƒ› d|› d|› d|d›�¡ dS )a]  
        [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.
        z*[EOD] Force-liquidating all open positionsÚ
END_OF_DAY)r*  r°   Ú eod_liquidation_max_wait_secondséx   r±   z[EOD] z& position(s) still not resolved after zLs of retries -- will fall through to next session's startup reconciliation: c                 s   s*   � | ]}|  d ¡dkr|  dd¡V  qdS )ÚstatusÚclosedÚ
current_plr   Nr?  ©rz   rå   r   r   r   rñ   ¼  s   €( z>SessionOrchestrator._end_of_day_liquidation.<locals>.<genexpr>c                 s   s2   � | ]}|  d ¡dkr|  dd¡dkrdV  qdS ©rN  rO  rP  r   r!   Nr?  rQ  r   r   r   rñ   ½  ó   €0 c                 s   s2   � | ]}|  d ¡dkr|  dd¡dkrdV  qdS rR  r?  rQ  r   r   r   rñ   ¾  rS  z[EOD] Session summary: trades=z wins=z losses=z total_P/L=$r:  N)r%   r&   r;   Úliquidate_allr8   r   rh   r·   rÇ   Úget_closing_symbolsri   rª   r€   Úsortedr<   rX   r?   Úload_today_tradesÚwrite_trades_summaryÚsum)	r   Úmax_waitÚpoll_intervalÚdeadlineÚ
still_openÚtradesÚtotal_plÚwinsÚlossesr   r   r   rc   |  s:   


õÿþ

ÿz+SessionOrchestrator._end_of_day_liquidationN)r)   r*   r+   r   rP   rT   rd   r]   r^   r_   rš   r„   r¢   r…   r™   rˆ   r`   ra   rb   r¸   r/   rÏ   r»   r¾   r¿   rc   r   r   r   r   r1   }   s.    <*C- D1 )r1   c               
   C   sˆ   z$t  ¡ � t  ¡  tƒ } |  ¡  W d   ƒ W d S 1 sw   Y  W d S  tyC } zt t|ƒ¡ t	 
d¡ W Y d }~d S d }~ww )Nr!   )r?   Úacquire_singleton_lockÚensure_dirsr1   rd   ÚRuntimeErrorr%   r¬   r/   ÚsysÚexit)ÚorchestratorÚer   r   r   ÚmainÃ  s   

&ý€þri  Ú__main__)$r,   re  rh   rN   rR   r›   r   r   r   Úconfig_loaderr   r   r   Úlogger_setupr   r[   r?   r~   rŠ   rÍ   r5   Úfast_pipeliner	   r
   r   Úposition_managerr   r<   r   Úalpaca_clientr   r%   r   r1   ri  r)   r   r   r   r   Ú<module>   s@    -8      L
ÿ