U
    ^¨j†.  ã                   @   s    d dl mZ d dlmZmZmZ d dlmZmZ d dlm	Z	m
Z
mZmZmZ dZdZdZdd	„ Zd
d„ Zdd„ Zdd„ Zddd„Zdd„ Zdd„ Zddd„ZdS )é    )ÚMongoClient)ÚPyMongoErrorÚCollectionInvalidÚBulkWriteError)ÚdatetimeÚ	timedelta)ÚDictÚAnyÚListÚUnionÚOptionalÚFLUX__TSi,  z --c                 C   sP   |  ¡ › d�}||  ¡ krHz| j|dddœd� W n tk
rF   Y nX | | S )uF   
    Helper para obtener/crear la colecciÃ³n TimeSeries correcta.
    Ú_flux_tsÚ	timestampÚminutes)Ú	timeFieldÚgranularity)Ú
timeseries)ÚlowerÚlist_collection_namesÚcreate_collectionr   )ÚdbÚ	centro_idÚcollection_name© r   ú7/var/www/itgmanager/itgmanager/nanoox/nanoox_bolt_ts.pyÚget_flux_collection   s    þþ
r   c                 C   s|   |   d¡}|stdƒ‚zt |d¡}W n" tk
rH   td|› �ƒ‚Y nX |   d¡}||   d¡|   d¡|   d¡t ¡ |d	œS )
zN
    Helper para validar y transformar un raw json en un documento mongo.
    r   zCampo 'timestamp' faltanteú%Y-%m-%d %H:%M:%Su   Formato fecha invÃ¡lido: Úidz	Flow ratezCumulant float point nummberÚ	connected)r   Ú	flow_rateÚcumulantr   Úreceived_atÚoriginal_id)ÚgetÚ
ValueErrorr   ÚstrptimeÚutcnow)ÚitemZraw_timestampZ	dt_objectZexternal_idr   r   r   Úparse_flux_document#   s    

úr)   c              
   C   sj   z*| t  }t||ƒ}t|ƒ}| |¡ W dS  tk
rd } ztdt|ƒ› �ƒ W Y ¢dS d}~X Y nX dS )uŽ   
    Procesa UN solo registro (Modo SÃ­ncrono) con manejo de errores.
    - Si tiene Ã©xito: Retorna None.
    - Si falla: Retorna False.
    Nz[ERROR] save_flux_single: F)ÚFLUX_DATABASE_NAMEr   r)   Ú
insert_oneÚ	ExceptionÚprintÚstr)ÚclientÚdatar   r   Ú
collectionÚdocÚer   r   r   Úsave_flux_single=   s    

r4   c                 C   sò   | t  }t||ƒ}g }g }|D ]„}zFt|ƒ}| |¡ |d }	|	dk	rft|	tƒr\| |	¡ n
| |	¡ W q tk
r  }
 ztd|
› �ƒ W Y ¢qW 5 d}
~
X Y qX q|rîz|j	|dd� W n2 t
k
rì } ztd|j› �ƒ W 5 d}~X Y nX |S )u…   
    Procesa un ARRAY de registros (Modo AsÃ­ncrono).
    Siempre guarda. Retorna lista PLANA UNIDIMENSIONAL de IDs encontrados.
    r#   Nz3Advertencia: Saltando item corrupto en batch Flux: F)ÚorderedzAdvertencia Bulk: )r*   r   r)   ÚappendÚ
isinstanceÚlistÚextendr%   r-   Úinsert_manyr   Údetails)r/   Ú	data_listr   r   r1   Údocs_to_insertZ
ids_to_ackr(   r2   Úoidr3   Úbwer   r   r   Úsave_flux_batchW   s,    


"r@   éh  é   c              
   C   s�  |g g g t gdœi}�z:||  ¡ kr*|W S | | }t ¡ }|jddddd�}|tdd� }	d||	dœi}
