
    FIjU                        U d Z ddlZddlZddlZddlZddlZddlZddlZddlm	Z	  e	e
      j                         j                  dz  Z e	e
      j                         j                  dz  ZdZddd	d
ddZeeef   ed<   ddddZeed<   dej*                  fdZ ej.                         Z G d dej2                        Z e       Zej9                  ej:                          ej<                         j?                  e        G d d      Z  G d d      Z!dedefdZ"dedefdZ#	 	 	 d2dede$dedz  dedz  de%e   f
d Z&d!ededz  fd"Z'	 	 	 	 d3d!ed#edz  d$edz  dedz  de$de%e   fd%Z(d!ede)fd&Z*d4ded'ede$fd(Z+d)ededz  fd*Z,d+ej*                  d!ededdfd,Z-d5ded-e$dz  ddfd.Z.d6d/Z/d7d0e$dej`                  fd1Z1y)8u   
Reusable bot observability module (DVI-824).

Records runs/steps/logs/artifacts in bot_diagnostics.db (SQLite).
Bot-agnostic — any future bot plugs in with DiagnosticRun().
No bot-specific logic lives here.
    N)Pathzbot_diagnostics.dbdiagnosticsag  
CREATE TABLE IF NOT EXISTS runs (
    id           TEXT PRIMARY KEY,
    bot          TEXT NOT NULL,
    trigger      TEXT NOT NULL,
    target_date  TEXT,
    started_at   REAL NOT NULL,
    finished_at  REAL,
    status       TEXT NOT NULL DEFAULT 'running',
    summary_json TEXT,
    error        TEXT
);
CREATE INDEX IF NOT EXISTS idx_runs_bot_started ON runs (bot, started_at DESC);
CREATE INDEX IF NOT EXISTS idx_runs_bot_status  ON runs (bot, status);

CREATE TABLE IF NOT EXISTS steps (
    id          TEXT PRIMARY KEY,
    run_id      TEXT NOT NULL,
    seq         INTEGER NOT NULL,
    name        TEXT NOT NULL,
    status      TEXT NOT NULL DEFAULT 'running',
    started_at  REAL NOT NULL,
    finished_at REAL,
    detail_json TEXT,
    FOREIGN KEY (run_id) REFERENCES runs (id)
);
CREATE INDEX IF NOT EXISTS idx_steps_run ON steps (run_id, seq);

CREATE TABLE IF NOT EXISTS logs (
    id      INTEGER PRIMARY KEY AUTOINCREMENT,
    run_id  TEXT NOT NULL,
    step_id TEXT,
    ts      REAL NOT NULL,
    level   TEXT NOT NULL,
    message TEXT NOT NULL,
    FOREIGN KEY (run_id) REFERENCES runs (id)
);
CREATE INDEX IF NOT EXISTS idx_logs_run_level ON logs (run_id, level);
CREATE INDEX IF NOT EXISTS idx_logs_run_step  ON logs (run_id, step_id);

CREATE TABLE IF NOT EXISTS artifacts (
    id           TEXT PRIMARY KEY,
    run_id       TEXT NOT NULL,
    step_id      TEXT,
    name         TEXT NOT NULL,
    kind         TEXT NOT NULL,
    path         TEXT NOT NULL,
    size_bytes   INTEGER,
    content_type TEXT,
    created_at   REAL NOT NULL,
    FOREIGN KEY (run_id) REFERENCES runs (id)
);
CREATE INDEX IF NOT EXISTS idx_artifacts_run ON artifacts (run_id);

