o
    yë¡j<g  ã                   @   s  d Z ddlmZmZ ddlmZmZ ddlmZ ddlm	Z	m
Z
mZmZmZmZmZmZmZmZmZ i dd“dd	“d
d“dd“dd“dd“dd“dd“dd“dd“dd“dddddœ“dd“d d	“d!d"“d#d$“ZeG d%d&„ d&ƒƒZdZd(ed)ed*ed+efd,d-„Zd.ed+efd/d0„Zd.ed1ed+efd2d3„Zd4ed5ed+efd6d7„Zd[d.ed1ed9ed+efd:d;„Zd<ed=ed+efd>d?„Z d.ed@ed+e!fdAdB„Z"dCZ#dDZ$dEZ%dFZ&dGZ'dHZ(dIZ)dJe*dKedLedMedNe*dOed+efdPdQ„Z+	'	'	R	'	'd\d(ed.edSedTedUd&dVedWedOed+efdXdY„Z,d'S )]a!  
stream_features.py

[SIMULATION-ONLY 2026-09-05] Pure computation layer over stream.py's
rolling buffers (1-min bars, 30s-default sub-bars with tick counters,
latest quote) producing one normalized StreamFeatures snapshot per
symbol. NOT consumed anywhere in the live bot yet -- see config.json's
stream_features._note. This is the foundation trend_engine.py /
regime_engine.py / prediction_engine.py will build on next.

No state, no file I/O, no network calls -- a pure function of whatever
bars/bars_sub/quote/prior are handed to it, same contract as
intraday_health.compute_health(). Nearly every calculation here is a
thin wrapper over indicators.py functions that already exist and are
already used elsewhere in this project (vwap, ema_series, atr, rsi,
normalized_slope_pct, consolidation_tightness, spread_pct) -- this
module's job is assembling them into the richer feature set the
prediction engine needs, not inventing new math.

TWO DIFFERENT KINDS OF "ACCELERATION" -- deliberately kept separate:

  1. Cross-horizon comparison (slope.* below): slope_30s vs slope_1m vs
     slope_3m vs slope_5m computed from the SAME snapshot of bars, no
     history required. If shorter horizons show a progressively
     stronger move in the same direction than longer horizons, that's
     "accelerating"; the reverse pattern is "decelerating." This is
     what the project's original design brief's worked examples (5m
     +0.20%, 3m +0.35%, 1m +0.50%, 30s +0.70% => accelerating) actually
     describe -- a shape comparison at one instant, not a time
     derivative. See acceleration.momentum_acceleration.

  2. True time-derivative (acceleration.velocity_acceleration /
     slope_acceleration): this reading's velocity/slope minus the
     PRIOR reading's, normalized by elapsed time. Needs the caller to
     hold and pass in the previous StreamFeatures for this symbol
     (`prior` param) -- this module stays stateless; state ownership
     lives with the caller, same pattern as
     intraday_health.compute_health()/update_health_state()'s
     `persisted` dict.
é    )Ú	dataclassÚfield)ÚdatetimeÚtimezone)Ú
get_config)ÚvwapÚvwap_seriesÚ
vwap_slopeÚextension_from_vwap_pctÚ
ema_seriesÚatrÚrsiÚnormalized_slope_pctÚconsolidation_tightnessÚ
spread_pctÚclassify_slopeÚmin_bars_1mé   Úmin_bars_subé   Úslope_sub_lookback_bucketsé   Úvwap_time_window_barsé   Úema9_periodé	   Úema20_periodÚema_slope_lookback_barsé   Ú
rsi_periodé   Ú
atr_periodÚrange_expansion_lookback_barsé   Ú#pressure_velocity_scale_pct_per_secgš™™™™™©?Úpressure_weightsgš™™™™™Ù?gffffffÖ?g      Ð?)ÚuptickÚvelocityÚvolumeÚvwap_min_hold_bars_for_supportÚvwap_min_rejections_for_failureÚvwap_max_extension_atrg      @Úvwap_flat_slope_threshold_pctg{®Gáz”?c                   @   s&  e Zd ZU eed< eed< dZeed< ee	d�Z
e	ed< eed�Zeed< eed�Zeed< eed�Zeed	< eed�Zeed
< eed�Zeed< eed�Zeed< eed�Zeed< eed�Zeed< eed�Zeed< eed�Zeed< eed�Zeed< eed�Zeed< eed�Zeed< dS )ÚStreamFeaturesÚsymbolÚcomputed_atFÚinsufficient_data)Údefault_factoryÚmissingÚmetaÚpricer'   ÚslopeÚaccelerationr(   r   ÚemaÚ
volatilityr   ÚquoteÚ
trade_flowÚpressureN)Ú__name__Ú
__module__Ú__qualname__ÚstrÚ__annotations__r   r0   Úboolr   Úlistr2   Údictr3   r4   r'   r5   r6   r(   r   r7   r8   r   r9   r:   r;   © rD   rD   ústream_features.pyr-   I   s$   
 r-   Nr.   r2   r/   Úreturnc                 C   s   t | |p	t tj¡d|d�S )NT)r.   r/   r0   r2   )r-   r   Únowr   Úutc)r.   r2   r/   rD   rD   rE   Ú_insufficient_   s   þrI   Úbarsc                 C   s   dd„ | D ƒS )Nc                 S   ó   g | ]}|d  ‘qS )ÚcrD   ©Ú.0ÚbrD   rD   rE   Ú
