o
    A%¸jp_  ã                   @   sª   d 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mZ ddlmZ ddlmZ ed	ƒZG d
d„ dƒZdefdd„ZG dd„ dƒZdS )aJ  
stream.py

Manages the Alpaca SIP real-time stream for a set of symbols. Buffers
incoming trades/quotes/bars into an in-memory per-symbol rolling window
that monitor.py reads and hands to the rule modules (rules_api.MarketView).

Runs the stream in a background thread so monitor.py's main loop stays
simple synchronous polling logic (poll the buffer, not the network).

Handles:
- reconnect with backoff (config: streaming.reconnect_backoff_seconds)
- stale data detection (config: streaming.stale_data_seconds)
- symbol subscribe/unsubscribe as the watchlist changes through the day
é    N)ÚdefaultdictÚdequeÚOrderedDict)ÚdatetimeÚtimezoneÚ	timedelta)Ú
get_config)Ú
get_logger)Ú
get_client)Úsession_open_dtÚstreamc                   @   s  e Zd ZdZd4dd„ZdZdedefd	d
„Zdd„ Z	d5defdd„Z
defdd„Zdefdd„Zdd„ Zdd„ Zd6defdd„Zdedefdd„Zdedefd d!„Zd"ed#efd$d%„Zd"ed#efd&d'„Zd(d)„ Zd*d+„ Zd,d-„ Zd6defd.d/„Zd0edefd1d2„Zd3S )7ÚSymbolBuffera	  Rolling window of 1-min bars + latest quote for one symbol, plus a
    second, parallel sub-minute rolling window (bars_sub / get_bars_sub())
    at a configurable bucket width for finer-grained stream-feature work.
    See config.json's streaming._note_sub_minute.é†  é   c                 C   sV   || _ tƒ | _d | _d | _d | _d | _|| _|| _tƒ | _	d | _
tƒ | _t ¡ | _d S ©N)Úmaxlenr   ÚbarsÚlatest_quoteÚlatest_trade_priceÚlast_updateÚ_current_barÚsub_bucket_secondsÚ
sub_maxlenÚbars_subÚ_current_bar_subr   Ú
_trade_logÚ	threadingÚRLockÚ_lock)Úselfr   r   r   © r    ú"/var/www/screener/trade1/stream.pyÚ__init__%   s   zSymbolBuffer.__init__éZ   ÚpriceÚreturnc                 C   sf   | j du rdS | j \}}}|r|sdS ||krdS ||krdS || d }||kr+dS ||k r1dS dS )z@+1 buyer-initiated, -1 seller-initiated, 0 unknown/no quote yet.Nr   é   éÿÿÿÿé   )r   )r   r$   ÚbidÚaskÚ_Úmidr    r    r!   Ú_classify_tradeg   s   
zSymbolBuffer._classify_tradec                 C   sZ   |t | jd� }| jr'| jd d |k r+| j ¡  | jr)| jd d |k sd S d S d S d S )N©Úsecondsr   )r   Ú_TRADE_LOG_MAX_AGE_SECONDSr   Úpopleft)r   Únow_tsÚcutoffr    r    r!   Ú_trim_trade_logy   s   
(ÿzSymbolBuffer._trim_trade_logç      4@Úwindow_secondsc           
      C   sÂ   | j �H | js	 W d  ƒ dS | jd d }|t|d� }d }}t| jƒD ]\}}}||k r2 n|dkr;||7 }q'|dk rC||7 }q'W d  ƒ n1 sNw   Y  || }	|	r_|| |	 S dS )zî(buy_vol - sell_vol) / (buy_vol + sell_vol) over the last
        window_seconds of classified trades, or None if there's nothing
        to compute it from (fails open, same convention as spread_pct
        when a quote isn't available).Nr'   r   r.   g        )r   r   r   Úreversed)