CREATE TABLE IF NOT EXISTS settings (
    bot            TEXT PRIMARY KEY,
    max_runs       INTEGER NOT NULL DEFAULT 100,
    max_log_lines  INTEGER NOT NULL DEFAULT 2000,
    retention_days INTEGER
);
zAapplication/vnd.openxmlformats-officedocument.spreadsheetml.sheetz.application/vnd.ms-excel.sheet.macroEnabled.12zapplication/pdfzapplication/jsonztext/csv)xlsxxlsmpdfjsoncsv_CONTENT_TYPESd     max_runsmax_log_linesretention_days_SETTING_DEFAULTSreturnc                      t        j                  t        t              d      } t         j                  | _        | j                  d       | j                  t               | S )zGOpen bot_diagnostics.db, create schema on first use. Caller must close.F)check_same_threadzPRAGMA journal_mode=WAL)	sqlite3connectstrDB_PATHRowrow_factoryexecuteexecutescript_SCHEMA)conns    &/var/www/html/togen/bot_diagnostics.py_open_dbr    _   sB    ??3w<5AD{{DLL*+wK    c                   4    e Zd ZdZdej
                  ddfdZy)_DiagHandlerzECaptures log records emitted inside a DiagnosticRun step into the DB.recordr   Nc                     t        t        dd       }|y t        t        dd       }	 |j                  ||j                  | j	                  |             y # t
        $ r Y y w xY w)Nrunstep_id)getattr_tl_db_log	levelnameformat	Exception)selfr$   r&   r'   s       r   emitz_DiagHandler.emitr   s\    &-c5$&?;%c9d;	KK!1!14;;v3FG 		s   ,A 	A A )__name__
__module____qualname____doc__logging	LogRecordr/    r!   r   r#   r#   o   s    O7,,  r!   r#   c                   P    e Zd ZdZdddeddfdZddZdefd	Zd
e	ddfdZ
ddZy)_StepCtxz:Context manager for one named step within a DiagnosticRun.r&   DiagnosticRunnamer   Nc                 r    || _         || _        t        t        j                               | _        d| _        y )NF)_run_namer   uuiduuid4id_warned)r.   r&   r:   s      r   __init__z_StepCtx.__init__   s)    	
4::<("r!   c           
         | j                   }|j                  5  |xj                  dz  c_        |j                  }|j                  j	                  d| j
                  |j
                  || j                  t        j                         f       |j                  j                          d d d        | j
                  t        _
        | S # 1 sw Y    xY w)N   z_INSERT INTO steps (id, run_id, seq, name, status, started_at) VALUES (?, ?, ?, ?, 'running', ?))r<   _lock	_step_seq_connr   r@   r=   timecommitr)   r'   )r.   r&   seqs      r   	__enter__z_StepCtx.__enter__   s    iiYY 	MMQM--CII5#&&#tzz499;?
 II	 gg	 	s   BCCc                 R   |rd}n| j                   rd}nd}| j                  j                  5  |re| j                  j                  j	                  d|t        j
                         t        j                  dt        |      i      | j                  f       nE| j                  j                  j	                  d|t        j
                         | j                  f       | j                  j                  j                          d d d        d t        _        y# 1 sw Y   d t        _        yxY w)NerrorwarnokzBUPDATE steps SET status=?, finished_at=?, detail_json=? WHERE id=?z3UPDATE steps SET status=?, finished_at=? WHERE id=?F)rA   r<   rE   rG   r   rH   r   dumpsr   r@   rI   r)   r'   )r.   exc_typeexc_valexc_tbstatuss        r   __exit__z_StepCtx.__exit__   s    F\\FFYY__ 	%		''XTYY[$**gs7|5L*MtwwW 		''ITYY[$''2 IIOO""$	% 	% s   CDD&dc                 ,   | j                   j                  5  | j                   j                  j                  dt	        j
                  |      | j                  f       | j                   j                  j                          d d d        y # 1 sw Y   y xY w)Nz)UPDATE steps SET detail_json=? WHERE id=?)r<   rE   rG   r   r   rP   r@   rI   r.   rV   s     r   
set_detailz_StepCtx.set_detail   se    YY__ 	%IIOO##;A( IIOO""$	% 	% 	%s   A*B

Bc                    d| _         | j                  j                  5  | j                  j                  j	                  d| j
                  f       | j                  j                  j                          d d d        y # 1 sw Y   y xY w)NTz)UPDATE steps SET status='warn' WHERE id=?)rA   r<   rE   rG   r   r@   rI   r.   s    r   rN   z_StepCtx.warn   s`    YY__ 	%IIOO##;dggZ IIOO""$		% 	% 	%s   AA==B)r   r8   r   N)r0   r1   r2   r3   r   rB   rK   boolrU   dictrY   rN   r6   r!   r   r8   r8      sH    D#O #3 #4 #T 0%D %T %%r!   r8   c                       e Zd ZdZ	 	 ddedededz  dedz  ddf