<listcomp>g   ó    z_closes.<locals>.<listcomp>rD   )rJ   rD   rD   rE   Ú_closesf   s   rR   Ún_barsc                 C   s   t | ƒ|kr| | d … S | S ©N)Úlen)rJ   rS   rD   rD   rE   Ú_windowj   s   rV   rG   Úthenc                 C   s   |dkrdS | | | d S )Nr   ç        ç      Y@rD   )rG   rW   rD   rD   rE   Ú_pct_changen   s   rZ   ç      N@Úbar_secondsc                 C   s8   t | t|dƒƒ}t|ƒdk rdS tt|ƒƒ}|d|  S )a7  normalized_slope_pct() over the last n_bars closes, converted to
    a %-per-MINUTE basis via (60 / bar_seconds) -- 'slope over an N-bar
    window,' the same standard interpretation intraday_health.py already
    uses for its own slope_lookback_bars (documented there as roughly an
    N-minute window at 1-min bars).

    [BUGFIX 2026-09-05] normalized_slope_pct() returns %-change PER BAR,
    not per unit time -- comparing that raw value between a 60s-bar
    series and a 30s-bar series is comparing different units (a "slope
    of 8 per bar" on 30s bars is a FASTER underlying rate than "8 per
    bar" on 60s bars, since each 30s bar covers half the time). Confirmed
    live in this module's own testing: an explicitly accelerating
    synthetic price series produced a NEGATIVE momentum_acceleration
    (the cross-horizon waterfall read as decelerating) purely because
    slope_sub (30s bars) was being compared unscaled against slope_1m/
    3m/5m (60s bars). Every horizon must be expressed on the same
    per-minute basis before any cross-horizon comparison (this function's
    output, and everything downstream: acceleration.momentum_acceleration
    here and trend_engine.py's flat-threshold classification) is
    meaningful.

    Floors the window at 2 bars: a regression line through a single
    point is undefined and normalized_slope_pct() correctly returns 0.0
    for it per its own guard -- silently reading that 0.0 as "flat" would
    be indistinguishable from an actual flat reading. n_bars=1 ("slope
    over the last 1 minute") is inherently meaningless as a fitted line;
    the last 2 bars is the smallest window that has an actual slope, and
    for exactly 2 points a least-squares fit reduces to the plain
    two-point slope anyway, so this doesn't quietly change behavior for
    any n_bars >= 2 call site.r   rX   r[   )rV   ÚmaxrU   r   rR   )rJ   rS   r\   ÚwindowÚper_barrD   rD   rE   Ú_horizon_slope_pctt   s
   r`   Úbar_aÚbar_bc                 C   s   t |d | d   ¡ dƒS )NÚtç      ð?)r]   Útotal_seconds)ra   rb   rD   rD   rE   Ú_bucket_seconds_betweenš   s   rf   Úwindow_barsc                    s€   t ˆ ƒ‰d}ttˆ ƒd ddƒD ]}ˆ | d ˆ| kr!|d7 }q tdtˆ ƒ| ƒ}t‡ ‡fdd„t|tˆ ƒƒD ƒƒ}||fS )a·  
    hold_duration_bars: current consecutive streak of closes above the
        running VWAP, counted backward from the latest bar (0 if the latest
        close is at/below VWAP right now).
    rejection_count: within the last `window_bars`, how many bars had their
        HIGH reach or cross the running VWAP but still CLOSED below it --
        an intrabar rejection, distinct from simply trading below VWAP the
        whole bar.
    r   é   éÿÿÿÿrL   c                 3   s@   � | ]}ˆ | d  ˆ|   k rˆ | d krn ndV  qdS )rL   Úhrh   NrD   )rN   Úi©rJ   ÚseriesrD   rE   Ú	<genexpr>»   s   € ,ÿþz,_vwap_hold_and_rejections.<locals>.<genexpr>)r   ÚrangerU   r]   Úsum)rJ   rg   Úhold_durationrk   ÚstartÚrejection_countrD   rl   rE   Ú_vwap_hold_and_rejections¨   s   

ÿrt   ÚSTRONG_VWAP_SUPPORTÚVWAP_RECLAIMÚ	VWAP_HOLDÚVWAP_EXTENSIONÚVWAP_FAILUREÚ	VWAP_FLATÚVWAP_DECLININGÚaboveÚslope_classrq   rs   ÚextendedÚcfgc                 C   sd   | r|rt S | r|dkrtS | r||d krtS | r |dkr tS | r$tS |dks.||d kr0tS tS )aú  
    Exhaustive, priority-ordered, mutually exclusive classification -- see
    module docstring / config.json's stream_features._note_vwap_classes for
    the full reasoning behind each branch. A symbol repeatedly holding above
    a positively-sloped VWAP (STRONG_VWAP_SUPPORT) should read as
    meaningfully more confident than one merely above a flat VWAP this
    instant (VWAP_HOLD) -- that's the whole point of this classification
    existing on top of the raw price_vs_vwap_pct/slope fields.
    Únegativer)   r   r*   )rx   r{   ÚVWAP_STRONG_SUPPORTrv   rw   ry   rz   ©r|   r}   rq   rs   r~   r   rD   rD   rE   Ú_classify_vwapË   s   rƒ   é   Úbars_subÚavg_vol_baselineÚpriorÚsub_bucket_secondsÚas_ofc	           T         s‚
  i t ¥|ptƒ  di ¡¥}|pt tj¡}	g }
t|ƒ|d k r(t| g d¢|	d�S |d }|d }d|i}|d d	 }||d
< |rH|| | d nd|d< dD ]!\}}t|ƒ|kr`|d|  d nd||< || du ro|
 	|¡ qNt|ƒdkr|d d |d< n	d|d< |
 	d¡ i }t|ƒdkr¥|d }t
