o
    #¹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ZddlmZmZ ddl	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mZ ddlmZ ddlmZ dd	lmZmZmZ ed
ƒZej ej e ¡¡Z!ej "e!d¡Z#G dd„ dƒZ$G dd„ dƒZ%G dd„ dƒZ&dd„ Z'e(dkr“e'ƒ  dS dS )aõ  
monitor.py

Bot CORE orchestrator for screener/trade1. Contains NO entry or exit
strategy: every trade decision comes from the rule modules named in
config.json "rules" (breakout_rules.py, reversal_rules.py,
exit_rules.py -- contract in rules_api.py). Work on entries/exits
happens in those files only.

    START -> wait for schedule.premarket_scan_time
          -> universe.py + scanner.py (today's candidate list)
          -> wait for market open, then subscribe the stream (only after
             the open, so premarket prints never land in session bars)
          -> every schedule.poll_interval_seconds:
               apply a finished background rescan (scanner.rescan_*)
               each OPEN position   -> exit module.evaluate() -> sell?
               each WATCHED symbol  -> entry modules in order  -> buy?
                  (while a slot is free, before no_new_entries_after)
          -> force_liquidate_time: sell everything, write the summary

What stays in the core (not strategy -- plumbing and safety):
  scanning/rescans, streaming, order placement, position sizing
  (trading.account_risk_pct_per_trade to the rule's stop, notional cap),
  max_positions, end-of-day liquidation, and a fallback that sells at
  the entry stop ONLY if the exit module can't load or raises.

Run directly:  python monitor.py
A singleton lock (data_store.acquire_singleton_lock) stops two copies
running from this folder. NOTE: screener/trade uses the SAME Alpaca
paper account -- never run both bots at the same time.
é    N)ÚdatetimeÚtimezone)Ú
get_configÚget_env)Ú
get_logger)ÚPositionManager)ÚStreamManager)Ú
get_client)Ú
MarketViewÚEntryDecisionÚExitDecisionÚmonitorzmonitor.pidc                   @   sˆ   e Zd ZdZdedefdd„Zdd„ Zdefd	d
„Zdd„ Z	dedefdd„Z
dddœdedefdd„Zdefdd„Zdefdd„ZdS )Ú
RuleModulezËOne loaded rule file. Wraps every call so a rule error is logged
    (once per symbol+method, to keep the log readable) and never raises
    into the core. Optionally re-imports the file when it changes.ÚnameÚ
hot_reloadc                 C   s,   || _ || _d | _d | _tƒ | _|  ¡  d S ©N)r   r   ÚmodÚ_mtimeÚsetÚ_errors_loggedÚload)Úselfr   r   © r   ú#/var/www/screener/trade1/monitor.pyÚ__init__?   s   zRuleModule.__init__c                 C   s(   t | jdd ƒ}|ptj t| j› d�¡S )NÚ__file__z.py)Úgetattrr   ÚosÚpathÚjoinÚBASE_DIRr   )r   Úfr   r   r   Ú_pathG   s   zRuleModule._pathÚreturnc                 C   s–   z.| j d u rt | j¡| _ nt | j ¡| _ tj |  ¡ ¡| _	| j
 ¡  t d| j› �¡ W dS  tyJ   t d| j› �| j d urCdnd ¡ Y dS w )Nz[RULES] loaded Tz[RULES] failed to load z  -- keeping the previous versionÚ F)r   Ú	importlibÚimport_moduler   Úreloadr   r   Úgetmtimer"   r   r   ÚclearÚlogÚinfoÚ	ExceptionÚ	exception©r   r   r   r   r   K   s   

ÿýzRuleModule.loadc                 C   sh   | j sd S z
tj |  ¡ ¡}W n
 ty   Y d S w || jkr2|| _t d| j	› d�¡ |  