dZdd	Zdefd
Zdede	fdZ
deddfdZddZddZdddZdedz  dededdfdZdededdfdZ	 ddededede	dz  def
dZy)r9   u,  
    Context manager recording one bot pipeline execution.

    Usage::

        with DiagnosticRun("winston", "scheduler", "2026-06-30") as run:
            with run.step("fetch_data"):
                log.info("Fetching…")   # captured automatically
            run.set_summary({"rows": 42})
    Nbottriggertarget_daterun_idr   c                    || _         || _        || _        |xs t        t	        j
                               | _        t        j                         | _	        d| _
        d| _        d | _        d| _        t        d   | _        t        d   | _        y )Nr   runningr   r   )r`   ra   rb   r   r>   r?   r@   	threadingLockrE   rF   _statusrG   
_log_countr   _max_log_lines	_max_runs)r.   r`   ra   rb   rc   s        r   rB   zDiagnosticRun.__init__   st     &2TZZ\!2^^%
 04
/@*:6r!   c           	         t               | _        | j                  j                  d| j                  f      j	                         }|r|d   | _        |d   | _        | j                  j                  d| j                  | j                  | j                  | j                  t        j                         f       | j                  j                          | t        _        d t        _        | S )Nz8SELECT max_runs, max_log_lines FROM settings WHERE bot=?r   r   zfINSERT INTO runs (id, bot, trigger, target_date, started_at, status) VALUES (?, ?, ?, ?, ?, 'running'))r    rG   r   r`   fetchonerk   rj   r@   ra   rb   rH   rI   r)   r&   r'   )r.   rows     r   rK   zDiagnosticRun.__enter__   s    Z
jj  F

(* 	  _DN"%o"6D

1WWdhhd.>.>		L	

 	

r!   c                 t   |r| j                   dk(  rd| _         n|s| j                   dk(  rd| _         | j                  5  |r| j                   dk(  rt        |      nd }| j                  r`| j                  j	                  d| j                   t        j
                         || j                  f       | j                  j                          d d d        d t        _	        d t        _
        | j                  d c}| _        |r|j                          t        | j                  | j                         y# 1 sw Y   fxY w)Nre   rM   donez;UPDATE runs SET status=?, finished_at=?, error=? WHERE id=?F)rh   rE   r   rG   r   rH   r@   rI   r)   r&   r'   close_pruner`   rk   )r.   rQ   rR   rS   err_msgr   s         r   rU   zDiagnosticRun.__exit__   s    	1"DLdlli7!DLZZ 	$&-$,,'2Ic'ltGzz

""Q\\499;A 

!!#	$ ::tdjJJLtxx(	$ 	$s   BD..D7r:   c                     t        | |      S N)r8   )r.   r:   s     r   stepzDiagnosticRun.step  s    d##r!   rV   c                    | j                   5  | j                  rU| j                  j                  dt        j                  |      | j
                  f       | j                  j                          d d d        y # 1 sw Y   y xY w)Nz)UPDATE runs SET summary_json=? WHERE id=?)rE   rG   r   r   rP   r@   rI   rX   s     r   set_summaryzDiagnosticRun.set_summary  s`    ZZ 	$zz

""?ZZ]DGG, 

