
    /qjG                    *   d Z ddlmZ ddlZddlZddlZddlZddlmZmZm	Z	 ddl
mZ ddlmZ  ed      Zej                  j!                  dd	       dd
ZddZddZ	 	 	 	 	 	 d	 	 	 	 	 	 	 	 	 	 	 	 	 	 	 	 	 ddZg dZdddZdddZdddZdddZy) a  Aria Audit-Log + DLQ (KAR-230 + KAR-233).

Append-only Audit-Log + Dead-Letter-Queue in SQLite.
Spaeter Migration nach Supabase EU geplant (separat).

aria.audit_log: pro datenfluss-event ein record (ts, actor, action,
target_kind, target_id, payload_hash, status, error).

aria.dlq: fail-events mit retry-schedule (exp-backoff).

Usage:
    from aria_audit import audit, dlq_push, dlq_pop_ready

    audit("aria-akp-deep", "promote", "video", "abc123", payload={...}, status="ok")
    dlq_push("aria-akp-deep", payload, error_msg)
    )annotationsN)datetimetimezone	timedelta)Path)Anyz/root/aria/state/audit.sqliteT)parentsexist_okc                 b    t        j                  t              } t         j                  | _        | S )N)sqlite3connectDB_PATHRowrow_factorycons    /root/aria/lib/aria_audit.py_connr      s     
//'
"CkkCOJ    c                 d    t               5 } | j                  d       ddd       y# 1 sw Y   yxY w)zIdempotent schema init.a  
            CREATE TABLE IF NOT EXISTS audit_log (
              id INTEGER PRIMARY KEY AUTOINCREMENT,
              ts TEXT NOT NULL,
              actor TEXT NOT NULL,
              action TEXT NOT NULL,
              target_kind TEXT,
              target_id TEXT,
              payload_hash TEXT,
              status TEXT NOT NULL,
              error TEXT,
              metadata_json TEXT
            );
            CREATE INDEX IF NOT EXISTS ix_audit_ts ON audit_log(ts);
            CREATE INDEX IF NOT EXISTS ix_audit_actor ON audit_log(actor, ts);
            CREATE INDEX IF NOT EXISTS ix_audit_action ON audit_log(action);

            CREATE TABLE IF NOT EXISTS dlq (
              id INTEGER PRIMARY KEY AUTOINCREMENT,
              ts_created TEXT NOT NULL,
              ts_next_retry TEXT NOT NULL,
              actor TEXT NOT NULL,
              action TEXT NOT NULL,
              payload_json TEXT,
              error TEXT,
              retry_count INTEGER NOT NULL DEFAULT 0,
              max_retries INTEGER NOT NULL DEFAULT 5,
              status TEXT NOT NULL DEFAULT 'pending'
            );
            CREATE INDEX IF NOT EXISTS ix_dlq_next ON dlq(status, ts_next_retry);
        N)r   executescriptr   s    r   init_schemar   $   s1    	 C  	  s   &/c                    | syt        j                  | dt              j                         }t	        j
                  |      j                         d d S )N T)	sort_keysdefault   )jsondumpsstrencodehashlibsha256	hexdigest)payloadraws     r   _payload_hashr'   H   sB    
**Wc
:
A
A
CC>>#((*3B//r   c                N   t                t               5 }|j                  dt        j                  t
        j                        j                         | |||t        |      |||rt        j                  |t              ndf	      }	|	j                  cddd       S # 1 sw Y   yxY w)z)Append-only audit record. Returns row id.zINSERT INTO audit_log (ts, actor, action, target_kind, target_id, payload_hash, status, error, metadata_json) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)r   N)r   r   executer   nowr   utc	isoformatr'   r   r   r    	lastrowid)
actoractiontarget_kind	target_idr%   statuserrormetadatar   curs
             r   auditr7   O   s     M	 Ckk1 X\\*446g&5=

8S14

 }}!  s   A<BB$)<   i,  i  i   iQ c                |   t                t        j                  t        j                        }|t        t        d         z   }t               5 }|j                  d|j                         |j                         | |t        j                  |t              ||f      }|j                  cddd       S # 1 sw Y   yxY w)z2Push failure to DLQ. Schedules first retry in 60s.r   secondszyINSERT INTO dlq (ts_created, ts_next_retry, actor, action, payload_json, error, max_retries) VALUES (?, ?, ?, ?, ?, ?, ?)r)   N)r   r   r+   r   r,   r   BACKOFF_SECONDSr   r*   r-   r   r   r    r.   )	r/   r0   r%   r4   max_retriesr+   
