
    L7`j	'                    \   S r SSKJr  SSKrSSKrSSKrSSKJrJrJr  SSK	J
r
  SSKJr  SrSS jrSS	 jrSS
 jr            SS jrSS.         SS jjrS S jrS!S jrS"S jrSS.         S#S jjrSr\S.S$S jjrSS.             S%S jjrSS jrS&S'S jjrg)(u@   SQLite-хранилище: схема, upsert, бэкап, runs.    )annotationsN)datetime	timedeltatimezone)Path)Iterablea  
PRAGMA foreign_keys = ON;

CREATE TABLE IF NOT EXISTS channels (
  username             TEXT PRIMARY KEY,
  tg_id                INTEGER,
  title                TEXT,
  subscribers          INTEGER,
  last_pulled_at       TEXT
);

CREATE TABLE IF NOT EXISTS posts (
  channel              TEXT NOT NULL REFERENCES channels(username),
  post_id              INTEGER NOT NULL,
  date                 TEXT NOT NULL,
  text                 TEXT NOT NULL DEFAULT '',
  link                 TEXT,
  media_type           TEXT,
  views                INTEGER,
  forwards             INTEGER,
  reposts              INTEGER,
  comments             INTEGER,
  reactions            INTEGER,
  er                   REAL,
  err                  REAL,
  views_growth         TEXT,
  is_deleted           INTEGER NOT NULL DEFAULT 0,
  first_seen_at        TEXT NOT NULL,
  last_updated_at      TEXT NOT NULL,
  stats_updated_at     TEXT,
  PRIMARY KEY (channel, post_id)
);

CREATE INDEX IF NOT EXISTS idx_posts_date ON posts(channel, date);
CREATE INDEX IF NOT EXISTS idx_posts_stats_updated ON posts(stats_updated_at);

CREATE TABLE IF NOT EXISTS runs (
  started_at           TEXT PRIMARY KEY,
  finished_at          TEXT,
  status               TEXT,
  succeeded_channels   TEXT,
  failed_channels      TEXT,
  error_summary        TEXT
);

CREATE TABLE IF NOT EXISTS channel_snapshots (
  channel      TEXT NOT NULL REFERENCES channels(username),
  date         TEXT NOT NULL,
  subscribers  INTEGER,
  captured_at  TEXT NOT NULL,
  PRIMARY KEY (channel, date)
);
c                 d    [         R                  " [        R                  S9R	                  S5      $ )Ntz%Y-%m-%dT%H:%M:%SZ)r   nowr   utcstrftime     /home/rasp/tgstat-puller/db.pynow_isor   B   s!    <<8<<(112FGGr   c                ~    [         R                  " U 5      nUR                  S5        [         R                  Ul        U$ )NzPRAGMA foreign_keys = ON)sqlite3connectexecuteRowrow_factorydb_pathconns     r   r   r   F   s.    ??7#DLL+,{{DKr   c                    U R                   R                  SSS9  [        U 5       nUR                  [        5        UR                  S5        UR                  5         S S S 5        g ! , (       d  f       g = f)NT)parentsexist_okzPRAGMA user_version = 1)parentmkdirr   executescriptSCHEMAr   commitr   s     r   init_dbr%   M   sS    NN5		T6"./ 
		s   7A&&
A4c          
         [        U 5       nUR                  SXX4[        5       45        UR                  5         S S S 5        g ! , (       d  f       g = f)Nal  
            INSERT INTO channels (username, tg_id, title, subscribers, last_pulled_at)
            VALUES (?, ?, ?, ?, ?)
            ON CONFLICT(username) DO UPDATE SET
              tg_id = excluded.tg_id,
              title = excluded.title,
              subscribers = excluded.subscribers,
              last_pulled_at = excluded.last_pulled_at
            r   r   r   r$   )r   usernametg_idtitlesubscribersr   s         r   upsert_channelr,   U   sD     
	T e')<	
 	 
		s   .A
A)on_datec                   [        5       nU=(       d    USS n[        U 5       nUR                  SXX$45        UR                  5         SSS5        g! , (       d  f       g= f)u  Записать срез числа подписчиков.

on_date='YYYY-MM-DD' — для исторического бэкфилла; None = сегодня (UTC).
Один срез на канал на день; повтор в ту же дату обновляет строку.N
   a  
            INSERT INTO channel_snapshots (channel, date, subscribers, captured_at)
            VALUES (?, ?, ?, ?)
            ON CONFLICT(channel, date) DO UPDATE SET
              subscribers = excluded.subscribers,
              captured_at = excluded.captured_at
            )r   r   r   r$   )r   channelr+   r-   r   dayr   s          r   insert_snapshotr2   m   sW     )C

S"XC		T ;,		
 	 
		s   %A
A c                t   [        5       n[        U 5       nUR                  SUUS   US   UR                  SS5      UR                  S5      UR                  S5      UR                  S5      UR                  S	5      UR                  S
5      UR                  S5      UR                  S5      UR                  S5      UR                  S5      UR                  S5      (       a  [        R
                  " US   SS9OSUU45      nUR                  5         UR                  sSSS5        $ ! , (       d  f       g= f)u   Вставить пост, если такого ещё нет. Возвращает 1, если вставлено, иначе 0.ac  
            INSERT OR IGNORE INTO posts (
              channel, post_id, date, text, link, media_type, views, forwards,
              reposts, comments, reactions, er, err, views_growth,
              is_deleted, first_seen_at, last_updated_at, stats_updated_at
            ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, 0, ?, ?, NULL)
            post_iddatetext link
media_typeviewsforwardsrepostscomments	reactionsererrviews_growthFensure_asciiN)r   r   r   getjsondumpsr$   rowcount)r   r0   postr   r   curs         r   insert_postrJ      s    
)C		Tll YV$ &!$#$%HLQ_H`H`

4/eDfj!
6 	||; 
		s   D	D))
D7c                   [        5       nUR                  S5      (       a  [        R                  " US   SS9OS n[	        U 5       nUR                  SUR                  S5      UR                  S5      UR                  S5      UR                  S5      UR                  S	5      UR                  S
5      UR                  S5      UUUUU45        UR                  5         S S S 5        g ! , (       d  f       g = f)NrA   FrB   a  
            UPDATE posts SET
              views = ?, forwards = ?, reposts = ?, comments = ?, reactions = ?,
              er = ?, err = ?, views_growth = ?,
              last_updated_at = ?, stats_updated_at = ?
            WHERE channel = ? AND post_id = ?
            r:   r;   r<   r=   r>   r?   r@   )r   rD   rE   rF   r   r   r$   )r   r0   r4   statsr   growth_jsonr   s          r   update_post_statsrN      s    
)CKP99UcKdKd$**U>2GjnK		T 		'"		*%		)$		*%		+&		$		% 	
. 	1 
		s   BC''
C5c                    [        U 5       nUR                  S[        5       X45        UR                  5         S S S 5        g ! , (       d  f       g = f)NzVUPDATE posts SET is_deleted = 1, last_updated_at = ? WHERE channel = ? AND post_id = ?r'   )r   r0   r4   r   s       r   mark_deletedrP      s<    		TdY)	
 	 
		s   -A
A)
now_iso_tsc               4   [         R                  " U=(       d
    [        5       R                  SS5      5      nU[	        US9-
  R                  S5      n[        U 5       nUR                  SX45      R                  5       sSSS5        $ ! , (       d  f       g= f)up   Посты канала за последние `days` дней, не помеченные удалёнными.Zz+00:00)daysr   z
            SELECT * FROM posts
            WHERE channel = ? AND is_deleted = 0 AND date >= ?
            ORDER BY date DESC
            N)	r   fromisoformatr   replacer   r   r   r   fetchall)r   r0   rT   rQ   now_dtcutoffr   s          r   posts_in_windowrZ      s{     ##Z%<79$E$Ec8$TUFyd++556JKF		T||
 
 (* 
		s   !B		
B   )stale_after_hoursc               l   [         R                  " [        R                  S9nUR	                  S5      nU[        US9-
  R	                  S5      n[        U 5       nUR                  S[        5       U45        UR                  SU45        UR                  5         SSS5        U$ ! , (       d  f       U$ = f)uZ  Начать новый run. Заодно закрывает зависшие 'running'-строки старше
`stale_after_hours` как failed — иначе жёсткое падение процесса оставляет
их без finished_at навсегда, и они не попадают в recent_runs() (алерт молчит).r
   z%Y-%m-%dT%H:%M:%S.%fZ)hoursu:  
            UPDATE runs SET
              finished_at = ?, status = 'failed',
              error_summary = 'stale run: процесс не завершился штатно (crash/kill), помечено при следующем старте'
            WHERE status = 'running' AND started_at < ?
            z;INSERT INTO runs (started_at, status) VALUES (?, 'running')N)
r   r   r   r   r   r   r   r   r   r$   )r   r\   r   
started_atrY   r   s         r   	start_runr`      s     ,,(,,
'C56JI$566@@AXYF		T Y	
 	IM	
 	 
  
	 s   A B$$
B3r7   )error_summaryc               :   [        U 5       nUR                  S[        5       U[        R                  " USS9[        R                  " USS9UU45      nUR                  5         UR                  S:X  a  [        SU< 35      e S S S 5        g ! , (       d  f       g = f)Nz
            UPDATE runs SET
              finished_at = ?, status = ?,
              succeeded_channels = ?, failed_channels = ?, error_summary = ?
            WHERE started_at = ?
            FrB   r   z(finish_run: no run found for started_at=)r   r   r   rE   rF   r$   rG   RuntimeError)r   r_   status	succeededfailedra   r   rI   s           r   
finish_runrg      s     
	Tll 	

959

66
  	<<1!I*XYY % 
		s   A6B
Bc                    U R                  5       (       d  g [        R                  " X R                  U R                  S-   5      5        g )Nz.bak)existsshutilcopy2with_suffixsuffix)r   s    r   	backup_dbrn     s2    >>
LL--gnnv.EFGr   c                    [        U 5       nUR                  SU45      R                  5       sS S S 5        $ ! , (       d  f       g = f)NzQSELECT * FROM runs WHERE finished_at IS NOT NULL ORDER BY started_at DESC LIMIT ?)r   r   rW   )r   limitr   s      r   recent_runsrq   $  s5    		T||_H
 (*	 
		s	   !7
A)returnstr)r   r   rr   zsqlite3.Connection)r   r   rr   None)r   r   r(   rs   r)   
int | Noner*   
str | Noner+   ru   rr   rt   )
r   r   r0   rs   r+   ru   r-   rv   rr   rt   )r   r   r0   rs   rH   dictrr   int)
r   r   r0   rs   r4   rx   rL   rw   rr   rt   )r   r   r0   rs   r4   rx   rr   rt   )
r   r   r0   rs   rT   rx   rQ   rv   rr   list[sqlite3.Row])r   r   r\   floatrr   rs   )r   r   r_   rs   rd   rs   re   	list[str]rf   r{   ra   rs   rr   rt   )   )r   r   rp   rx   rr   ry   )__doc__
__future__r   rE   rj   r   r   r   r   pathlibr   typingr   r#   r   r   r%   r,   r2   rJ   rN   rP   rZ   STALE_RUN_HOURSr`   rg   rn   rq   r   r   r   <module>r      si   F "    2 2  4
nH  	
   
:   	
  
6 F< IM*-;E"  <K @ ZZ Z 	Z
 Z Z Z 
Z>Hr   