!!#	$ 	$ 	$s   A"A88Bc                 r    | j                   5  | j                  dk(  rd| _        d d d        y # 1 sw Y   y xY w)Nre   rN   rE   rh   r[   s    r   rN   zDiagnosticRun.warn  s1    ZZ 	&||y(%	& 	& 	&s   -6c                 T    | j                   5  d| _        ddd       y# 1 sw Y   yxY w)ziMark the run as cancelled (user-requested stop). Terminal status
        that __exit__ will not override.	cancelledNrz   r[   s    r   cancelzDiagnosticRun.cancel  s'     ZZ 	'&DL	' 	' 	's   'c                    | j                   5  d| _        |rW| j                  rK| j                  j                  dt	        |      | j
                  f       | j                  j                          d d d        y # 1 sw Y   y xY w)NrM   z"UPDATE runs SET error=? WHERE id=?)rE   rh   rG   r   r   r@   rI   )r.   errs     r   failzDiagnosticRun.fail$  se    ZZ 	$"DLtzz

""8Xtww' 

!!#	$ 	$ 	$s   A!A77B r'   levelmessagec           	         | j                   5  | j                  
	 ddd       y| j                  j                  d| j                  |t	        j                         ||f       | xj
                  dz  c_        | j
                  | j                  kD  r8| j                  j                  d| j                  f       | j                  | _        | j                  j                          ddd       y# 1 sw Y   yxY w)z3Write one log record; lock-safe, trims if over cap.NzMINSERT INTO logs (run_id, step_id, ts, level, message) VALUES (?, ?, ?, ?, ?)rD   zZDELETE FROM logs WHERE id = (  SELECT id FROM logs WHERE run_id=? ORDER BY id ASC LIMIT 1))rE   rG   r   r@   rH   ri   rj   rI   )r.   r'   r   r   s       r   r*   zDiagnosticRun._db_log.  s    ZZ 	 zz!	  	  JJ*'499;w?
 OOq O!4!44

"" WWJ	 #'"5"5JJ%	  	  	 s   C*B=C**C3c                 d    | j                  t        t        dd      |j                         |       y)zEDirectly insert a log line, independent of the Python logging module.r'   N)r*   r(   r)   upper)r.   r   r   s      r   logzDiagnosticRun.logD  s!    WS)T2EKKM7Kr!   datakindrv   c                 ~   t        t        j                               }|r|j                  nd}t        | j
                  z  | j                  z  }|j                  dd       |j                         }|| d| z  }	|	j                  |       t        j                  |d      }
| j                  5  | j                  rm| j                  j                  d|| j                  |||t        |	      t        |      |
t        j                         f	       | j                  j!                          ddd       |S # 1 sw Y   |S xY w)zGPersist bytes to disk and record an artifacts row. Returns artifact id.NT)parentsexist_ok.zapplication/octet-streamzINSERT INTO artifacts (id, run_id, step_id, name, kind, path, size_bytes, content_type, created_at) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?))r   r>   r?   r@   ARTIFACTS_DIRr`   mkdirlowerwrite_bytesr
   getrE   rG   r   lenrH   rI   )r.   r:   r   r   rv   art_idr'   art_dirextfpathcts              r   add_artifactzDiagnosticRun.add_artifactH  s    TZZ\"!$''t$((*TWW4dT2jjlVHAcU++$%?@ZZ 		$zz

"": TWWgtTZTB		=	 

!!#		$ 		$ s   .A:D22D<)NN)r   r9   r\   ru   )r0   r1   r2   r3   r   rB   rK   r]   rU   r8   rv   r^   rx   rN   r}   r   r*   r   bytesr   r6   r!   r   r9   r9      s   	 #'!77 7 4Z	7
 d
7 
7&$T *$ $ $$T $d $&
'$ sTz  #      ,L Ls Lt L !%  	
 o 
r!   r9   r`   c                     t               }	 |j                  d| f      j                         }|rt        |      nt        t              	 |j                          S # |j                          w xY w)NzHSELECT max_runs, max_log_lines, retention_days FROM settings WHERE bot=?)r    r   rm   r^   r   rq   )r`   r   rn   s      r   get_settingsr   g  sY    :DllVF
 (* 	  tCyT*;%<<



s   <A A,c           	         t        |       }dD ]  }||v s||   ||<    t               }	 |j                  d| |d   |d   |j                  d      f       |j	                          |j                          |S # |j                          w xY w)Nr   zINSERT INTO settings (bot, max_runs, max_log_lines, retention_days) VALUES (?, ?, ?, ?) ON CONFLICT(bot) DO UPDATE SET   max_runs=excluded.max_runs,   max_log_lines=excluded.max_log_lines,   retention_days=excluded.retention_daysr   r   r   )r   r    r   r   rI   rq   )r`   kwargscurkr   s        r   set_settingsr   s  s    
s
C< ;AYCF :D8 #j/3#7AQ9RS	
 	

J 	

