o
    Tc£jÀ@  ã                   @   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	 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G dd„ dƒZdS )a7  
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 fast_pipeline.py and position_manager.py read from.

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)Ú
get_config)Ú
get_logger)Ú
get_client)Úsession_open_dtÚstreamc                   @   s¬   e Zd ZdZd%dd„Zdefdd„Zdefd	d
„Zdd„ Zdd„ Z	d&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e
fdd „Zd!edefd"d#„Zd$S )'Ú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   sD   || _ tƒ | _d | _d | _d | _d | _|| _|| _tƒ | _	d | _
d S ©N)Úmaxlenr   ÚbarsÚlatest_quoteÚlatest_trade_priceÚlast_updateÚ_current_barÚsub_bucket_secondsÚ
sub_maxlenÚbars_subÚ_current_bar_sub)Úselfr   r   r   © r   ú	stream.pyÚ__init__%   s   
zSymbolBuffer.__init__Úbarc                 C   óT   |d }|| j |< | j  |¡ t| j ƒ| jkr(| j jdd� t| j ƒ| jksd S d S ©NÚtF)Úlast)r   Úmove_to_endÚlenr   Úpopitem©r   r   Úkeyr   r   r   Ú_upsert_barR   ó   
ÿzSymbolBuffer._upsert_barc                 C   r   r    )r   r#   r$   r   r%   r&   r   r   r   Ú_upsert_bar_subY   r)   zSymbolBuffer._upsert_bar_subc                 C   s"   | j }|j| | }|j|dd�S )Nr   ©ÚsecondÚmicrosecond)r   r,   Úreplace)r   ÚtsÚ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.
        Nr!   r   )	r!   ÚoÚhÚlÚcÚvÚ
tick_countÚupticksÚ	downticksr4   r5   r6   r7   r8   é   r9   r:   )r2   r   r*   ÚmaxÚmin)r   ÚpriceÚsizer/   Ú
prev_priceÚbucketÚis_new_bucketÚbr   r   r   Ú_accumulate_bar_sube   s,   

þüz SymbolBuffer._accumulate_bar_subTÚreturnc                 C   ó,   t | j ¡ ƒ}|r| jd ur|| jg }|S r   )Úlistr   Úvaluesr   ©r   Úinclude_formingr   r   r   r   Úget_bars_subˆ   ó   zSymbolBuffer.get_bars_subr>   r?   c                 C   s4   | j }|| _ || _|  |||¡ |  ||||¡ d S r   )r   r   Ú_accumulate_barrD   )r   r>   r?   r/   r@   r   r   r   Úon_tradeŽ   s
   zSymbolBuffer.on_tradeÚbidÚaskc                 C   s   |||f| _ || _d S r   )r   r   )r   rO   rP   r/   r   r   r   Úon_quote•   s   
zSymbolBuffer.on_quotec              	   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 )zBDirect minute-bar ingestion (Alpaca also streams aggregated bars).r.   r   r+   ©r!   r3   r4   r5   r6   r7   Nr!   )Úhasattrr.   r(   r   r   )r   r3   r4   r5   r6   r7   r/   Úminute_bucketr   r   r   Úon_bar™   s   
ÿzSymbolBuffer.on_barc                 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   r+   r!   rR   r4   r5   r6   r7   )r.   r   r(   r<   r=   )r   r>   r?   r/   rT   rC   r   r   r   rM   ¥   s   
ÿzSymbolBuffer._accumulate_barc                 C   rF   r   )rG   r   rH   r   rI   r   r   r   Úget_bars´   rL   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   r.   Útotal_seconds)r   rW   rY   r"   r   r   r   Úis_staleº   s   

zSymbolBuffer.is_staleN)r   r   r   )T)Ú__name__Ú
__module__Ú__qualname__Ú__doc__r   Údictr(   r*   r2   rD   rG   rK   ÚfloatrN   rQ   rU   rM   rV   ÚintÚboolr\   r   r   r   r   r      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d„ Z	dd„ Z
dd„ Z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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	ƒ | _
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   rk   r   r   Ú<lambda>Ì   s    z(StreamManager.__init__.<locals>.<lambda>)r   Úcfgr	   ÚclientÚgetr<   rc   r   ÚbuffersÚsetÚ_subscribedÚ_streamÚ_threadÚ_loopÚ	threadingÚEventÚ_stop_eventÚ
_connected)r   Úsub_buffer_minutesr   rk   r   r   Å   s   ÿ
zStreamManager.__init__Úsymbolsc           
      C   s  t ƒ }t tj¡}|| tj¡krdS |D ]r}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œ¡ qJt	 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 z: r   r+   rR   z[STREAM] Backfilled z with z  historical bars (session open $z.4fz) before live subscription)r
   r   rY   r   rZ   Ú
