o
    |~£jœ  ã                   @   sä   d 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G dd„ dƒƒZ			
ddededededededefdd„Zddededefdd„Zdefdd„Z			
ddedededededededefdd„Zd
S )aá  
fast_pipeline.py

[FEATURE 2026-09-11] Live wiring for the "fast_prediction" decision
engine: stream_features -> fast_prediction_engine -> fast_entry_gate.
Same role as prediction_pipeline.py plays for the "prediction_pipeline"
decision engine -- a thin orchestrator, all real logic lives in the
engines themselves.

Selected via config.json's entry.decision_engine == "fast_prediction".
See fast_prediction_engine.py's module docstring for why this path
exists (2026-09-10's zero-trade session, the structure_score cold-start
+ extension-cliff interaction it was built to route around) and for
confirmation that stream_features.py itself is completely unchanged --
this only changes how its output is consumed downstream.

CALLER-HELD STATE: monitor.py keeps one persistence dict per symbol
(Monitor._fast_engine_state) across poll cycles and passes it in as
`persistence_state`; this function returns the updated value for the
caller to store back. A fresh {} for a symbol is a cold start.
é    )Ú	dataclassÚfieldÚasdict)ÚdatetimeÚtimezone)Úcompute_features)Úcompute_fast_prediction)Úevaluate_fast_entryc                   @   sÂ   e Zd ZU eed< eed< eed< eed�Z	eed< eed�Z
eed< dZeed< dZeed	< dZeed
< dZeed< dZeed< dZeed< dZeed< dZeed< dZeed< dZeed< dS )ÚFastPipelineDecisionÚsymbolÚshould_enterÚconfirmation_score)Údefault_factoryÚreasons_forÚreasons_againstNÚ	directionÚ
confidenceÚmomentumÚvolume_stateÚ
vwap_stateÚ
ema9_stateÚ
resistanceÚ	extensionÚstateg        Úpersistence_seconds_elapsed)Ú__name__Ú
__module__Ú__qualname__ÚstrÚ__annotations__ÚboolÚfloatr   Úlistr   r   r   r   r   r   r   r   r   r   r   r   © r#   r#   úfast_pipeline.pyr
      s    
 r
   é   Nr   ÚbarsÚbars_subÚavg_vol_baseliner   Úsub_bucket_secondsÚpremarket_resultc	              
   C   sd   |pt  tj¡}| dd¡ | di ¡ t| |||||d ||d�}	|	|d< t| |	||d�}
|
|	fS )aX  
    [FEATURE 2026-09-11 v2] Features + fast_prediction only -- NOT the
    entry gate. Split out of evaluate_fast_pipeline() so a caller can
    rank a wider pool of candidates by fast_prediction_engine's own
    Direction/Confidence reading BEFORE deciding which ones are worth
    spending a full evaluate_fast_entry() persistence-gate call on.

    Why this exists: monitor.py._scan_for_entries() used to pick its
    top-5 entry shortlist by intraday_health.health_score, computed
    before any fast_prediction reading existed for that cycle -- a
    chicken-and-egg problem for ranking by confidence instead. Per
    explicit instruction ("the confidence value should be the value
    used for selecting the top_5"), monitor.py now calls this for
    EVERY health-eligible candidate first, ranks by prediction.confidence,
    and only calls evaluate_fast_entry_only() below for the top 5.

    Returns (prediction, features) -- both fully computed, so a
    subsequent evaluate_fast_entry_only() call for a chosen candidate
    never recomputes either.

    bars / bars_sub / quote / avg_vol_baseline: same shapes monitor.py
        already hands prediction_pipeline.evaluate_entry_pipeline().
    state: caller-held per-symbol dict, see module docstring. Holds
        {"prior_features": ..., "persistence_state": ...}. Updates
        state["prior_features"] as a side effect (unconditionally, even
        for a candidate that doesn't end up in the top 5 -- its features
        genuinely evolved this cycle regardless of ranking, and the next
        cycle's acceleration math needs that real prior reading).
    Úprior_featuresNÚpersistence_state)r(   Úpriorr)   Úas_of)r*   )r   Únowr   ÚutcÚ
setdefaultr   r   )r   r&   r'   Úquoter(   r   r)   r.   r*   ÚfeaturesÚ
predictionr#   r#   r$   Úcompute_fast_ranking3   s    þr5   Úreturnc                 C   sr   |pt  tj¡}t| |||d |d�}|j|d< t| |j|j|j	|j
|j|j|j|j|j|j|j|j|j|jd�S )a1  
    Finishes evaluate_fast_pipeline()'s work given an already-computed
    (prediction, features) pair from compute_fast_ranking() above --
    avoids a redundant compute_features() call for a candidate already
    scored this cycle. Same return shape/contract as
    evaluate_fast_pipeline() below.
    r,   ©r.   )r   r   r   r   r   r   r   r   r   r   r   r   r   r   r   )r   r/   r   r0   r	   Úpersistencer
   r   r   r   r   r   r   r   r   r   r   r   r   r   )r   r4   r3   r   r.   Údecisionr#   r#   r$   Úevaluate_fast_entry_only`   s   

úr:   c                 C   s.   t | ƒ}| jr| j ¡ nd|d< |t |ƒdœS )a  
    [FEATURE 2026-09-11 v3] Full raw-indicator snapshot for logging --
    not just the classified labels (Strong/Expanding/Bullish/...) that
    the [FAST] log line and FastPipelineDecision carry, but every
    underlying numeric/raw value stream_features.py and
    fast_prediction_engine.py computed this cycle (slope_5m,
    momentum_acceleration, relative_volume, volume_acceleration, vwap
    value/slope/classification, price_vs_ema9_pct/ema9_slope, atr_pct,
    vwap_distance_atr, resistance_level/distance_pct, spread_pct, rsi,
    etc.). Per explicit instruction ("make sure ... we have all the
    indicator values logged during the life cycle of the symbol") --
    the classified labels alone aren't enough to reconstruct exactly
    what the engine saw at a given moment; this is the reproducible
    version. Dumps `features` (a stream_features.StreamFeatures) and
    `prediction` (a fast_prediction_engine.FastPredictionReading)
    wholesale via dataclasses.asdict() rather than hand-picking fields,
    so a new indicator added to either dataclass is captured
    automatically without this function needing to be updated too.
    NÚcomputed_at)r3   r4   )r   r;   Ú	isoformat)r3   r4   Úfr#   r#   r$   Úsnapshot_dictv   s   r>   c	                 C   s@   |pt  tj¡}t| ||||||||d�	\}	}
t| |	|
||d�S )aL  
    Convenience wrapper combining compute_fast_ranking() +
    evaluate_fast_entry_only() for a caller that doesn't need to rank
    against other candidates first (e.g. a one-off/test evaluation).
    monitor.py._scan_for_entries() calls the two split functions
    directly instead -- see compute_fast_ranking()'s docstring.
    )r)   r.   r*   r7   )r   r/   r   r0   r5   r:   )r   r&   r'   r2   r(   r   r)   r.   r*   r4   r3   r#   r#   r$   Úevaluate_fast_pipeline�   s   

þr?   )r%   NN)N)Ú__doc__Údataclassesr   r   r   r   r   Ústream_featuresr   Úfast_prediction_enginer   Úfast_entry_gater	   r
   r   r"   r!   ÚdictÚintr5   r:   r>   r?   r#   r#   r#   r$   Ú<module>   sD    þÿÿÿ
þ-þÿÿÿþþ