s   ;A8 8B
limitrT   qc                    t               }	 d}| g}|r|dz  }|j                  |       |r%|dz  }|j                  d| dd| dd| dg       |dz  }|j                  t        |d             |j	                  ||      j                         D cg c]  }t        |       c}|j                          S c c}w # |j                          w xY w)NzpSELECT id, bot, trigger, target_date, started_at, finished_at, status, summary_json, error FROM runs WHERE bot=?z AND status=?z< AND (summary_json LIKE ? OR error LIKE ? OR trigger LIKE ?)%z! ORDER BY started_at DESC LIMIT ?i  )r    appendextendminr   fetchallr^   rq   )r`   r   rT   r   r   sqlparamsrs           r   	list_runsr     s     :D% 	
 u?"CMM&!QQCMMQqc8q1X1#Qx8922c%o&!%c6!:!C!C!EFAQF

 G

s   BB< B7$B< 7B< <Crc   c                 |   t               }	 |j                  d| f      j                         }|s	 |j                          y t	        |      }|j                  d| f      j                         D cg c]  }t	        |       c}|d<   |j                  d| f      j                         D cg c]  }t	        |       c}|d<   |j                  d| f      j                         }t        |      D cg c]  }t	        |       c}|d<   ||j                          S c c}w c c}w c c}w # |j                          w xY w)NzSELECT * FROM runs WHERE id=?z/SELECT * FROM steps WHERE run_id=? ORDER BY seqstepsz~SELECT id, run_id, step_id, name, kind, size_bytes, content_type, created_at FROM artifacts WHERE run_id=? ORDER BY created_at	artifactszbSELECT id, run_id, step_id, ts, level, message FROM logs WHERE run_id=? ORDER BY id DESC LIMIT 500logs)r    r   rm   rq   r^   r   reversed)rc   r   rn   resultsar   ls           r   get_runr     s$   :Dll:VIFOOQ. 	

- c!\\AF9hj
DG
w "\\E	 hj
DG
{ ||CI
 (*	 	
 ,4D>:a$q':v

+


 ; 	

s;   #D) .D) /D)D) *D<4D) 0D$D) D) )D;r   r'   c                    t               }	 d}| g}|r$|dz  }|j                  |j                                |r|dz  }|j                  |       |r|dz  }|j                  d| d       |dz  }|j                  t        |d             |j	                  ||      j                         D cg c]  }t        |       c}|j                          S c c}w # |j                          w xY w)NzGSELECT id, run_id, step_id, ts, level, message FROM logs WHERE run_id=?z AND level=?z AND step_id=?z AND message LIKE ?r   z ORDER BY id ASC LIMIT ?i  )r    r   r   r   r   r   r^   rq   )	rc   r   r'   r   r   r   r   r   r   s	            r   get_run_logsr     s     :D( 	 x>!CMM%++-(##CMM'"((CMMAaS(#))c%&'!%c6!:!C!C!EFAQF

 G

s   B!C -C?C C C)c                    t               }	 t        j                         }|j                  d|| f      }|j                  d|| f       |j                          |j                  dkD  |j                          S # |j                          w xY w)u  Directly mark a still-running run 'cancelled' in the DB (terminal).

    Unlike cooperative cancellation (which needs the worker thread to observe a
    flag), this updates the row unconditionally so the UI reflects the stop even
    when the worker is gone — e.g. a run orphaned by a process restart, or one
    wedged inside a blocking step. Returns True if a running row was updated.
    zzUPDATE runs SET status='cancelled', finished_at=?, error=COALESCE(error,'Stopped by user') WHERE id=? AND status='running'RUPDATE steps SET status='error', finished_at=? WHERE run_id=? AND status='running'r   )r    rH   r   rI   rowcountrq   )rc   r   nowr   s       r   force_cancelr     sz     :DiikllW&M

 	3&M	

 	||a



