Files
polit-scraper/analyzer.py

147 lines
4.9 KiB
Python

import sqlite3
import json
import os
import subprocess
import time
import sys
import logging
import re
import requests
from datetime import datetime
from dotenv import load_dotenv
# --- KONFIGURATION & LOGGING ---
load_dotenv()
BASE_DIR = os.path.dirname(__file__)
DEBUG_LOG = os.getenv("DEBUG_LOG", "false").lower() == "true"
LOG_LEVEL = logging.DEBUG if DEBUG_LOG else logging.INFO
MY_PID = os.getpid()
def sende_admin(titel, nachricht):
admin_url = f"{os.getenv('NTFY_BASE')}/polit-scraper-admin"
try:
requests.post(
admin_url,
data=nachricht.encode("utf-8"),
headers={"Title": titel.encode("utf-8"), "Priority": "4"},
auth=(os.getenv("ADMIN_USER"), os.getenv("ADMIN_PW")),
timeout=5
)
except: pass
logging.basicConfig(
level=LOG_LEVEL,
format='%(asctime)s [%(levelname)s] [%(process)d] %(message)s',
handlers=[
logging.FileHandler(os.path.join(BASE_DIR, "analyzer.log"), encoding='utf-8'),
logging.StreamHandler(sys.stdout)
]
)
logger = logging.getLogger(f"analyzer_{MY_PID}")
GEMINI_MODEL = "gemini-2.5-flash-lite"
DB_PATH = os.path.join(BASE_DIR, "polit_scraper.db")
INSTRUCTION_PATH = os.path.join(BASE_DIR, "gemini_instructions.md")
def extract_json(text):
start = text.find("{")
ende = text.rfind("}") + 1
if start != -1 and ende > start:
json_str = text[start:ende]
json_str = re.sub(r'```json\s*', '', json_str)
json_str = re.sub(r'\s*```', '', json_str)
return json_str.strip()
return None
def analysiere_mit_gemini(text):
sicherer_text = text[:8000]
anweisung = f"Analysiere diesen Text.\n\nTEXT:\n{sicherer_text}"
env = os.environ.copy()
env["GEMINI_SYSTEM_MD"] = INSTRUCTION_PATH
env["NODE_ENV"] = "production"
env["CI"] = "true"
env["FORCE_COLOR"] = "0"
befehl = ["gemini", "-m", GEMINI_MODEL, "-p", anweisung]
try:
prozess = subprocess.run(
befehl,
capture_output=True,
text=True,
encoding="utf-8",
timeout=300,
env=env,
stdin=subprocess.DEVNULL
)
stdout = prozess.stdout.strip()
stderr = prozess.stderr.strip()
if DEBUG_LOG:
logger.debug(f"DEBUG_STDOUT: {stdout}")
if stderr: logger.debug(f"DEBUG_STDERR: {stderr}")
if "Error: setRawMode EIO" in stderr:
logger.error(f"TTY_ERROR detected in stderr: {stderr}")
return "TTY_ERROR"
if "exhausted your capacity" in stderr or "exhausted your capacity" in stdout:
return "RATE_LIMIT"
json_str = extract_json(stdout)
if json_str:
return json_str
logger.error(f"PARSING_FAILED. Roh-Output war: {stdout}")
return None
except Exception as e:
logger.error(f"SUBPROCESS_CRASH: {e}")
return None
def main():
logger.info(f"Worker {MY_PID} bereit.")
while True:
try:
conn = sqlite3.connect(DB_PATH, timeout=60)
cursor = conn.cursor()
cursor.execute("""
UPDATE analysis_queue
SET status='PROCESSING', worker_id=?, last_update=CURRENT_TIMESTAMP
WHERE id=(SELECT id FROM analysis_queue WHERE status='PENDING' ORDER BY id ASC LIMIT 1)
""", (MY_PID,))
conn.commit()
cursor.execute("SELECT id, url, raw_text FROM analysis_queue WHERE status='PROCESSING' AND worker_id=?", (MY_PID,))
job = cursor.fetchone()
if job:
job_id, url, text = job
logger.info(f"Analysiere Job {job_id}: {url}")
ergebnis = analysiere_mit_gemini(text)
if ergebnis == "RATE_LIMIT" or ergebnis == "TTY_ERROR":
logger.warning(f"Job {job_id} ({url}) returned {ergebnis}. Re-queueing and waiting.")
cursor.execute("UPDATE analysis_queue SET status='PENDING', worker_id=NULL WHERE id=?", (job_id,))
conn.commit()
time.sleep(60)
elif ergebnis:
cursor.execute("UPDATE analysis_queue SET status='ANALYZED', result_json=?, worker_id=NULL WHERE id=?", (ergebnis, job_id))
logger.info(f"Job {job_id} ANALYZED.")
conn.commit()
else:
cursor.execute("UPDATE analysis_queue SET status='FAILED', worker_id=NULL WHERE id=?", (job_id,))
logger.error(f"Job {job_id} FAILED (Parser returned None).")
conn.commit()
conn.close()
time.sleep(1)
except sqlite3.OperationalError:
time.sleep(5)
except Exception as e:
logger.error(f"Main Loop Fehler: {e}")
time.sleep(10)
if __name__ == "__main__": main()