
    5j                      T   d Z ddlZddlZddlZddlZddlmZ ddlmZ  eej                  j                  d ee
      j                         j                  dz              Z ej                         Z ej                   d      Zd Zd	 Zd
 ZddZd ZddZd Zd Zd ZddZd Zy)u  
scada_store.py — SQLite telemetry store for SCADA sensor readings (DVI-1095 P1).

Durable store decoupling SCADA views from live mailbox re-parsing: the
background collector (and, later, a LAN SCADA Agent posting to /scada/ingest)
writes readings here, and /scada/status, /scada/history, and reporting read
them back even when the mailbox is unreachable.

Flask-independent (mirrors scada_report.py) so app.py, the report scheduler,
and tests can all import it. The DB lives alongside the app's other state
files (scada.db next to scada_config.json).

Schema
------
readings(sensor, ts, value, raw_value, source) — one row per sensor sample.
    ts is a sortable text timestamp ("YYYY-MM-DD HH:MM[:SS]"); value is the
    float when the sample is numeric, raw_value always keeps the original
    string ("72.4", "On", "Off"). UNIQUE(sensor, ts, source) makes ingestion
    idempotent — collector cycles re-parse overlapping mailbox windows.
meta(key, value) — collector health + ingest bookkeeping for UI surfacing.
    N)datetime)PathSCADA_DB_FILEzscada.dbz=^(\d{1,2})/(\d{1,2})/(\d{4})\s+(\d{1,2}):(\d{2})(?::(\d{2}))?c                      t        j                  t        t              d      } | j	                  d       | j	                  d       | S )N   )timeoutzPRAGMA journal_mode=WALzPRAGMA busy_timeout=10000)sqlite3connectstrr   executeconns    "/var/www/html/togen/scada_store.py_connectr   (   s6    ??3}-r:DLL*+LL,-K    c                      t         5  t               5 } | j                  d       | j                  d       | j                  d       | j                  d       ddd       ddd       y# 1 sw Y   xY w# 1 sw Y   yxY w)zECreate tables/indexes if missing. Idempotent; call at import/startup.a  CREATE TABLE IF NOT EXISTS readings (
                   id INTEGER PRIMARY KEY AUTOINCREMENT,
                   sensor TEXT NOT NULL,
                   ts TEXT NOT NULL,
                   value REAL,
                   raw_value TEXT NOT NULL,
                   source TEXT NOT NULL DEFAULT 'email',
                   created_at TEXT NOT NULL DEFAULT (datetime('now')),
                   UNIQUE(sensor, ts, source)
               )zOCREATE INDEX IF NOT EXISTS idx_scada_readings_sensor_ts ON readings(sensor, ts)z@CREATE INDEX IF NOT EXISTS idx_scada_readings_ts ON readings(ts)zBCREATE TABLE IF NOT EXISTS meta (key TEXT PRIMARY KEY, value TEXT)N)_write_lockr   r   r   s    r   init_scada_dbr   /   s    	 Rhj RD	
	 	'	( 	N	PP	R#R R R R R Rs#   A4AA(A4(A1	-A44A=c           	         | xs dj                         }|s|S t        j                  |      }|rS|j                         \  }}}}}}| dt	        |      ddt	        |      ddt	        |      dd| 	}	|	|rd| z   S dz   S 	 t        j                  |j                  dd            }
|
j                  |
j                  rd      S d	      S # t        $ r |cY S w xY w)
a%  Normalize a sample timestamp to sortable "YYYY-MM-DD HH:MM[:SS]" text.

    SCADA emails carry US-style "MM/DD/YYYY HH:MM:SS" stamps while Graph
    receivedDateTime is ISO; both must land on one sortable axis. Unparseable
    input is returned as-is (stored verbatim, sorts best-effort).
     -02d :Zz+00:00%Y-%m-%d %H:%M:%Sz%Y-%m-%d %H:%M)strip	_US_DT_REmatchgroupsintr   fromisoformatreplacestrftimesecond
ValueError)rawsmmodayrhhmissoutdts              r   normalize_tsr2   F   s     
AA!"BBBAc"gc]!CGC=#b'#atD""h--"--##AIIc8$<={{")).RRAQRR s   AC C CCc                    g }| xs g D ]  }t        |j                  d      xs d      j                         }t        |j                  d      |j                  d      nd      j                         }t        t        |j                  d      xs d            }|r|r|s	 t	        |      }|j                  |||||f        |syt        5  t               5 }|j                  }	|j                  d|       |j                  |	z
  cddd       cddd       S # t
        $ r d}Y yw xY w# 1 sw Y   nxY wddd       y# 1 sw Y   yxY w)zInsert reading dicts [{sensor, ts, value}] idempotently; return the
    number of NEW rows stored (duplicates on (sensor, ts, source) are ignored).
    ``value`` may be any scalar; the float form is derived when possible.sensorr   valueNtsr   z\INSERT OR IGNORE INTO readings (sensor, ts, value, raw_value, source) VALUES (?, ?, ?, ?, ?))r   getr   r2   floatr&   appendr   r   total_changesexecutemany)
rowssourcepreparedrr4   	raw_valuer6   numr   befores
             r   insert_readingsrC   \   s@    HZR 