s   AA7 7B	notec                 n   t               }	 t        j                         }|j                  d| f      j                         }|D ]1  }|j                  d|||d   f       |j                  d||d   f       3 |j	                          t        |      |j                          S # |j                          w xY w)a#  Mark every still-'running' run for a bot as 'cancelled'.

    A run only exists while its worker thread is alive; no thread survives a
    process restart, so any row left 'running' at startup is an orphan/zombie
    that would otherwise spin forever in the UI. Returns the count swept.
    z4SELECT id FROM runs WHERE bot=? AND status='running'zUUPDATE runs SET status='cancelled', finished_at=?, error=COALESCE(error,?) WHERE id=?r@   r   )r    rH   r   r   rI   r   rq   )r`   r   r   r   rowsr   s         r   sweep_runningr     s     :Diik||BSF

(* 	  
	ALL6dAdG$
 LL7ag
	 	4y



s   BB" "B4artifact_idc                     t               }	 |j                  d| f      j                         }|rt        |      nd 	 |j	                          S # |j	                          w xY w)Nz"SELECT * FROM artifacts WHERE id=?)r    r   rm   r^   rq   )r   r   rn   s      r   get_artifactr     sP    :Dll0;.

(* 	  tCyT)



s   /A Ar   c                    | j                  d|f       | j                  d|f       | j                  d|f       | j                  d|f       t        |z  |z  }|j                         r!t        j                  t        |      d       y y )Nz$DELETE FROM logs      WHERE run_id=?z$DELETE FROM steps     WHERE run_id=?z$DELETE FROM artifacts WHERE run_id=?z DELETE FROM runs      WHERE id=?T)ignore_errors)r   r   existsshutilrmtreer   )r   rc   r`   r   s       r   _delete_runr   )  sv    LL7&CLL7&CLL7&CLL3&Cc!F*G~~c'l$7 r!   r   c                 x   	 t               }	 |1|j                  d| f      j                         }|r|d   nt        d   }|j                  d| |f      j	                         }|D ]  }t        ||d   |         |j                  d| f      j                         }|r[|d   rVt        j                         |d   dz  z
  }|j                  d	| |f      j	                         }|D ]  }t        ||d   |         |j                          |j                          y# |j                          w xY w# t        $ r Y yw xY w)
zNRemove oldest runs beyond max_runs cap and expired runs beyond retention_days.Nz)SELECT max_runs FROM settings WHERE bot=?r   zJSELECT id FROM runs WHERE bot=? ORDER BY started_at DESC LIMIT -1 OFFSET ?r@   z/SELECT retention_days FROM settings WHERE bot=?r   iQ z2SELECT id FROM runs WHERE bot=? AND started_at < ?)
r    r   rm   r   r   r   rH   rI   rq   r-   )	r`   r   r   rn   excessr   ret_rowcutoffolds	            r   rr   rr   3  sS   #z	ll?#(*  /23z?7H7T \\>h hj	 
  0D!D'3/0 llAC6hj  7#34w/?'@5'HHllH&M (*   4Aags34 KKMJJLDJJL s)   
D- C:D D- D**D- -	D98D9c                     	 t               } 	 | j                  d      j                         D cg c]  }|d   	 }}| j                          |D ]  }t	        |        yc c}w # | j                          w xY w# t
        $ r Y yw xY w)u8   Prune every bot — called by the periodic sweep thread.zSELECT DISTINCT bot FROM runsr   N)r    r   r   rq   rr   r-   )r   r   botsr`   s       r   	prune_allr   [  s    	z	"&,,/N"O"X"X"Z[QAaD[D[JJL 	C3K	 \JJL  s7   
A9 !A$ AA$ "A9 A$ $A66A9 9	BBinterval_secondsc                 d     d fd}t        j                  |dd      }|j                          |S )zCStart a daemon thread that calls prune_all() on the given interval.c                  F    	 t        j                          t                 ru   )rH   sleepr   )r   s   r   _loopz#start_periodic_prune.<locals>._loopk  s    JJ'(K r!   Tz
diag-prune)targetdaemonr:   r\   )rf   Threadstart)r   r   ts   `  r   start_periodic_pruner   i  s,    
 	dFAGGIHr!   )2   NN)NNNr   )zInterrupted by restartru   r\   )i  )2r3   r   r4   r   r   rf   rH   r>   pathlibr   __file__resolveparentr   r   r   r
   r^   r   __annotations__r   
Connectionr    localr)   Handlerr#   _diag_handlersetLevelDEBUG	getLogger
addHandlerr8   r9   r   r   intlistr   r   r   r]   r   r   r   r   rr   r   r   r   r6   r!   r   <module>r      s          
x.
 
 
"
)
),@
@X&&(//-?<~ P<"S#X  (+TUYZ 4 Z'$$  ioo
7??     w}} %      } -
=% =%D[ [@	c 	d 	c  4 		 $J Tz	
 
$Z6C D4K B : 4Z Tz	
  
$Z<  6s # S <c dTk 8g(( 8# 8C 8D 8% %sTz %T %P	3 	):J:J 	r!   