¡  d S d S )Nú[RULES] z .py changed on disk -- reloading)r   r   r   r(   r"   ÚOSErrorr   r*   r+   r   r   )r   Úmtimer   r   r   Úmaybe_reloadZ   s   ÿ
ýzRuleModule.maybe_reloadÚmethodc                 C   s   | j d uott| j |d ƒƒS r   )r   Úcallabler   )r   r3   r   r   r   Úhasf   s   zRuleModule.hasr$   N)ÚsymbolÚdefaultr6   c                G   sv   |   |¡s|S z	t| j|ƒ|Ž W S  ty:   ||f}|| jvr6| j |¡ t d| j› d|› d|› d�¡ | Y S w )Nr/   Ú.ú(z,) raised -- ignored (logged once per symbol))	r5   r   r   r,   r   Úaddr*   r-   r   )r   r3   r6   r7   ÚargsÚkeyr   r   r   Úcalli   s   

 úzRuleModule.callc                 C   s   t ƒ  | ji ¡S r   )r   Úgetr   r.   r   r   r   Úcfgv   s   zRuleModule.cfgc                 C   s   | j d uo|  ¡  dd¡S )NÚenabledT)r   r?   r>   r.   r   r   r   r@   y   ó   zRuleModule.enabled)Ú__name__Ú
__module__Ú__qualname__Ú__doc__ÚstrÚboolr   r"   r   r2   r5   r=   Údictr?   r@   r   r   r   r   r   :   s    r   c                   @   s@   e Zd ZdZdd„ Zdd„ Zdd„ Zdd	„ Zd
d„ Zdd„ Z	dS )Ú_TickFanoutz³Stream listener that forwards raw ticks to whichever rule modules
    define on_trade/on_quote/on_bar/seed_bars (looked up per call, so a
    hot reload takes effect immediately).c                 C   s
   || _ d S r   )Úmodules)r   rJ   r   r   r   r   ‚   s   
z_TickFanout.__init__c                 G   sB   | j D ]}| |¡r|j|g|¢R d|rt|d ƒndiŽ qd S )Nr6   r   r$   )rJ   r5   r=   rF   )r   r3   r;   Úmr   r   r   Ú_fan…   s
   

(€þz_TickFanout._fanc                 G   ó   | j dg|¢R Ž  d S )NÚon_trade©rL   ©r   Úar   r   r   rN   Š   ó   z_TickFanout.on_tradec                 G   rM   )NÚon_quoterO   rP   r   r   r   rS   �   rR   z_TickFanout.on_quotec                 G   rM   )NÚon_barrO   rP   r   r   r   rT   �   rR   z_TickFanout.on_barc                 G   rM   )NÚ	seed_barsrO   rP   r   r   r   rU   “   rR   z_TickFanout.seed_barsN)
rB   rC   rD   rE   r   rL   rN   rS   rT   rU   r   r   r   r   rI   }   s    rI   c                   @   s  e Zd ZdZ		d7dd„Zdd„ Zdd	„ Zd
d„ Zd8dede	fdd„Z
de	f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efd d!„Zd9d"ed#ed$edefd%d&„Zd:d(ed"ed)ed*ed+e	f
d,d-„Zd"ed.ed*efd/d0„Zd1d2„ Zd3d4„ Zd5d6„ ZdS );ÚSessionOrchestratorzÇLive by default. simulate.py builds the SAME class with a replay
    clock, a simulated broker (position_mgr) and a replayed stream, so a
    backtest runs exactly the decision code that trades live.NTc           	         sh  t ƒ | _|| _|rtƒ  ¡  |p|rtƒ nd | _|ptƒ | _|p"t	ƒ | _
