/
makk
/
Stone_parser
Обзор
Документация
Войти
/
makk
/
Stone_parser
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
master
dedup.py
154 строки
4 KB
makk
stone_monitor
28 май 2026, 09:45
Верифицирован
28 май 2026, 09:45
06d9882
Код
Авторство
О чём код?
""" Stone Monitor — SQLite-дедупликация просмотренных сообщений. """ import hashlib import logging import sqlite3 import threading from datetime import datetime, timezone, timedelta from typing import List from config import CACHE_DB_PATH, CACHE_TTL_DAYS logger = logging.getLogger("stone_monitor.dedup") _lock = threading.Lock() def _get_connection() -> sqlite3.Connection: conn = sqlite3.connect(CACHE_DB_PATH) conn.execute("PRAGMA journal_mode=WAL") conn.execute("PRAGMA synchronous=NORMAL") return conn def init_cache(): """Создаёт таблицу кэша.""" with _lock: conn = _get_connection() conn.execute(""" CREATE TABLE IF NOT EXISTS seen_messages ( msg_hash TEXT PRIMARY KEY, msg_id TEXT, channel TEXT, added_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP ) """) conn.execute("CREATE INDEX IF NOT EXISTS idx_seen_added ON seen_messages(added_at)") conn.commit() conn.close() logger.info("Cache DB initialized: %s", CACHE_DB_PATH) def _clean_old_entries(conn: sqlite3.Connection): """Удаляет записи старше CACHE_TTL_DAYS.""" cutoff = (datetime.now(timezone.utc) - timedelta(days=CACHE_TTL_DAYS)).isoformat() deleted = conn.execute("DELETE FROM seen_messages WHERE added_at < ?", (cutoff,)).rowcount if deleted: logger.info("Cache cleanup: %d old entries removed", deleted) _last_cleanup = None def clean_old_if_needed(): """Очистка не чаще раза в час.""" global _last_cleanup now = datetime.now(timezone.utc) if _last_cleanup and (now - _last_cleanup).seconds < 3600: return _last_cleanup = now with _lock: conn = _get_connection() _clean_old_entries(conn) conn.commit() conn.close() def _msg_hash(item: dict) -> str: """Хэш сообщения: channel + msg_id или channel + первые 100 символов текста.""" msg_id = item.get("msg_id", "") channel = item.get("channel", "") text = item.get("text", "") if msg_id: raw = f"{channel}:{msg_id}" else: raw = f"{channel}:{text[:100]}" return hashlib.sha256(raw.encode()).hexdigest() def is_seen(item: dict) -> bool: """Проверка: видели ли уже это сообщение.""" h = _msg_hash(item) with _lock: conn = _get_connection() row = conn.execute("SELECT 1 FROM seen_messages WHERE msg_hash = ?", (h,)).fetchone() conn.close() return row is not None def mark_seen(item: dict): """Помечает сообщение как обработанное.""" h = _msg_hash(item) with _lock: conn = _get_connection() conn.execute( "INSERT OR IGNORE INTO seen_messages (msg_hash, msg_id, channel, added_at) VALUES (?, ?, ?, ?)", (h, item.get("msg_id", ""), item.get("channel", ""), datetime.now(timezone.utc).isoformat()) ) conn.commit() conn.close() def dedup_items(items: List[dict], skip_cache: bool = False) -> List[dict]: """Фильтрует список: убирает уже виденные сообщения. skip_cache=True — для тестов: не проверяет кэш, не пишет в кэш. """ if not items: return [] if skip_cache: logger.info("Dedup: skip_cache=True, returning all %d items", len(items)) return items clean_old_if_needed() result = [] skipped = 0 for item in items: if is_seen(item): skipped += 1 continue result.append(item) mark_seen(item) if skipped: logger.info("Dedup: kept %d, skipped %d of %d total", len(result), skipped, len(items)) return result def clear_cache(): """Полная очистка кэша.""" with _lock: conn = _get_connection() deleted = conn.execute("DELETE FROM seen_messages").rowcount conn.commit() conn.close() logger.info("Cache cleared: %d entries", deleted) def get_cache_stats() -> dict: """Возвращает статистику кэша.""" with _lock: conn = _get_connection() count = conn.execute("SELECT COUNT(*) FROM seen_messages").fetchone()[0] conn.close() return {"seen_messages": count, "ttl_days": CACHE_TTL_DAYS} init_cache()