r   r6   r2   r3   Úbuy_volÚsell_volÚtsÚsideÚsizeÚtotalr    r    r!   Úget_trade_imbalance~   s$   þ
€ôz SymbolBuffer.get_trade_imbalanceÚbarc                 C   sˆ   | j �7 |d }|| j|< | j |¡ t| jƒ| jkr2| jjdd� t| jƒ| jksW d   ƒ d S W d   ƒ d S 1 s=w   Y  d S ©NÚtF)Úlast)r   r   Úmove_to_endÚlenr   Úpopitem©r   r?   Úkeyr    r    r!   Ú_upsert_bar“   s   
ÿü"üzSymbolBuffer._upsert_barc                 C   sT   |d }|| j |< | j  |¡ t| j ƒ| jkr(| j jdd� t| j ƒ| jksd S d S r@   )r   rC   rD   r   rE   rF   r    r    r!   Ú_upsert_bar_sub›   s   
ÿzSymbolBuffer._upsert_bar_subc                 C   s"   | j }|j| | }|j|dd�S )Nr   ©ÚsecondÚmicrosecond)r   rK   Úreplace)r   r:   ÚwidthÚbucket_secondr    r    r!   Ú_sub_bucket_for¢   s   zSymbolBuffer._sub_bucket_forc              
   C   sô   |   |¡}| jdu p| jd |k}|r+| jdur|  | j¡ ||||||ddddœ	| _| j}|sNt|d |ƒ|d< t|d |ƒ|d< ||d< |d  |7  < |d	  d
7  < |durv||krh|d  d
7  < dS ||k rx|d  d
7  < dS dS dS )aÎ  
        [SIMULATION-ONLY 2026-09-05] prev_price is the trade price
        immediately before this one (captured by on_trade() BEFORE it
        overwrites self.latest_trade_price) -- compared against the new
        price for tick direction (upticks/downticks), deliberately
        continuous across bucket boundaries rather than resetting
        direction tracking to "unknown" at the start of every new bucket,
        so a bucket's uptick_ratio isn't artificially diluted by treating
        its first tick as directionless. tick_count/upticks/downticks are
        plain O(1) counters -- no raw ticks are stored, so this adds no
        unbounded memory regardless of how many trades print per bucket.
        NrA   r   )	rA   ÚoÚhÚlÚcÚvÚ
tick_countÚupticksÚ	downticksrR   rS   rT   rU   rV   r&   rW   rX   )rP   r   rI   ÚmaxÚmin)r   r$   r<   r:   Ú
prev_priceÚbucketÚis_new_bucketÚbr    r    r!   Ú_accumulate_bar_sub§   s,   