|p)dd„ | _|p/tj| _| j di ¡}t| dd¡ƒ‰ ‡ fdd„| dg ¡D ƒ| _| d	¡rZt|d	 ˆ ƒnd | _| j| jrf| jgng  }| j
j t|ƒ¡ || _g | _i | _i | _i | _i | _i | _i | _d| _d | _ d| _!d | _"d | _#t$ %¡ | _&|r²t' 't'j(| j)¡ t' 't'j*| j)¡ d S d S )
Nc                   S   s   t  tj¡S r   )r   Únowr   Úutcr   r   r   r   Ú<lambda>¦   s    z.SessionOrchestrator.__init__.<locals>.<lambda>Úrulesr   Fc                    s   g | ]}t |ˆ ƒ‘qS r   )r   )Ú.0Ún©Úhotr   r   Ú
<listcomp>«   s    z0SessionOrchestrator.__init__.<locals>.<listcomp>Úentry_modulesÚexit_module)+r   r?   Úliver   Úvalidater	   Úclientr   Úposition_mgrr   ÚstreamÚ_nowÚ
data_storeÚappend_decision_recordÚ_decision_sinkr>   rG   r`   r   ra   Ú	listenersÚappendrI   Úall_modulesÚ
candidatesÚ_scan_by_symbolÚ_levelsÚ_avg_vol_baselineÚ_entry_stateÚ_exit_stateÚ_last_loggedÚ	_shutdownÚ_last_scan_startedÚ_early_rescan_doneÚ_rescan_threadÚ_pending_scanÚ	threadingÚLockÚ_pending_lockÚsignalÚSIGTERMÚ_handle_signalÚSIGINT)	r   rd   re   rf   ÚclockÚdecision_sinkrb   ÚrcÚall_modsr   r]   r   r   œ   sB   

þzSessionOrchestrator.__init__c                 C   s   t  d|› d�¡ d| _d S )Nz[SHUTDOWN] Received signal z, shutting down gracefullyT)r*   r+   ru   )r   ÚsignumÚframer   r   r   r   Ç   s   
z"SessionOrchestrator._handle_signalc                 C   s¾   t  d| jd d › ddd„ | jD ƒ› d| jr| jjnd › �¡ | jd u s+| jjd u r0t  d¡ t 	¡ s;t  d	¡ d S |  
¡  | jrDd S |  ¡  |  ¡  |  ¡  |  ¡  |  ¡  t  d
¡ d S )Nz1[START] monitor.py (trade1 core) starting | mode=ÚmodeÚexecution_modez | entry modules=c                 S   s   g | ]}|j ‘qS r   )r   )r[   rK   r   r   r   r_   Î   s    z+SessionOrchestrator.run.<locals>.<listcomp>z | exit module=zr[RULES] NO working exit module -- positions will only be sold at their entry stop (core fallback) or at end of dayz[START] Not a weekday, exitingz[SHUTDOWN] Session complete)r*   r+   r?   r`   ra   r   r   ÚerrorÚmarket_timeÚ
is_weekdayÚ_wait_for_premarket_scanru   Ú	_run_scanÚ_wait_for_market_openÚ_start_streamingÚ_main_trading_loopÚ_end_of_day_liquidationr.   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_timeru   ÚtimeÚsleepr.   r   r   r   rŒ   å   ó   
ÿz,SessionOrchestrator._wait_for_premarket_scanr$   Úcandidates_file_suffixr#   c              	   C   s¢   t  | j¡}t dt|ƒ› �¡ i i i }}}t j| j||||d�}t dt|ƒ› d�¡ i i }}tj| j|||d�}	tj	| j|||||	|d�}
|
|||dœS )zePure scan -- returns results without touching live state, so it
        can run on the rescan thread.z,[SCAN] Universe size after asset filtering: )Úbaseline_outÚprev_day_high_outÚprev_close_outz[SCAN] Prefiltered to z symbols)Údaily_atr_outÚ	ref5d_out)ÚprefilteredÚvolume_baselinesÚprev_day_highsÚprev_closesÚrange_20d_highsr™   )rn   Ú	baselinesÚ	daily_atrÚref5d)