||ƒ}t||d ƒ| |d< n	d|d< |
 	d¡ t|ƒdkrÉ|d }t|dƒ}t||d ƒ| |d< n	d|d< |
 	d¡ t|dƒt|dƒt|dƒdœ}t|ƒ|d k�rttt||d ƒƒƒ}|dt|dƒ  |d< n	d|d< |
 	d¡ |d |d |d g}|d du�r#| 	|d ¡ |d dk�r,dn
|d dk �r5dnd}|dk�rF|d |d  | n|d |d  }dt|dƒi}|du�rÈ|j�sÈt|jtƒ�rot|	|j  ¡ d ƒnd }| d¡}|j�r€|j d¡nd}|du�r˜|du�r˜t|| | d!ƒ|d"< nd|d"< | d¡}|j�r«|j d¡nd}|du�rÃ|du�rÃt|| | d!ƒ|d#< nd|d#< nd|d"< d|d#< |d$ td%d&„ t|dƒD ƒƒtd'd&„ t|dƒD ƒƒd(œ} t|ƒdk�rü|d d$ | d)< nd| d)< d*d+„ |D ƒ}!t|!ƒdk�r9t|!ƒd }"t|!d|"… ƒ|" �p d,}#t|!|"d… ƒt|!ƒ|"  }$t|$|# dƒ| d-< n	d| d-< |
 	d-¡ |�rQtt|!ƒ| dƒ| d.< nd| d.< t|ƒ‰t|ƒdk�rmt|tdt|ƒd ƒd/�nd0}%t||d1 ƒ}&t‡fd2d&„|&D ƒƒt|&ƒ d }'tˆdƒt|%dƒtt|ˆƒdƒt|'dƒtd|' dƒd3œ}(t|ƒ})t|)|d4 ƒ}*t|)|d5 ƒ}+|*�r¼|*d nd},t|d6 t|*ƒƒ}-|-dk�rÕt|*|- d… ƒnd0}.|,du�rát|,dƒnd|+�rìt|+d dƒndt|.dƒ|,�rütt||,ƒdƒndd7œ}/t|t|d8 tdt|ƒd ƒƒd9�}0|�r|0| d nd0}1t|tdt|ƒƒd/�}2|d: }3t|dtdt|ƒƒ … �p;||3d/�}4|4dk�oK|2|4 |4 d }5t|0dƒt|1dƒ|0�r`t|ˆ |0 dƒnd|5d;u�rkt|5dƒndd<œ}6|0�r€|,du�r€t||, |0 dƒnd|/d=< t ||d1 ƒ\}7}8ˆ�r–|%ˆ d nd0}9t!|9|d> ƒ}:|6 d?¡};|;du�o°t"|;ƒ|d@ k}<|7|(dA< |8|(dB< t#|ˆk|:|7|8|<|dC�|(dD< |dE }=t$||=d9�}>t|ƒ|=d k}?|?�sá|
 	dF¡ t|>dƒ|?dGœ}@|�r|\}A}B}C|A|Bt|B|A dƒtt%|A|BƒdƒdHœ}DndddddHœ}D|
 	dI¡ |�rM|d }E|E dJd¡}F|E dKd¡}G|E dLd¡}H|Ft|Ft|dƒ dƒ|F�r=t|G|F dƒnd|F�rHt|H|F dƒnddMœ}IndddddMœ}I|
 	dN¡ |dO ‰i ‰ |I dP¡|I dQ¡}J}K|Jdu�r€|Kdu�r€tdRtd |J|K ƒƒˆ dS< | d¡}L|Ldu�r�|dT �p�d,}MtdRtd |L|M ƒƒˆ dU< |  d-¡}N|Ndu�r³tdRtd |Nd  ƒƒˆ dV< ˆ �rÕt‡fdWd&„ˆ D ƒƒ}Ott‡ ‡fdXd&„ˆ D ƒƒ|O d dƒ}Pnd}P|
 	dY¡ d}Q|Pdu�r|du�r|j�s|j&�r|j& dZ¡}R|Rdu�rt|P|R dƒ}Q|P|Qd[œ}St'dii d\| “d]|	“d^d;“d_|
