diff --git a/src/database.py b/src/database.py index b54e145..e6ff7bf 100644 --- a/src/database.py +++ b/src/database.py @@ -938,6 +938,22 @@ async def init_db(): await db.commit() logger.info("Migration: user_activity_days angelegt (MAU/DAU-Erfassung)") + # Migration: System-Status (Key-Value). Der Monitor meldet hier z.B. den + # Telegram-Session-Status, das Verwaltungsportal liest ihn nur an. + cursor = await db.execute( + "SELECT name FROM sqlite_master WHERE type='table' AND name='system_status'" + ) + if not await cursor.fetchone(): + await db.execute(""" + CREATE TABLE system_status ( + key TEXT PRIMARY KEY, + value TEXT, + updated_at TEXT DEFAULT CURRENT_TIMESTAMP + ) + """) + await db.commit() + logger.info("Migration: system_status angelegt (u.a. Telegram-Session-Status)") + # Migration: Token-Usage-Monatstabelle cursor = await db.execute("SELECT name FROM sqlite_master WHERE type='table' AND name='token_usage_monthly'") if not await cursor.fetchone(): diff --git a/src/main.py b/src/main.py index 58a0f45..a1592ee 100644 --- a/src/main.py +++ b/src/main.py @@ -222,6 +222,41 @@ async def daily_source_health_check(): finally: await db.close() + +async def update_telegram_status(): + """Prüft die Telegram-Session und meldet den Status in system_status. + + Das Verwaltungsportal zeigt den Eintrag im Reiter Recherche-Zugänge an. + Nutzt den eigenen Telethon-Client des Prozesses, ein Fehler darf den + Betrieb nie stören (nur Status + Log). + """ + db = await get_db() + try: + status = {"ok": False, "account": None, "error": None} + try: + from feeds.telegram_parser import TelegramParser + client = await TelegramParser()._get_client() + if client: + me = await client.get_me() + status["ok"] = True + status["account"] = ((me.first_name or "").strip() or None) if me else None + else: + status["error"] = "Session fehlt oder nicht autorisiert" + except Exception as e: + status["error"] = str(e) + await db.execute( + "INSERT INTO system_status (key, value, updated_at) " + "VALUES ('telegram_session', ?, CURRENT_TIMESTAMP) " + "ON CONFLICT(key) DO UPDATE SET value = excluded.value, updated_at = CURRENT_TIMESTAMP", + (json.dumps(status),), + ) + await db.commit() + logger.info(f"Telegram-Status gemeldet: ok={status['ok']}") + except Exception as e: + logger.error(f"Telegram-Status-Update fehlgeschlagen: {e}", exc_info=True) + finally: + await db.close() + async def cleanup_expired(): """Bereinigt abgelaufene Lagen basierend auf retention_days.""" db = await get_db() @@ -344,8 +379,12 @@ async def lifespan(app: FastAPI): scheduler.add_job(check_auto_refresh, "interval", minutes=1, id="auto_refresh") scheduler.add_job(cleanup_expired, "interval", hours=1, id="cleanup") scheduler.add_job(daily_source_health_check, "cron", hour=4, minute=0, id="source_health") + scheduler.add_job(update_telegram_status, "cron", hour=4, minute=30, id="telegram_status") scheduler.start() + # Telegram-Status einmal beim Start melden (asynchron, blockiert den Start nicht) + asyncio.create_task(update_telegram_status()) + logger.info("OSINT Lagemonitor gestartet") yield