ÚuniverseÚget_universe_symbolsrd   r*   r+   ÚlenÚprefilter_by_snapshotÚscannerÚcompute_range_20d_highÚscan)r   r™   Úall_symbolsr¤   r¡   r¢   rŸ   r¥   r¦   r£   rn   r   r   r   Ú_compute_scané   s,   ü
ÿûz!SessionOrchestrator._compute_scanÚresultc              	   C   sÞ   | j  |d ¡ |d D ]?}|d }| j |d | d¡| d¡| d¡dœ¡}| d	i ¡ |d ¡|d	< | | d
i ¡ |d i ¡¡ || j|d < qdd„ | jD ƒ}|d | _t dt	| jƒ› ddd„ | jD ƒ› �¡ |S )zÒMain-thread only. Swaps in the new candidate list; keeps each
        already-known symbol's levels from the first scan it appeared in
        (a mid-session scan's 'premarket_high' is really the last-6h high).r¤   rn   Úmetricsr6   Úpremarket_highÚprevious_day_highÚrange_20d_high)r²   Úprev_day_highr´   r¥   r¦   c                 S   ó   g | ]}|d  ‘qS ©r6   r   ©r[   Úcr   r   r   r_     ó    z3SessionOrchestrator._apply_scan.<locals>.<listcomp>z[SCAN] z candidates selected: c                 S   r¶   r·   r   r¸   r   r   r   r_     rº   )
rq   Úupdaterp   Ú
setdefaultr>   ro   rn   r*   r+   r©   )r   r°   r¹   rK   ÚlvÚold_symsr   r   r   Ú_apply_scan  s"   ý
ÿzSessionOrchestrator._apply_scanc                 C   s   t   ¡ | _|  |  ¡ ¡ d S r   )r–   rv   r¿   r¯   r.   r   r   r   r�     s   
zSessionOrchestrator._run_scanc              
      s"  ˆj d  dd¡}|sdS ˆj� ˆjd}ˆ_W d  ƒ n1 s"w   Y  |durŠˆ |¡‰dd„ ˆjD ƒ‰tˆj ¡ ƒ‰‡fdd„ˆD ƒ}‡‡‡fdd„ˆj	j
D ƒ‰ t d	t|ƒ› d
|› dtˆ ƒ› d
ˆ › �¡ ˆj	 |¡ ˆj	 ˆ ¡ ‡ fdd„ˆjD ƒD ]	}ˆj |d¡ q€t ¡ r�dS ˆjdurœˆj ¡ rœdS ˆj d }t ¡ }| dd¡}| dd¡}|rÄd|  kr½dk rÄn nt||ƒ}|oÍˆj oÍ||k}	|	rÔdˆ_nˆjrät ¡ ˆj |d k rädS t ¡ ˆ_dt ¡  d¡ ‰‡‡fdd„}
t d|› d�¡ tj|
dd�ˆ_ˆj ¡  dS )zêRefresh the candidate list every scanner.rescan_interval_minutes
        (0 = off) until no_new_entries_after. Dropped symbols are
        unsubscribed unless a position is open in them; new ones are
        backfilled and subscribed.r«   Úrescan_interval_minutesr   Nc                 S   r¶   r·   r   r¸   r   r   r   r_   *  rº   z5SessionOrchestrator._maybe_rescan.<locals>.<listcomp>c                    ó   g | ]}|ˆ vr|‘qS r   r   ©r[   Ús)r¾   r   r   r_   ,  ó    c                    s,   g | ]}|ˆ vr|ˆvr|ˆ  ¡ vr|‘qS r   )Ú_benchmarksrÂ   )Únew_symsÚ	open_symsr   r   r   r_   -  s    ÿz[RESCAN] applied: +Ú z / -c                    s   g | ]
}|d  ˆ v r|‘qS )é   r   )r[   Úk)Údroppedr   r   r_   2  s    Úfirst_rescan_after_open_minutesÚ"rescan_interval_first_hour_minutesé<   TÚ_z%H%Mc                     sb   zˆ j ˆd�} W n ty   t d¡ Y d S w ˆ j� | ˆ _W d   ƒ d S 1 s*w   Y  d S )N)r™   z4[RESCAN] scan failed; keeping current candidate list)r¯   r,   r*   r-   r|   ry   )Úres)r   Úsuffixr   r   ÚworkerK  s   
þ"ÿz1SessionOrchestrator._maybe_rescan.<locals>.workerz%[RESCAN] starting background rescan (z-min interval))ÚtargetÚdaemon) r?   r>   r|   ry   r¿   rn   r   re   Úget_open_symbolsrf   Ú_subscribedr*   r+   r©   Úadd_symbolsÚremove_symbolsrr   ÚpoprŠ   Úis_new_entries_cutoffrx   Úis_aliveÚminutes_since_openÚminrw   rv   r–   Únow_etÚstrftimerz   ÚThreadÚstart)r   Úintervalr°   Úaddedr<   ÚscÚ	mins_openÚearlyÚ
first_hourÚ	due_earlyrÒ   r   )rË   rÆ   r¾   rÇ   r   rÑ   r   Ú_maybe_rescan  sL   ÿ
*


	z!SessionOrchestrator._maybe_rescanc                 C   s   t | j di ¡ dg ¡ƒS )NÚ	streamingÚbenchmark_symbols)Úlistr?   r>   r.   r   r   r   rÅ   X  rA   zSessionOrchestrator._benchmarksc                    sb   dd„ | j D ƒ‰ ˆ ‡ fdd„|  ¡ D ƒ7 ‰ t|  ¡ ƒdd„ | j D ƒ | j_ˆ r/| j ˆ ¡ d S d S )Nc                 S   r¶   r·   r   r¸   r   r   r   r_   \  rº   z8SessionOrchestrator._start_streaming.<locals>.<listcomp>c                    rÁ   r   r   ©r[   Úb©Úsymbolsr   r   r_   ]  rÄ   c                 S   s   h | ]}|d  ’qS r·   r   r¸   r   r   r   Ú	<setcomp>^  rº   z7SessionOrchestrator._start_streaming.<locals>.<setcomp>)rn   rÅ   r   rf   Ú	bars_onlyrá   r.   r   rï   r   r�   [  s    ÿz$SessionOrchestrator._start_streamingc                 C   r“   )Né   )rŠ   Úis_past_market_openru   r–   r—   r.   r   r   r   rŽ   b  r˜   z)SessionOrchestrator._wait_for_market_openc                 C   s¤   | j d d }| jsLt ¡ sNz!| jD ]}| ¡  q|  ¡  |  ¡  | j 	¡ r.t 