þüz SymbolBuffer._accumulate_bar_subTc                 C   óx   | j �/ t| j ¡ ƒ}|r"| jd ur*|t| jƒg }W d   ƒ |S W d   ƒ |S W d   ƒ |S 1 s5w   Y  |S r   )r   Úlistr   Úvaluesr   Údict©r   Úinclude_formingr   r    r    r!   Úget_bars_subÊ   ó   
ýþ
þþ
þüzSymbolBuffer.get_bars_subr<   c                 C   ó:   | j � |  |||¡ W d   ƒ d S 1 sw   Y  d S r   )r   Ú_on_trade_locked)r   r$   r<   r:   r    r    r!   Úon_tradeÑ   ó   "ÿzSymbolBuffer.on_tradec                 C   sZ   | j }|| _ || _|  |||¡ |  ||||¡ |  |¡}| j |||f¡ |  |¡ d S r   )r   r   Ú_accumulate_barr_   r-   r   Úappendr4   )r   r$   r<   r:   r[   r;   r    r    r!   ri   Õ   s   
zSymbolBuffer._on_trade_lockedr)   r*   c                 C   rh   r   )r   Ú_on_quote_locked©r   r)   r*   r:   r    r    r!   Úon_quoteß   rk   zSymbolBuffer.on_quotec                 C   s   |||f| _ || _d S r   )r   r   ro   r    r    r!   rn   ã   s   
zSymbolBuffer._on_quote_lockedc              	   C   s@   | j � |  ||||||¡ W d  ƒ dS 1 sw   Y  dS )zBDirect minute-bar ingestion (Alpaca also streams aggregated bars).N)r   Ú_on_bar_locked)r   rQ   rR   rS   rT   rU   r:   r    r    r!   Úon_barç   s   "ÿzSymbolBuffer.on_barc              	   C   sd   t |dƒr|jddd�n|}|  ||||||dœ¡ || _| jd ur.| jd |kr0d | _d S d S d S )NrM   r   rJ   ©rA   rQ   rR   rS   rT   rU   rA   )ÚhasattrrM   rH   r   r   )r   rQ   rR   rS   rT   rU   r:   Úminute_bucketr    r    r!   rq   ì   s   
ÿzSymbolBuffer._on_bar_lockedc                 C   sš   |j ddd�}| jd u s| jd |kr*| jd ur|  | j¡ ||||||dœ| _d S | j}t|d |ƒ|d< t|d |ƒ|d< ||d< |d  |7  < d S )	Nr   rJ   rA   rs   rR   rS   rT   rU   )rM   r   rH   rY   rZ   )r   r$   r<   r:   ru   r^   r    r    r!   rl   ÷   s   
ÿzSymbolBuffer._accumulate_barc                 C   r`   r   )r   ra   r   rb   r   rc   rd   r    r    r!   Úget_bars  rg   zSymbolBuffer.get_barsÚstale_secondsc                 C   sH   | j d u rdS t tj¡}| j }|jd u r|jtjd�}||  ¡ |kS )NT)Útzinfo)r   r   Únowr   Úutcrx   rM   Útotal_seconds)r   rw   ry   rB   r    r    r!   Úis_stale  s   

zSymbolBuffer.is_staleN)r   r   r   ©r5   )T)Ú__name__Ú
__module__Ú__qualname__Ú__doc__r"   r0   ÚfloatÚintr-   r4   r>   rc   rH   rI   rP   r_   ra   rf   rj   ri   rp   rn   rr   rq   rl   rv   Úboolr|   r    r    r    r!   r      s*    
@#
r   r%   c                 C   s2   | j t| jƒt| jƒt| jƒt| jƒt| jƒdœS )Nrs   )Ú	timestampr‚   ÚopenÚhighÚlowÚcloseÚvolume)r^   r    r    r!   Ú	_bar_dict  s   ÿr‹   c                   @   sî   e Zd Zdd„ Zdefdd„Zdefdd„Z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efdd„Zdefdd„Zdedefdd„Zdedefdd„Zdefdd „Zd*ded"efd#d$„Zdedefd%d&„Zdefd'd(„Zd)S )+ÚStreamManagerc                    sš   t ƒ d | _tƒ | _| j dd¡‰ | j dd¡}tdt|d ˆ  ƒƒ‰t‡ ‡fdd	„ƒ| _t	ƒ | _
g | _t	ƒ | _d | _d | _d | _t ¡ | _t ¡ | _d S )
NÚ	streamingÚsub_minute_bucket_secondsr   Úsub_minute_buffer_minutesé   r&   é<   c                      s   t ˆ ˆd�S )N©r   r   )r   r    r’   r    r!   Ú<lambda>$  s    z(StreamManager.__init__.<locals>.<lambda>)r   Úcfgr
   ÚclientÚgetrY   rƒ   r   ÚbuffersÚsetÚ_subscribedÚ	listenersÚ	bars_onlyÚ_streamÚ_threadÚ_loopr   ÚEventÚ_stop_eventÚ
_connected)r   Úsub_buffer_minutesr    r’   r!   r"     s    ÿ
zStreamManager.__init__Úsymbolsc           
      C   s,  t ƒ }t tj¡}|| tj¡krdS |D ]~}z| jj|||d�}W n ty? } zt	 
d|› d|› �¡ W Y d}~qd}~ww |sCq| j| }|D ]%}|jjddd�}	| |	t|jƒt|jƒt|jƒt|jƒt|jƒdœ¡ qJ|  d|d	d
„ |D ƒ¡ t	 d|› dt|ƒ› dt|d jƒd›d�¡ qdS )a´  
        [BUGFIX 2026-09-09] stream_features.compute_features() treats
        bars[0] as "the session's first 1-min bar" (session_open in its
        extension_from_open_pct calc) on the assumption a symbol's
        buffer always spans the whole session from 09:30. That's only
        true if the symbol has been subscribed since before the open --
        false for (a) any symbol added mid-day via add_symbols() (no
        history, subscribe_trades/bars only deliver ticks going
        forward), and (b) EVERY symbol after any intraday process
        restart, since SymbolBuffer is in-memory only and starts empty
        again. Confirmed root cause of a real miss: SRAD, added to the
        pool at 12:13 on 2026-09-09 (well after its 11:55 reversal low
        and post a 12:27 restart besides), had its "session open" read
        as ~$12.65 (the price when the buffer happened to start) instead
        of the real $12.69 09:30 open -- silently understating how
        extended it actually was, undermining the whole extension gate.
        Fix: before subscribing a symbol to the live stream, pull its
        real 1-min bars since market_open_time via REST and pre-load the
        buffer, so bars[0] is always genuinely the session's first bar
        regardless of when/why the subscription started. No-op before
        the open (nothing to backfill yet -- the live stream will
        deliver the real bar 0 itself, same as today).
        N)ÚstartÚendz[STREAM] Backfill failed for ú: r   rJ   rs   Ú	seed_barsc                 S   s   g | ]}t |ƒ‘qS r    )r‹   )Ú.0r^   r    r    r!   Ú
<listcomp>d  s    z+StreamManager._backfill.<locals>.<listcomp>z[STREAM] Backfilled z with z  historical bars (session open $z.4fz) before live subscription)r   r   ry   r   rz   Ú
astimezoner•   Úget_minute_barsÚ	ExceptionÚlogÚwarningr—   r…   rM   rH   r‚   r†   r‡   rˆ   r‰   rŠ   Ú_notifyÚinforD   )
r   r£   Úopen_dtry   Úsymbolr   ÚeÚbufr?   ru   r    r    r!   Ú	_backfill9  s4   €þ

þÿðzStreamManager._backfillc                 C   s<   |   |¡ tj| j|fdd�| _| j ¡  | jjdd� d S )NT)ÚtargetÚargsÚdaemoné
   )Útimeout)rµ   r   ÚThreadÚ_run_foreverr�   r¤   r¡   Úwait)r   r£   r    r    r!   r¤   i  s   

zStreamManager.startc                 C   sN   | j  ¡  | jr#| jr%zt | j ¡ | j¡ W d S  ty"   Y d S w d S d S r   )r    r˜   rž   rœ   ÚasyncioÚrun_coroutine_threadsafeÚstop_wsr¬   ©r   r    r    r!   Ústopp  s   
ÿýzStreamManager.stopc              
      sx  ‡ fdd„|D ƒ}|rOˆ j durOˆ  jt|ƒ8  _zˆ j jˆ jg|¢R Ž  ˆ j jˆ jg|¢R Ž  W n tyN } zt 	d|› d|› �¡ W Y d}~nd}~ww ‡ fdd„|D ƒ}|r_ˆ j du radS ˆ  
|¡ ˆ j |¡ z1ˆ j jˆ jg|¢R Ž  ˆ j jˆ jg|¢R Ž  ˆ j jˆ jg|¢R Ž  t dt|ƒ› d|› �¡ W dS  ty» } zt 	d	|› d|› �¡ W Y d}~dS d}~ww )
aß  [FIX 2026-09-23] Used to schedule _subscribe_async() onto
        self._loop -- but self._stream.run() runs on alpaca-py's OWN loop
        (asyncio.run inside run()), so self._loop never ran and the
        coroutine never executed: added symbols were silently never
        subscribed. alpaca-py's subscribe_*/unsubscribe_* are themselves
        thread-safe (they post onto the stream's running loop and wait),
        so they're called directly from the caller's thread here.c                    s$   g | ]}|ˆ j v r|ˆ jv r|‘qS r    )r›   r™   ©r¨   ÚsrÁ   r    r!   r©   ‚  s   $ z-StreamManager.add_symbols.<locals>.<listcomp>Nz#[STREAM] tick subscribe failed for r¦   c                    ó   g | ]	}|ˆ j vr|‘qS r    ©r™   rÃ   rÁ   r    r!   r©   Š  ó    z[STREAM] Subscribed z added symbols: z [STREAM] add_symbols failed for )rœ   r›   r˜   Úsubscribe_tradesÚ	_on_tradeÚsubscribe_quotesÚ	_on_quoter¬   r­   r®   rµ   r™   ÚupdateÚsubscribe_barsÚ_on_barr°   rD   )r   r£   Úpromoter³   Únew_symsr    rÁ   r!   Úadd_symbolsx  s0   
"€ÿ
 $€ÿzStreamManager.add_symbolsc              
      sÊ   ‡ fdd„|D ƒ}|sdS ˆ j  |¡ ˆ jdurJzˆ jj|Ž  ˆ jj|Ž  ˆ jj|Ž  W n tyI } zt d|› d|› �¡ W Y d}~nd}~ww |D ]	}ˆ j	 
|d¡ qLt dt|ƒ› d|› �¡ dS )z«Unsubscribes and drops the buffer for symbols no longer
        watched, so the subscription set stays ~one scanner list wide
        instead of growing with every rescan.c                    s   g | ]	}|ˆ j v r|‘qS r    rÆ   rÃ   rÁ   r    r!   r©   ›  rÇ   z0StreamManager.remove_symbols.<locals>.<listcomp>Nz#[STREAM] remove_symbols failed for r¦   z[STREAM] Unsubscribed z dropped symbols: )r™   Údifference_updaterœ   Úunsubscribe_tradesÚunsubscribe_quotesÚunsubscribe_barsr¬   r­   r®   r—   Úpopr°   rD   )r   r£   Úgoner³   Úsymr    rÁ   r!   Úremove_symbols—  s    
"€ÿzStreamManager.remove_symbolsc                 Ã   sP   �| j |j }| t|jƒt|jƒ|j¡ |  d|j|jt|jƒt|jƒ¡ d S )Nrj   )r—   r²   rj   r‚   r$   r<   r…   r¯   )r   Útrader´   r    r    r!   rÉ   «  s   €&zStreamManager._on_tradec              
   Ã   s|   �| j |j }|jr:|jr<| t|jƒt|jƒ|j¡ |  d|j|jt|jƒt|jƒt|jp/dƒt|j	p5dƒ¡ d S d S d S )Nrp   r   )
r—   r²   Ú	bid_priceÚ	ask_pricerp   r‚   r…   r¯   Úbid_sizeÚask_size)r   Úquoter´   r    r    r!   rË   °  s   € ÿþzStreamManager._on_quotec                 Ã   sZ   �| j |j }| t|jƒt|jƒt|jƒt|jƒt|jƒ|j	¡ |  
d|jt|ƒ¡ d S )Nrr   )r—   r²   rr   r‚   r†   r‡   rˆ   r‰   rŠ   r…   r¯   r‹   )r   r?   r´   r    r    r!   rÎ   ·  s   €ÿzStreamManager._on_barÚmethodc                 G   sŒ   | j D ]@}t||d ƒ}|d u rqz||Ž  W q tyC   t|ddƒsAt dt|ƒj› d|› d�¡ zd|_W n	 ty@   Y nw Y qw d S )NÚ_stream_error_loggedFz[STREAM] listener Ú.z failed (ignored)T)rš   Úgetattrr¬   r­   Ú	exceptionÚtyper~   rá   )r   rà   r·   ÚlstÚfnr    r    r!   r¯   ½  s"   

ÿ€ûúzStreamManager._notifyc           	         s¨  ˆ j d }ˆ j d }d}ˆ j ¡ sÇ||k rÇzct ¡ ˆ _t ˆ j¡ ˆ j ¡ ˆ _	ˆ j
s/t|ƒˆ _
tˆ j
ƒ}‡ fdd„|D ƒ}|rUˆ j	jˆ jg|¢R Ž  ˆ j	jˆ jg|¢R Ž  ˆ j	jˆ jg|¢R Ž  t dt|ƒ› d�¡ ˆ j ¡  d}ˆ j	 ¡  W nE ty½ } z9ˆ j ¡ r‹W Y d }~n<|t|t|ƒd ƒ }t d	|› d
|› d|d › d|› d�	¡ t |¡ |d7 }W Y d }~nd }~ww ˆ j ¡ sÇ||k s||krÒt d¡ d S d S )NÚreconnect_backoff_secondsÚmax_reconnect_attemptsr   c                    rÅ   r    )r›   )r¨   ÚxrÁ   r    r!   r©   Ý  rÇ   z.StreamManager._run_forever.<locals>.<listcomp>z#[STREAM] Connecting SIP stream for z symbolsr&   z[STREAM] Disconnected (z); reconnecting in zs (attempt ú/ú)zf[STREAM] Max reconnect attempts exceeded. Stream is DOWN. monitor.py should fall back to REST polling.)r”   r    Úis_setr¾   Únew_event_looprž   Úset_event_loopr•   Ú
new_streamrœ   r™   r˜   ÚsortedrÈ   rÉ   rÊ   rË   rÍ   rÎ   r­   r°   rD   r¡   Úrunr¬   rZ   r®   ÚtimeÚsleepÚerror)	r   r£   ÚbackoffsÚmax_attemptsÚattemptÚcurrentÚticksr³   Údelayr    rÁ   r!   r¼   Í  sJ   






ÿ
ÿ
€ùëÿzStreamManager._run_foreverr²   r%   c                 C   ó   | j |  ¡ S r   )r—   rv   ©r   r²   r    r    r!   rv   õ  s   zStreamManager.get_barsc                 C   rü   )z½Sub-minute (configurable bucket width, default 30s) rolling
        bars -- see SymbolBuffer.get_bars_sub() / config.json's
        streaming._note_sub_minute. Not consumed by anything yet.)r—   rf   rý   r    r    r!   rf   ø  s   zStreamManager.get_bars_subc                 C   s   | j | jS r   )r—   r   rý   r    r    r!   Ú	get_quoteþ  s   zStreamManager.get_quoter5   r6   c                 C   s   | j |  |¡S r   )r—   r>   )r   r²   r6   r    r    r!   r>     s   z!StreamManager.get_trade_imbalancec                 C   s   | j |  | jd ¡S )NÚstale_data_seconds)r—   r|   r”   rý   r    r    r!   Úis_symbol_stale  ó   zStreamManager.is_symbol_stalec                 C   s   | j  ¡ o
| j ¡  S r   )r¡   rí   r    rÁ   r    r    r!   Ú
is_healthy  r  zStreamManager.is_healthyNr}   )r~   r   r€   r"   ra   rµ   r¤   rÂ   rÑ   rÙ   rÉ   rË   rÎ   Ústrr¯   r¼   rv   rf   rþ   r‚   r>   r„   r   r  r    r    r    r!   rŒ     s$    0(rŒ   )r�   r¾   r   ró   Úcollectionsr   r   r   r   r   r   Úconfig_loaderr   Úlogger_setupr	   Úalpaca_clientr
   Úmarket_timer   r­   r   rc   r‹   rŒ   r    r    r    r!   Ú<module>   s     y