diff --git a/poller.py b/poller.py index 952739f..83f9c71 100755 --- a/poller.py +++ b/poller.py @@ -298,6 +298,9 @@ class JsonlSenke: def zustand(self, qid): return None, None + def rollback(self): + self.db.rollback() + def schreibe(self, quelle, items, slot): frisch, neu, geaendert = [], 0, 0 for it in items: @@ -357,6 +360,9 @@ class PostgresSenke: import psycopg self.conn = psycopg.connect(dsn, autocommit=False) + def rollback(self): + self.conn.rollback() + def sync_quellen(self, quellen): with self.conn.cursor() as cur: for q in quellen: @@ -553,9 +559,18 @@ def durchlauf(quellen, senke, slot): except Exception as e: # eine Quelle darf nicht fehler = f"{type(e).__name__}: {e}" # den Durchlauf killen status = getattr(e, "code", None) + # Ohne Rollback bleibt eine abgebrochene Transaktion stehen, und + # das notiere() gleich darunter scheitert an ihr statt den Fehler + # festzuhalten - der eigentliche Grund waere dann verloren. + senke.rollback() dauer = int((time.monotonic() - t0) * 1000) - senke.notiere(q["id"], slot, dauer, status, gesehen, neu, geaendert, - uebersprungen, bytes_roh, fehler) + try: + senke.notiere(q["id"], slot, dauer, status, gesehen, neu, geaendert, + uebersprungen, bytes_roh, fehler) + except Exception as e: # das Protokoll darf den Lauf nicht killen + senke.rollback() + print(f" ! {q['id']:14} Protokoll fehlgeschlagen: " + f"{type(e).__name__}: {e}", flush=True) ges_neu += neu ges_geaendert += geaendert ges_gesehen += gesehen