next_retryr   r6   s	            r   dlq_pushr?   r   s    M
,,x||
$Cy);<<J	 Ckk+ $$&

7C0
 }}  s   AB22B;c           	     D   t                t        j                  t        j                        j                         }t               5 }| r't        |j                  d|| |f            cddd       S t        |j                  d||f            cddd       S # 1 sw Y   yxY w)z Get DLQ entries ready for retry.zlSELECT * FROM dlq WHERE status='pending' AND ts_next_retry <= ? AND actor = ? ORDER BY ts_next_retry LIMIT ?Nz^SELECT * FROM dlq WHERE status='pending' AND ts_next_retry <= ? ORDER BY ts_next_retry LIMIT ?)	r   r   r+   r   r,   r-   r   listr*   )r/   limitr+   r   s       r   dlq_pop_readyrC      s    M
,,x||
$
.
.
0C	 
C1eU# 
 
 CKKl%L
 
 
 
s    B0BBc                <   t                t               5 }|j                  d| f      j                         }|s
	 ddd       y|d   dz   }|r|j                  d|| f       	 ddd       y||d   k\  r|j                  d||| f       	 ddd       yt        t        |t        t              dz
           }t        j                  t        j                        t        |      z   j                         }|j                  d	|||| f       ddd       y# 1 sw Y   yxY w)
zQMark a retry attempt. Schedules next retry with exp-backoff or marks done/failed.z5SELECT retry_count, max_retries FROM dlq WHERE id = ?Nretry_count   z6UPDATE dlq SET status='done', retry_count=? WHERE id=?r=   zAUPDATE dlq SET status='failed', retry_count=?, error=? WHERE id=?r:   zAUPDATE dlq SET retry_count=?, ts_next_retry=?, error=? WHERE id=?)r   r   r*   fetchoner<   minlenr   r+   r   r,   r   r-   )dlq_idsuccessr4   r   rowrE   backoffr>   s           r   dlq_mark_retryrN      s   M	 
CkkQTZS\]ffh
 
 -(1,KKPS^`fRgh
 
 #m,,KK[^ikprx]yz
 
 "#k33G!3K"LMll8<<09W3MMXXZ
O*eV4	

 
 
s   %DD+DA7DDc           	        t                t        j                  t        j                        t        |      z
  j                         }t               5 }d}|g}| r|dz  }|j                  |        |j                  d| d|      j                         }ddd       D ci c]  }|d    d|d	    d|d
    |d    c}S # 1 sw Y   .xY wc c}w )u>   Quick stats — wieviele audit-events last N hours pro action.)hourszts >= ?z AND actor = ?zASELECT actor, action, status, COUNT(*) as n FROM audit_log WHERE z/ GROUP BY actor, action, status ORDER BY n DESCNr/   :r0   r3   n)r   r   r+   r   r,   r   r-   r   appendr*   fetchall)r/   rP   sincer   whereparamsrowsrs           r   audit_summaryrZ      s    M\\(,,')%*@@KKME	 
Cw%%EMM% {{OPUw W= >
 (*	 	
 JNNAqzl!AhK=!H+73?NN
 
 Os   AC!CC)returnzsqlite3.Connection)r[   None)r%   r   r[   r    )NNNokr   N)r/   r    r0   r    r1   
str | Noner2   r^   r%   r   r3   r    r4   r    r5   zdict | Noner[   int)r      )r/   r    r0   r    r%   r   r4   r    r=   r_   r[   r_   )N
   )r/   r^   rB   r_   r[   zlist[sqlite3.Row])r   )rJ   r_   rK   boolr4   r    r[   r\   )N   )r/   r^   rP   r_   r[   dict)__doc__
__future__r   r"   r   r   timer   r   r   pathlibr   typingr   r   parentmkdirr   r   r'   r7   r<   r?   rC   rN   rZ    r   r   <module>rm      s     #     2 2  
.
/   TD  1!H0 #    	
     	@ /,"
.Or   