“d`da|i“db|“dU|“dc|“dd|“dV| “de|(“df|/“dg|6“dh|@“dI|D“dN|I“dY|S“ŽS )ja
  
    bars: 1-min bars, oldest first (stream.StreamManager.get_bars()).
    bars_sub: sub-minute bars, oldest first, forming bucket included
              (stream.StreamManager.get_bars_sub()) -- each carries
              tick_count/upticks/downticks alongside o/h/l/c/v.
    quote: (bid, ask, timestamp) or None.
    avg_vol_baseline: same "expected normal volume" baseline
              entry_engine/intraday_health already take -- used for
              relative_volume-style comparisons. None is handled
              (volume.relative_volume simply omitted from output).
    prior: this symbol's previous StreamFeatures reading, or None on
           the first call -- see module docstring's two-kinds-of-
           acceleration note. Caller-held; this function never stores
           anything itself.
    sub_bucket_seconds: must match whatever streaming.
           sub_minute_bucket_seconds actually is, so time-normalized
           calculations (velocity, acceleration) use the real elapsed
           time rather than an assumed 30.
    as_of: [FEATURE 2026-09-05] the moment this reading is "as of," used
           for computed_at (and therefore the elapsed-time base for
           acceleration vs `prior`). Defaults to real wall-clock time
           (datetime.now()), correct for live use where a bar really
           does arrive at roughly the same real-world instant it's
           timestamped. A REPLAY/BACKTEST caller MUST pass the bar's own
           timestamp here instead -- otherwise computed_at would be the
           real time the backtest script happens to execute, completely
           unrelated to the simulated market time being replayed, which
           would corrupt every acceleration calculation that compares
           consecutive readings' elapsed time.
    Ústream_featuresr   )	r4   r'   r5   r(   r   r7   r8   r   r:   )r/   ri   rL   Úlastr   ÚoÚopenrY   NÚextension_from_open_pct))Úprice_1mrh   )Úprice_3mr   )Úprice_5mr   r   éþÿÿÿÚ	price_subÚvelocity_1mrh   Úvelocity_subr   r   )Úslope_1mÚslope_3mÚslope_5mr   r[   Ú	slope_subr˜   r—   r–   Úmomentum_accelerationr   rd   r#   Úvelocity_accelerationÚslope_accelerationÚvc                 s   ó   � | ]}|d  V  qdS ©r�   NrD   rM   rD   rD   rE   rn   }  ó   € z#compute_features.<locals>.<genexpr>c                 s   rž   rŸ   rD   rM   rD   rD   rE   rn   ~  r    )Ú	volume_1mÚ	volume_3mÚ	volume_5mÚ
volume_subc                 S   rK   )r�   rD   rM   rD   rD   rE   rP   …  rQ   z$compute_features.<locals>.<listcomp>g•Ö&è.>Úvolume_accelerationÚrelative_volume)ÚlookbackrX   r   c                 3   s    � | ]}|d  ˆ krdV  qdS )rL   rh   NrD   rM   )Úv_waprD   rE   rn   ˜  ó   € )Úvaluer5   Úprice_vs_vwap_pctÚtime_above_pctÚtime_below_pctr   r   r   )Úema9Úema20Ú
ema9_slopeÚprice_vs_ema9_pctr!   )Úperiodr"   F)r   Úatr_pctÚvwap_distance_atrÚrange_expansion_pctÚprice_vs_ema9_atrr,   r´   r+   Úhold_duration_barsrs   r‚   Úclassificationr   zFrsi (insufficient history -- neutral 50.0 default, not a real reading))rª   Úhas_full_history)ÚbidÚaskÚspreadr   r9   Ú
tick_countÚupticksÚ	downticks)Útick_count_subÚtrades_per_secondÚuptick_ratioÚdowntick_ratior:   r%   rÂ   rÃ   g      ð¿r&   r$   r'   r(   c                 3   s   � | ]}ˆ | V  qd S rT   rD   ©rN   Úk)ÚpwrD   rE   rn     r    c                 3   s    � | ]}ˆ | ˆ|  V  qd S rT   rD   rÄ   )Ú
componentsrÆ   rD   rE   rn     r©   r;   Úscore)rÈ   r5   r.   r/   r0   r2   r3   rˆ   r4   r5   r6   r   r7   r8   r   rD   )(ÚDEFAULT_CONFIGr   Úgetr   rG   r   rH   rU   rI   Úappendrf   rZ   r]   r`   r   rR   rV   Úroundr0   Ú
isinstancer/   re   r'   r5   rp   r   r	   Úminr
   r   r   r   rt   r   Úabsrƒ   r   r   r;   r-   )Tr.   rJ   r…   r9   r†   r‡   rˆ   r‰   r   Únow_tsr2   Únow_barÚ