¡ s.|  ¡  W n ty=   t d¡ Y nw t |¡ | jsPt ¡ rd S d S d S d S )NÚscheduleÚpoll_interval_secondsz$[LOOP] poll cycle failed; continuing)r?   ru   rŠ   Úis_force_liquidate_timerm   r2   ré   Ú_update_open_positionsre   Úhas_available_slotrÚ   Ú_scan_for_entriesr,   r*   r-   r–   r—   )r   râ   rK   r   r   r   r�   g  s   

€ÿ
óz&SessionOrchestrator._main_trading_loopc                 C   sT   |   ¡  t ¡ ¡}dd„ | jd d  d¡D ƒ\}}}||j|||dd�  ¡ d S )	Nc                 s   s   � | ]}t |ƒV  qd S r   )Úint)r[   Úxr   r   r   Ú	<genexpr>|  ó   € z:SessionOrchestrator._minutes_since_open.<locals>.<genexpr>rõ   Úmarket_open_timeú:r   )ÚhourÚminuteÚsecondÚmicrosecondg      N@)rg   Ú
astimezonerŠ   Ú_tzr?   ÚsplitÚreplaceÚtotal_seconds)r   rÞ   ÚhrK   Úsecr   r   r   Ú_minutes_since_openz  s   $z'SessionOrchestrator._minutes_since_openr6   ÚbarsÚmodulec                    sÈ   t ˆ j |i ¡ƒ}tdd„ |D ƒƒ|d< ˆ j |¡}t|ˆ  ¡ |d d |ˆ j |¡|r4|d |d fnd |ˆ j	 |¡ˆ j
 |i ¡ˆ  ¡ ˆ j ¡ ||rO| ¡ ni ‡ fdd	„ˆ  ¡ D ƒ|f‡ fd
d„	d�S )Nc                 s   s   � | ]}|d  V  qdS )r
  Nr   rí   r   r   r   rý   �  rþ   z,SessionOrchestrator._view.<locals>.<genexpr>Úsession_highéÿÿÿÿr¹   r   rÉ   c                    s   i | ]	}|ˆ j  |¡“qS r   )rf   Úget_barsrí   r.   r   r   Ú
