
    djyj                     0   S r SSKrSSKrSSKrSSKrSSKr\R                  R                  \R                  R                  \	5      5      r
\R                  R                  \
S5      r\4S jrS rSS jrSS jrSS jrS	 r\S
:X  Ga  \" \R*                  5      S:  a  \R*                  S   O\R                  R                  \
S5      r\" \R*                  5      S:  a  \R*                  S   OSr\R                  R1                  \5      (       d  \" S\5        \R4                  " S5        \" 5       r\" \\" \SS9R;                  5       \\" 5       5      r\R?                  S5      RA                  5       S   r!\" S\\\!4-  5        \RE                  5         gg)a  Ingest getData JSON into emanagement.sqlite.

Reusable by:
  * db_build.py         (initial build / re-seed)
  * fetch_emanagement.bat  ->  python db_ingest.py _raw.json bat
  * serve.py            (imports ingest_json() and logs every realtime poll)

Design: upsert the point (dimension) if unseen, then INSERT OR IGNORE one
reading per (point, server-measurement-time). Dedup is automatic.
    Nzemanagement.sqlitec                 v    [         R                  " U SSS9nUR                  S5        UR                  S5        U$ )N   F)timeoutcheck_same_threadzPRAGMA journal_mode=WALzPRAGMA foreign_keys=ON)sqlite3connectexecute)dbcons     %/var/www/html/energodata/db_ingest.pyr   r      s5     //"bE
BCKK)*KK()J    c                 (   U (       d  gU R                  S5      nUS   R                  5       (       a?  [        U5      S:  a0  US   < SUS   < 3US   SR                  USS 5      =(       d    S4$ SUS   SR                  USS 5      =(       d    S4$ )zA'01_ee_KGJ1_vykon_procenta' -> ('01_ee','KGJ1','vykon_procenta').)NNN_r            N)splitisdigitlenjoin)namepartss     r   
parse_namer      s    !JJsOEQxc%jAo 8U1X.a#((59:M:UQUVV%(CHHU12Y/7488r   c                 `   S=n=pS=n
=pSnSnU(       a  UR                  S5      =(       d    Sn
UR                  S5      =(       d    SnUR                  S5      =(       d    Sn[        UR                  SS5      =(       d    S5      n[	        UR                  S	S5      =(       d    S5      n[        U
5      u  pxn	S
X#U4-  nU R                  SXX4XXXXX45        U R                  SU45      R                  5       nUS   $ ! [        [        4 a    Sn Nf = f! [        [        4 a    Sn Nf = f)zOEnsure a points row exists; fill sid/metadata when available. Returns point_id.Ng      ?r   r   unittypmultr   decz
%d;%d;%d;0a  INSERT INTO points(key,layer,channel,register,zdroj,sid,name,unit,mult,decimals,typ,site_code,device,metric)
           VALUES(?,?,?,?,?,?,?,?,?,?,?,?,?,?)
           ON CONFLICT(key) DO UPDATE SET
             sid       = COALESCE(points.sid, excluded.sid),
             name      = COALESCE(points.name, excluded.name),
             unit      = COALESCE(NULLIF(points.unit,''), excluded.unit),
             mult      = CASE WHEN points.mult=1 AND excluded.mult<>1 THEN excluded.mult ELSE points.mult END,
             decimals  = CASE WHEN points.decimals=0 AND excluded.decimals<>0 THEN excluded.decimals ELSE points.decimals END,
             typ       = COALESCE(points.typ, excluded.typ),
             site_code = COALESCE(points.site_code, excluded.site_code),
             device    = COALESCE(points.device, excluded.device),
             metric    = COALESCE(points.metric, excluded.metric)z'SELECT point_id FROM points WHERE key=?)getfloat	TypeError
ValueErrorintr   r	   fetchone)curkeylayerchannelregistersidmetasitedevicemetricr   r   r   r   r   zdrojrows                    r   upsert_pointr1   "   s*   !!D!6D4Daxx'4xx'4xx&$$((61-23Ttxxq).Q/S)$/fEH55EKK	E 
WD3V\eg ++?#
H
Q
Q
SCq6M) :&2s2:&/a/s$   $$D  $D  DDD-,D-c                 \   U R                  5       nSnU=(       d    0 R                  5        GHW  u  pxUR                  S5      (       d  M  [        USS 5      n	UR                  5        GH  u  pU
R                  S5      (       d  M  [        U
SS 5      nUR	                  S5      nUR                  5        H  u  pUR                  S5      (       d  M  [        USS 5      nU=(       d    0 R	                  S	5      =(       d    0 nS
U;  a  MV  SXU4-  nU(       a  UR	                  U5      OSn[        UUXUUU5      nUR                  SUUR	                  S5      [        US
   5      45        XeR                  -  nM     GM     GMZ     UR                  SX#U45        U R                  5         U$ )zRdata = {'l1': {...}, 'l2': {...}} nested getData structure. Returns #new readings.r   lr   Nchr   r*   rc0valzl%d.ch%d.r%dz=INSERT OR IGNORE INTO readings(point_id,ts,raw) VALUES(?,?,?)dtz=INSERT INTO snapshots(server_date,source,n_new) VALUES(?,?,?))
cursoritems
startswithr#   r   r1   r	   r    rowcountcommit)r   dataserver_datesourcelabelsr%   newLchsr'   r4   bodyr(   r*   rkrvr)   r6   r&   r+   pids                        r   ingest_datarI   B   sr   
**,C
C:2$$&||C  AabE
		HB==&&"QR&kG((5/C**,}}S))r!"v;hB^^D)/R?$'AA*0vzz#d"3UXsDQ[ "&&,bi0@AC||# ' $	 ', KKOc*,JJLJr   c                     [         R                  " U5      nUS   S   S   n[        XR                  S0 5      UR                  S5      X#5      $ )z)text = full getData response JSON string.returnsr   r5   r>   date)jsonloadsrI   r   )r   textr@   rA   dr5   s         r   ingest_jsonrQ   a   sD    

4A	)QAsEE&"-quuV}fMMr   c                  T   [         R                  R                  [        S5      n [         R                  R	                  U 5      (       d  0 $ [        U SS9R                  5       R                  5       n[        R                  " XR                  S5      UR                  S5      S-    5      $ )Nz	labels.jsz	utf-8-sigencoding{}r   )ospathr   HEREexistsopenreadstriprM   rN   indexrindex)pts     r   _load_labelsrb   g   su    
T;'A77>>!	Q%**,224A::aQXXc]Q%6788r   __main__r   z	_raw.jsonr   batzno file:zutf-8rS   zSELECT COUNT(*) FROM readingsz=ingested %d new reading(s) [source=%s]; total readings now %d)NN)NimportN)re   N)#__doc__r   rM   rW   resysrX   dirnameabspath__file__rY   r   DBr   r   r1   rI   rQ   rb   __name__r   argvrawr@   rZ   printexitr   r[   r\   nr	   r$   totalclose r   r   <module>rv      sS  	 " ! !	wwrwwx01	ww||D./ 9@>N9 zSXX*#((1+T;0OCMA-SXXa[5F77>>#j#
)CCcG499;V\^TAKK78AACAFE	
IQPVX]L^
^_IIK r   