astimezonern   Úget_minute_barsÚ	ExceptionÚlogÚwarningrp   Ú	timestampr.   r(   rb   ÚopenÚhighÚlowÚcloseÚvolumeÚinfor$   )
r   r{   Úopen_dtrY   Úsymbolr   ÚeÚbufr   rT   r   r   r   Ú	_backfillÖ   s2   €þ

þÿñzStreamManager._backfillc                 C   s<   |   |¡ tj| j|fdd�| _| j ¡  | jjdd� d S )NT)ÚtargetÚargsÚdaemoné
   )Útimeout)rŽ   rv   ÚThreadÚ_run_foreverrt   r|   ry   Úwait©r   r{   r   r   r   r|     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   )rx   rq   ru   rs   ÚasyncioÚrun_coroutine_threadsafeÚstop_wsr€   ©r   r   r   r   Ústop  s   
ÿýzStreamManager.stopc              
      s–   ‡ fdd„|D ƒ}|rˆ j d u sˆ jd u rd S ˆ  |¡ ˆ j |¡ zt ˆ  |¡ˆ j¡ W d S  tyJ } zt	 
d|› �¡ W Y d }~d S d }~ww )Nc                    s   g | ]	}|ˆ j vr|‘qS r   )rr   )Ú.0Úsr›   r   r   Ú
<listcomp>  s    z-StreamManager.add_symbols.<locals>.<listcomp>zadd_symbols failed: )rs   ru   rŽ   rr   Úupdater˜   r™   Ú_subscribe_asyncr€   r�   r‚   )r   r{   Únew_symsrŒ   r   r›   r   Úadd_symbols  s   

ÿ€ÿzStreamManager.add_symbolsc                 Ã   sH   �| j j| jg|¢R Ž  | j j| jg|¢R Ž  | j j| jg|¢R Ž  d S r   )rs   Úsubscribe_tradesÚ	_on_tradeÚsubscribe_quotesÚ	_on_quoteÚsubscribe_barsÚ_on_barr—   r   r   r   r¡   !  s   €zStreamManager._subscribe_asyncc                 Ã   s.   �| j |j }| t|jƒt|jƒ|j¡ d S r   )rp   r‹   rN   rb   r>   r?   rƒ   )r   Útrader�   r   r   r   r¥   '  s   € zStreamManager._on_tradec                 Ã   sB   �| j |j }|jr|jr| t|jƒt|jƒ|j¡ d S d S d S r   )rp   r‹   Ú	bid_priceÚ	ask_pricerQ   rb   rƒ   )r   Úquoter�   r   r   r   r§   +  s
   € ÿzStreamManager._on_quotec                 Ã   sF   �| j |j }| t|jƒt|jƒt|jƒt|jƒt|jƒ|j	¡ d S r   )
rp   r‹   rU   rb   r„   r…   r†   r‡   rˆ   rƒ   )r   r   r�   r   r   r   r©   0  s
   €ÿzStreamManager._on_barc                 C   s‚  | j d }| j d }d}| j ¡ s´||k r´zPt ¡ | _t | j¡ | j ¡ | _	t
|ƒ| _| 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 ¡ rxW 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   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.)rm   rx   Úis_setr˜   Únew_event_loopru   Úset_event_looprn   Ú
new_streamrs   rq   rr   r¤   r¥   r¦   r§   r¨   r©   r�   r‰   r$   ry   Úrunr€   r=   r‚   ÚtimeÚsleepÚerror)r   r{   ÚbackoffsÚmax_attemptsÚattemptrŒ   Údelayr   r   r   r•   6  sB   





ÿ
ÿ
€ùòÿzStreamManager._run_foreverr‹   rE   c                 C   ó   | j |  ¡ S r   )rp   rV   ©r   r‹   r   r   r   rV   W  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.)rp   rK   r¿   r   r   r   rK   Z  s   zStreamManager.get_bars_subc                 C   s   | j | jS r   )rp   r   r¿   r   r   r   Ú	get_quote`  s   zStreamManager.get_quotec                 C   s   | j |  | jd ¡S )NÚstale_data_seconds)rp   r\   rm   r¿   r   r   r   Úis_symbol_stalec  ó   zStreamManager.is_symbol_stalec                 C   s   | j  ¡ o
| j ¡  S r   )ry   r²   rx   r›   r   r   r   Ú
is_healthyf  rÃ   zStreamManager.is_healthyN)r]   r^   r_   r   rG   rŽ   r|   rœ   r£   r¡   r¥   r§   r©   r•   ÚstrrV   rK   rÀ   rd   rÂ   rÄ   r   r   r   r   re   Ä   s     /!re   )r`   r˜   rv   r·   Úcollectionsr   r   r   r   r   Úconfig_loaderr   Úlogger_setupr   Úalpaca_clientr	   Úmarket_timer
   r�   r   re   r   r   r   r   Ú<module>   s     &