<dictcomp>‘  s    z-SessionOrchestrator._view.<locals>.<dictcomp>c                    s   ˆ j j|| d�S )N)Úwindow_seconds)rf   Úget_trade_imbalance)ÚwrÃ   r.   r   r   rY   ’  s    z+SessionOrchestrator._view.<locals>.<lambda>)r6   rW   Úpricer  Úbars_subÚquoteÚlevelsÚvolume_baseliner­   rÜ   Ú
slots_freeÚ
last_trader?   Ú
benchmarksÚ
_imbalance)rH   rp   r>   Úmaxrf   Ú	get_quoter
   rg   Úget_bars_subrq   ro   r  re   rù   r?   rÅ   )r   r6   r  r  r  r  Úqr   r.   r   Ú_view  s(   


ñzSessionOrchestrator._viewFÚkindÚstateÚreasonÚrecordc           	      C   sN   ||f}||f}|s| j  |¡|krdS || j |< |  ||||dœ|¥¡ dS )zDecision log (data/decisions/<date>.jsonl), written when a
        symbol's (state, reason) changes -- not every 5-second poll.N)r6   r$  r%  r&  )rt   r>   rj   )	r   r$  r6   r%  r&  r'  Úforcer<   Úsigr   r   r   Ú_log_decision•  s   

ÿ
ÿz!SessionOrchestrator._log_decisionr  c                 C   sf   t | jj |i ¡ƒ}| j |||¡ | j |¡s/| j |d ¡ | jD ]}|j	d||||d� q"d S d S )NÚon_exitr·   )
rH   re   Ú	positionsr>   Úexit_positionÚis_symbol_openrs   rÙ   rm   r=   )r   r6   r  r&  ÚpositionrK   r   r   r   Ú_sell   s   
ýzSessionOrchestrator._sellc           	   
   C   sV  | j  ¡ D ]£}| j |¡}|sq| j j| }|d d }| j}d }|d ur\| ¡ r\|  |||¡}|jd|t	|ƒ| j
 |i ¡|d�}|d ur\t|tƒs\t d|j› dt|ƒj› d�¡ d }|d u r‰| d¡}|d urˆ||krˆt d	|› d
|d›d|d›�¡ |  ||d|d›�¡ q| jd||j|jd|ji|jd� |jr¨|  |||jp¦|j¡ qd S )Nr  r¹   Úevaluater·   r/   ú.evaluate returned z, not ExitDecision -- ignoredÚ
stop_pricez[EXIT] z> core fallback: no exit decision from the exit module, price $ú.2fz <= entry stop $z.core fallback stop (exit module unavailable) $Úexitr±   ©r(  )re   rÕ   rf   r  r,  ra   r@   r#  r=   rH   rs   r¼   Ú
isinstancer   r*   r‰   r   ÚtyperB   r>   Úwarningr0  r*  r%  r&  r±   Úshould_exit)	r   r6   r  Úpr  ÚemÚdecisionÚviewÚstopr   r   r   rø   ¨  sB   ÿ 
ÿÿÿ€âz*SessionOrchestrator._update_open_positionsc                 C   s  i }| j  ¡ D ]}|||d < q| jD ]í}|d }| j  |¡r q| j  ¡ s( d S | j |¡}|s1q| jD ]Ë}| ¡ s;q4| j	 
|j|fi ¡}| j|||| |¡d�}|jd|||d�}	|	d u r_q4t|	tƒsut d|j› dt|	ƒj› d�¡ q4| jd|j› �||	j|	j|	j|	j|	jd	œ|	jd
� |	js’q4|	jd u s�|	j|jk s³t d|› d|j› d|	j› d|j› d�	¡ q4t|j d|jƒ}
| j j!||j|	j|	jpÉ|	jg|
|	jd�}|rÿ| jD ]}| j	 "|j|fd ¡ qÔi | j#|< t$| j j% |i ¡ƒ}| j&D ]}|jd|||d� qó qd S )Nr6   )r  r1  r·   r/   r2  z, not EntryDecision -- ignoredzentry:)ÚreasonsÚplanr±   r6  z[ENTRY] rÈ   z% wanted to buy with an invalid stop (z
 vs price z) -- skippedÚNAME)r@  ÚsetuprA  Úon_entry)'re   Úclosed_trades_todayrn   r.  rù   rf   r  r`   r@   rr   r¼   r   r#  r>   r=   r7  r   r*   r‰   r8  rB   r*  r%  r&  r@  rA  r±   Úshould_enterr?  r  r9  r   r   Úenter_positionrÙ   rs   rH   r,  rm   )r   r  Útr¹   r6   r  rK   r%  r>  ÚdrC  ÚenteredÚmmÚposr   r   r   rú   É  sd   



 þÿ
