a
    $Fwj                     @  s   d dl mZ d dlmZ d dlmZ d dlmZ d dlmZm	Z	 d dl
mZ ddlmZmZ dd	lmZ dd
lmZmZ ddlmZ eegdf ZG dd dZdS )    )annotations)defaultdict)datetime)isfinite)meanpstdev)Callable   )SymbolStateTick)OpportunityQueue)
PricePointRollingWindow)SmartSchedulerNc                   @  st   e Zd ZdZdddddddd	Zd
ddddZed
ddddZeddddddZedddddZ	dS )RealtimeGatewayzIncremental tick processor.

    This module does not open a market-data connection itself. It validates and
    consumes normalized ticks from any provider adapter.
    NzSmartScheduler | NonezOpportunityQueue | NonezDecisionCallback | NoneNone)	schedulerqueueon_candidatereturnc                 C  s6   |pt  | _|pt | _|| _i | _tdd | _d S )Nc                   S  s   t dddS )Ni,  )Z
max_pointsZmax_age_seconds)r    r   r   9/home/opc/hanz-radar/runtime/src/hanz_realtime/gateway.py<lambda>$       z*RealtimeGateway.__init__.<locals>.<lambda>)r   r   r   r   r   statesr   windows)selfr   r   r   r   r   r   __init__   s    zRealtimeGateway.__init__r   r
   )tickr   c           	      C  s  |  | |j|jf}| j|}|d u rDt|j|jd}|| j|< |j}|j}td|j	|j
 }| jd7  _| j|7  _|j|_|j	|_
|j|_|r|dkr|j| d d |_|r|j|kr|j|  }|dkr|r|j| | |_|jd ur*|jd ur*|jdkr*|j|j |j d |_| j| }|t|j|j|j	 | |||_|jddh |jd	kr| jj|j|j|j| ||jd
 | jd ur| | n|jdk r| j |j|j |S )N)symbolmarket        r	   r         ?      Y@pricevolumeg     K@)r   r    scorereason
updated_at     A@)!_validate_tickr    r   r   getr
   
last_pricelast_timestampmaxr%   last_volume
tick_countcumulative_volumer$   	timestampprice_change_pcttotal_secondsvelocitybidask
spread_pctr   addr   _anomaly_scoreanomaly_scoredirty_flagsupdater   Zupsert_reasonr   remove)	r   r   keystateZprevious_priceZprevious_timeZincremental_volumesecondswindowr   r   r   ingest'   sN    

$
zRealtimeGateway.ingestc                 C  s`   | j  std| j s$tdt| jr8| jdkr@tdt| jrT| jdk r\tdd S )Nzsymbol is requiredzmarket is requiredr   z!price must be positive and finitez&volume must be non-negative and finite)r   strip
ValueErrorr    r   r$   r%   )r   r   r   r   r*   X   s    

zRealtimeGateway._validate_tickr   float)rC   rA   r   c                 C  s  |   }|  }t|dk r dS |d }|d }|rHt|| d d nd}t|dkrht|d d n|d }|dkr|d | nd}t|rt|t| d nd}	t|jpdd d	}
t|d
 dtt|d dd	 d t|	d d	 tt|j	d d |
 }t
tdt|ddS )N   r!   r   r"   r#   r	   g       @g      4@g      2@r)   g      $@g       @   )pricesvolumeslenabsr   r   minr8   r.   r5   round)rC   rA   rK   rL   firstlastZmove_pctZaverage_volumeZvolume_ratioZ
volatilityZspread_penaltyr&   r   r   r   r:   c   s,    $ zRealtimeGateway._anomaly_scorestr)rA   r   c                 C  s    | j dkrdS | j dkrdS dS )NP   zANOMALI KUATA   ZMENGUATZPANTAU)r;   )rA   r   r   r   r>   ~   s
    

zRealtimeGateway._reason)NNN)
__name__
__module____qualname____doc__r   rD   staticmethodr*   r:   r>   r   r   r   r   r      s      1
r   )
__future__r   collectionsr   r   mathr   
statisticsr   r   typingr   modelsr
   r   r   r   Zrollingr   r   r   r   ZDecisionCallbackr   r   r   r   r   <module>   s   