o
    þ¥›j)6  ã                   @   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 edƒZG d	d
„ d
ƒZG dd„ dƒZdS )a6  
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 entry_engine.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Ú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   úB/var/www/screener/trade/premarket_backup_2026-09-08_2010/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_barQ   ó   
ÿzSymbolBuffer._upsert_barc                 C   r   r   )r   r"   r#   r   r$   r%   r   r   r   Ú_upsert_bar_subX   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Ú	downticksr3   r4   r5   r6   r7   é   r8   r9   )r1   r   r)   ÚmaxÚmin)r   ÚpriceÚsizer.   Ú
prev_priceÚbucketÚis_new_bucketÚbr   r   r   Ú_accumulate_bar_subd   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_barrC   )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   rN   rO   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    r2   r3   r4   r5   r6   Nr    )Úhasattrr-   r'   r   r   )r   r2   r3   r4   r5   r6   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    rQ   r3   r4   r5   r6   )r-   r   r'   r;   r<   )r   r=   r>   r.   rS   rB   r   r   r   rL   ¤   s   
ÿzSymbolBuffer._accumulate_barc                 C   rE   r   )rF   r   rG   r   rH   r   r   r   Úget_bars³   rK   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   ÚutcrW   r-   Útotal_seconds)r   rV   rX   r!   r   r   r   Úis_stale¹   s   

zSymbolBuffer.is_staleN)r   r   r   )T)Ú__name__Ú
__module__Ú__qualname__Ú__doc__r   Údictr'   r)   r1   rC   rF   rJ   ÚfloatrM   rP   rT   rL   rU   Ú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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   rj   r   r   Ú<lambda>Ë   s    z(StreamManager.__init__.<locals>.<lambda>)r   Úcfgr	   ÚclientÚgetr;   rb   r   ÚbuffersÚsetÚ_subscribedÚ_streamÚ_threadÚ_loopÚ	threadingÚEventÚ_stop_eventÚ
_connected)r   Úsub_buffer_minutesr   rj   r   r   Ä   s   ÿ
zStreamManager.__init__Úsymbolsc                 C   s2   t j| j|fdd�| _| j ¡  | jjdd� d S )NT)ÚtargetÚargsÚdaemoné
   )Útimeout)ru   ÚThreadÚ_run_foreverrs   Ústartrx   Úwait©r   rz   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   )rw   rp   rt   rr   ÚasyncioÚrun_coroutine_threadsafeÚstop_wsÚ	Exception©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yE } zt 	d|› �¡ W Y d }~d S d }~ww )Nc                    s   g | ]	}|ˆ j vr|‘qS r   )rq   )Ú.0Úsr‰   r   r   Ú
<listcomp>ä   s    z-StreamManager.add_symbols.<locals>.<listcomp>zadd_symbols failed: )
rr   rt   rq   Úupdater…   r†   Ú_subscribe_asyncrˆ   ÚlogÚwarning)r   rz   Únew_symsÚer   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   )rr   Ú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   )ro   ÚsymbolrM   ra   r=   r>   Ú	timestamp)r   ÚtradeÚbufr   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   )ro   r›   Ú	bid_priceÚ	ask_pricerP   ra   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   )
ro   r›   rT   ra   ÚopenÚhighÚlowÚcloseÚvolumerœ   )r   r   rž   r   r   r   rš   þ   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.)rl   rw   Úis_setr…   Únew_event_looprt   Úset_event_looprm   Ú
new_streamrr   rp   rq   r•   r–   r—   r˜   r™   rš   r�   Úinfor#   rx   Úrunrˆ   r<   r‘   ÚtimeÚsleepÚerror)r   rz   ÚbackoffsÚmax_attemptsÚattemptr“   Údelayr   r   r   r�     sB   





ÿ
ÿ
€ùòÿzStreamManager._run_foreverr›   rD   c                 C   ó   | j |  ¡ S r   )ro   rU   ©r   r›   r   r   r   rU   %  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.)ro   rJ   r¹   r   r   r   rJ   (  s   zStreamManager.get_bars_subc                 C   s   | j | jS r   )ro   r   r¹   r   r   r   Ú	get_quote.  s   zStreamManager.get_quotec                 C   s   | j |  | jd ¡S )NÚstale_data_seconds)ro   r[   rl   r¹   r   r   r   Úis_symbol_stale1  ó   zStreamManager.is_symbol_stalec                 C   s   | j  ¡ o
| j ¡  S r   )rx   r«   rw   r‰   r   r   r   Ú
is_healthy4  r½   zStreamManager.is_healthyN)r\   r]   r^   r   rF   r‚   rŠ   r”   r�   r–   r˜   rš   r�   ÚstrrU   rJ   rº   rc   r¼   r¾   r   r   r   r   rd   Ã   s    !rd   )r_   r…   ru   r±   Úcollectionsr   r   r   r   r   Úconfig_loaderr   Úlogger_setupr   Úalpaca_clientr	   r�   r   rd   r   r   r   r   Ú<module>   s     &