ddddddœ}| |
|¡ dd	¡ |¡}t	|ƒ}|s¤|W S t	t
|ƒƒ}|d
d
|… }g }g }g }g }|D ]b}| d¡}t|tƒrò| ¡ nt|ƒ}| | d¡¡ | | d¡¡ | | d¡¡ | |¡ qÒ||||dœ}||i}|W S  tk
�rŠ } ztd|› �ƒ | W Y ¢S d
}~X Y nX d
S )uµ   
    Obtiene datos histÃ³ricos de Flux.
    Si no hay datos, retorna la estructura vacÃ­a estandarizada.
    Si el Ãºltimo dato es > 5 min, agrega estado desconectado al final.
    ©r    r!   Zconnection_statusÚTSr   ©ÚhourÚminuteÚsecondÚmicrosecondrB   ©Údaysr   ©z$gtez$lt©r    r!   r   r   Ú_idéÿÿÿÿNr    r!   r   ú"Error obteniendo historicos Flux: )ÚDEFAULT_LAST_SYNC_TIMEr   r   ÚnowÚreplacer   ÚfindÚsortÚlimitr8   Úreversedr$   r7   Ú	isoformatr.   r6   r,   r-   )r   r   Ú
sensor_keyÚNÚstepÚempty_resultr1   ÚahoraÚ
inicio_hoyÚinicio_mananaÚqueryÚ
projectionÚcursorÚdocsÚdocs_ascÚdocs_sampledÚflow_rate_valsÚcumulant_valsÚconnected_valsÚts_valsr2   Zraw_tsÚts_strZvectors_packageÚ	resultador3   r   r   r   Úget_flux_historical_data   sh    	üÿ	 ÿû	
ü	 ÿrl   c           
   
   C   sà   t }zž||  ¡ kr|W S | | }t ¡ }|jddddd�}|tdd� }|jd||dœidgdddœd	�}|ržd|krž|d }t|tƒr”| d
¡W S t	|ƒW S |W S  t
k
rÚ }	 ztd|	› �ƒ | W Y ¢S d}	~	X Y nX dS )u‘   
    Busca el Ãºltimo registro recibido HOY para determinar el estado de conexiÃ³n.
    Retorna la fecha formateada o ' --' si no hay datos.
    r   rE   rB   rJ   r   rL   )r   rO   )r   rN   )rU   ra   r   z#Advertencia en get_flux_heartbeat: N)rQ   r   r   rR   rS   r   Úfind_oner7   Ústrftimer.   r,   r-   )
r   r   Zhora_defaultr1   r]   r^   r_   Zlast_heartbeat_docZrec_atr3   r   r   r   Úget_flux_heartbeatÞ   s0     ÿû

ro   c                 C   s>   t | tƒr|  ¡ S t | tƒr | S td| › dt| ƒ› d�ƒ‚dS )ud  
    @brief Verifica si un timestamp es un objeto datetime y lo convierte a ISO 8601.

    @param timestamp Puede ser un objeto datetime o una cadena.
    @return string con el timestamp en formato ISO 8601 si es datetime.
            Si ya es string, se devuelve sin cambios.
    @throws ValueError si el timestamp no es un datetime o string vÃ¡lido.
    u    Formato de timestamp invÃ¡lido: z (ú)N)r7   r   rX   r.   r%   Útype)r   r   r   r   Úserialize_timestamp	  s
    	

rr   c              
   C   sš  |dgdgdgt gdœi}�z>||  ¡ kr0|W S | | }t ¡ }|jddddd�}|tdd� }	d||	d	œi}
dddddd
œ}| |
|¡ dd¡ |¡}t	|ƒ}|sª|W S t	t
|ƒƒ}|dd|… }g }g }g }g }|D ]j}t| d¡tƒrú| d¡ ¡ nt| d¡ƒ}| | d¡¡ | | d¡¡ | | d¡¡ | |¡ qØ|||||dœi}|W S  tk
�r” } ztd|› �ƒ | W Y ¢S d}~X Y nX dS )um   
    Obtiene datos histÃ³ricos de Flux.
    Si no hay datos, retorna la estructura vacÃ­a estandarizada.
    NFrC   r   rE   rB   rJ   r   rL   rM   rO   r    r!   r   rP   )rQ   r   r   rR   rS   r   rT   rU   rV   r8   rW   r7   r$   rX   r.   r6   r,   r-   )r   r   rY   rZ   r[   r\   r1   r]   r^   r_   r`   ra   rb   rc   rd   re   rf   rg   rh   ri   r2   rj   rk   r3   r   r   r   Úlegacy_get_flux_historical_data  sd    üÿ	 ÿû	,üÿ	rs   N)rA   rB   )rA   rB   )Úpymongor   Úpymongo.errorsr   r   r   r   r   Útypingr   r	   r
   r   r   r*   ZUMBRAL_SEGUNDOSrQ   r   r)   r4   r@   rl   ro   rr   rs   r   r   r   r   Ú<module>   s   (
_+