last_pricer4   Úsession_openÚlabelÚnr'   Úprev_barÚelapsedÚprev_subÚelapsed_subr5   Ú
per_bucketÚhorizon_chainÚlong_horizon_signrš   r6   Úv_nowÚv_prevÚs_nowÚs_prevr(   Úvols_1mÚmidÚearlyÚrecentÚv_sloperg   Ú
time_aboveÚvwap_outÚ
closes_allÚema9_seriesÚema20_seriesÚema9_nowÚema_lookbackr°   Úema_outÚar³   Útightness_recentÚtightness_lookbackÚtightness_earlierÚrange_expansionr8   rq   rs   Úvwap_slope_pctr}   Úvwap_atr_distr~   r   Ú	rsi_valueÚrsi_has_full_historyÚrsi_outrº   r»   Ú_tsÚ	quote_outÚcurrent_subr½   r¾   r¿   r:   Úuptick_rÚ
downtick_rÚv_subÚscaleÚ	vol_accelÚ
weight_sumÚpressure_scoreÚpressure_slopeÚprior_scoreÚpressure_outrD   )rÇ   rÆ   r¨   rE   Úcompute_featuresè   sÌ  $
ÿ
ÿ$
€




ýÿ
(
ÿÿ
ÿÿ


ý
,"û	 ü$ÿ
ÿüÿ
ÿÿ


þ

ý
üÿ




(
$

ÿþýüûúúúúùùùùøøøør  rT   )r[   )NNr„   NN)-Ú__doc__Údataclassesr   r   r   r   Úconfig_loaderr   Ú
indicatorsr   r   r	   r
   r   r   r   r   r   r   r   rÉ   r-   r?   rB   rI   rR   ÚintrV   ÚfloatrZ   r`   rC   rf   Útuplert   r�   rv   rw   rx   ry   rz   r{   rA   rƒ   r  rD   rD   rD   rE   Ú<module>   sž    )4ÿþýüûúùø	÷
öõôñðïî&ÿÿ
ÿûÿþýüûû