ÿÿ


€Øz%SessionOrchestrator._scan_for_entriesc                 C   sn  t  d¡ | jD ]}| d¡ q| jd d }| jd d }t ¡ | }t ¡ |k r_| j ¡ }|s2nF|D ]}| j 	|¡}|rD|d d n| jj
| d }|  ||d	¡ q4t |¡ t ¡ |k s*| j ¡ }	|	rxt  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 )Nz*[EOD] Force-liquidating all open positionsÚon_session_endrõ   Ú eod_liquidation_max_wait_secondsrö   r  r¹   Úentry_priceÚ
END_OF_DAYz[EOD] z& position(s) still not resolved after zs: c                 s   s*   � | ]}|  d ¡dkr|  dd¡V  qdS )ÚstatusÚclosedÚ
pl_dollarsr   N©r>   ©r[   rH  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 ©rQ  rR  rS  r   rÉ   NrT  rU  r   r   r   rý     ó   €0 c                 s   s2   � | ]}|  d ¡dkr|  dd¡dkrdV  qdS rV  rT  rU  r   r   r   rý     rW  z[EOD] Session summary: trades=z wins=z losses=z total_P/L=$r4  )r*   r+   rm   r=   r?   r–   re   rÕ   rf   r  r,  r0  r—   r9  r©   Úsortedr?  rh   Úload_today_tradesÚwrite_trades_summaryÚsum)r   rK   Úmax_waitÚpoll_intervalÚdeadlineÚopen_symbolsr6   r  Úcurrent_priceÚ
still_openÚtradesÚtotal_plÚwinsÚlossesr   r   r   r‘   ù  s@   


 
ø

ÿÿ

ÿz+SessionOrchestrator._end_of_day_liquidation)NNNNNT)r$   r   )F)rB   rC   rD   rE   r   r   r’   rŒ   rF   rH   r¯   r¿   r�   ré   rì   rÅ   r�   rŽ   r�   Úfloatr  r   r
   r#  r*  r0  rø   rú   r‘   r   r   r   r   rV   —   s.    
ÿ+; !0rV   c                  C   s  zot  ¡ �` t  ¡  ttdƒ�} |  tt ¡ ƒ¡ W d   ƒ n1 s#w   Y  zt	ƒ }| 
¡  W zt t¡ W n tyA   Y nw zt t¡ W w  tyR   Y w w W d   ƒ W d S W d   ƒ W d S 1 shw   Y  W d S  tyŽ } zt t|ƒ¡ t d¡ W Y d }~d S d }~ww )Nr  rÉ   )rh   Úacquire_singleton_lockÚensure_dirsÚopenÚPID_PATHÚwriterF   r   ÚgetpidrV   r’   Úremover0   ÚRuntimeErrorr*   r‰   Úsysr5  )r!   ÚorchestratorÚer   r   r   Úmain  s8   
ÿ
ÿþÿÿ÷&õ€þrr  Ú__main__))rE   r%   r   ro  r–   r}   rz   r   r   Úconfig_loaderr   r   Úlogger_setupr   rŠ   rh   r§   r«   Úposition_managerr   rf   r   Úalpaca_clientr	   Ú	rules_apir
   r   r   r*   r   ÚdirnameÚabspathr   r    r   rj  r   rI   rV   rr  rB   r   r   r   r   Ú<module>   s>     C   
ÿ