>QUU8_*+113!%%.*Dg"MSSU	#aeeDk/R01Yb		"C 	S)V<=
> 	 +hj +D##&'/	1 !!F*+ + +  	C	
+ + + + +s6   "DD?-D*	D?D'&D'*D3	/D??Ec            	          t               5 } | j                  d      j                         }ddd       D cg c]  \  }}}|||d }}}}t        d |D        d      }||dS # 1 sw Y   =xY wc c}}}w )u"  Most recent stored reading per sensor: {"time": <max ts>, "readings":
    [{label, value, ts}]} — the /scada/status fallback shape when Graph is
    down. Per-reading ts lets the UI age each sensor independently (a dead
    sensor must show stale even while its neighbors keep reporting).a  SELECT r.sensor, r.raw_value, r.ts
                 FROM readings r
                WHERE r.id = (SELECT r2.id FROM readings r2
                               WHERE r2.sensor = r.sensor
                            ORDER BY r2.ts DESC, r2.id DESC LIMIT 1)
             ORDER BY r.sensorN)labelr5   r6   c              3   (   K   | ]
  \  }}}|  y w)N ).0_r6   s      r   	<genexpr>z"latest_snapshot.<locals>.<genexpr>   s     *Ar"*s   )default)timereadings)r   r   fetchallmax)r   r<   r(   vr6   rM   latests          r   latest_snapshotrR   v   s    
 
 .t||"# $,8: 	. DHHHxq!R!ar2HHH*T*D9F11. . Is    A'A3'A0c           	      \   g g }}| rA|j                  ddj                  dt        |       z         d       |j                  |        |r+|j                  d       |j                  t	        |             d}|r|ddj                  |      z   z  }|d	z  }i }t               5 }|j                  ||      D ];  \  }}	}
|j                  |g       }t        |      |k  s(|j                  |	|
d
       = 	 ddd       |j                         D ]  }|j                           |S # 1 sw Y   0xY w)a$  Time series per sensor: {sensor: [{time, value}, ...]} ascending by ts.

    ``sensors`` optionally restricts labels; ``since`` is an inclusive
    normalized-ts lower bound. Each series keeps its most recent
    ``limit_per_sensor`` points so one runaway sensor can't bloat the payload.
    zsensor IN (,?)zts >= ?z*SELECT sensor, ts, raw_value FROM readingsz WHERE z AND z ORDER BY sensor, ts DESC)rL   r5   N)
r9   joinlenextendr2   r   r   
setdefaultvaluesreverse)sensorssincelimit_per_sensorwhereparamssqlseriesr   r4   r6   r@   ptss               r   query_seriesre      s&    6E{388C#g,,>#?"@BCgYl5)*
6Cy7<<...&&CF	 =t%)\\#v%> 	=!FB	##FB/C3x**

B;<	==
 }} M= =s   !:D"D""D+c                      t               5 } | j                  d      j                         d   cd d d        S # 1 sw Y   y xY w)NzSELECT COUNT(*) FROM readingsr   r   r   fetchoner   s    r   reading_countri      s=    	 Kt||;<EEGJK K Ks	   "7A c           	          t         5  t               5 }|j                  d| |dn
t        |      f       d d d        d d d        y # 1 sw Y   xY w# 1 sw Y   y xY w)NzaINSERT INTO meta (key, value) VALUES (?, ?) ON CONFLICT(key) DO UPDATE SET value = excluded.valuer   )r   r   r   r   )keyr5   r   s      r   set_metarl      sY    	 8hj 8DE"3u:6	88 8 8 8 8 8s!   A"AAA	
AAc                     t               5 }|j                  d| f      j                         }d d d        r|d   S d S # 1 sw Y   xY w)Nz$SELECT value FROM meta WHERE key = ?r   rg   )rk   r   rows      r   get_metaro      sO    	 VtllAC6JSSUV3q6"d"V Vs	   ">Ac                    t        j                         j                  d      }t        d|       | r1t        d|       t        dd       t        d|       t        d|       y
t        dt	        |xs d             t        d	|       y
)zCPersist the outcome of one collector cycle for UI health surfacing.r   collector_last_runcollector_last_okcollector_last_errorr   collector_last_fetchedcollector_last_insertedzunknown errorcollector_last_error_atN)r   nowr$   rl   r   )okerrorfetchedinsertedrw   s        r   record_collector_runr|      ss    
,,.
!
!"5
6C!3'	$c*',)73*H5'U-Eo)FG*C0r   c                  r    t        d      t        d      t        d      xs dt        d      t               dS )zCollector health summary for API payloads: {last_run, last_ok,
    last_error, last_error_at, readings}. last_error is None when the most
    recent cycle succeeded.rq   rr   rs   Nrv   )last_runlast_ok
last_errorlast_error_atrM   )ro   ri   rG   r   r   collector_healthr      s<    
 12/056>$!";<!O r   )email)NNi  )Nr   r   )__doc__osrer	   	threadingr   pathlibr   environr7   __file__resolveparentr   Lockr   compiler   r   r   r2   rC   rR   re   ri   rl   ro   r|   r   rG   r   r   <module>r      s   , 
 	    JJNN?DN$:$:$<$C$Cj$PQS innBJJDF	R.,+42$:K8#1
r   