Files
imap-mail-filter-service/app/routers/filters.py
T
duffyduckandClaude Opus 4.7 1c30bffb5e add "reset all processed" button + endpoint
POST /api/filters/reset-all-processed drops and recreates the
processed_mails table — sekundenschnell auch bei Millionen Zeilen
via DROP+CREATE (same pattern as the fast clear-logs). Safe because
already-moved mails are no longer in the source folder, so a fresh
re-evaluation cannot re-move them, only pick up ones that were
skipped due to a rule bug (e.g. the has_attachment fix).

Button "Alle Regeln zurücksetzen" next to "Neue Regel" on the filter
list, with a confirmation dialog and result feedback.

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
2026-09-01 12:39:03 +02:00

215 lines
7.5 KiB
Python

from fastapi import APIRouter, Depends, HTTPException
from sqlalchemy import func, text
from sqlalchemy.orm import Session
import logging
from app.database import engine, get_db
from app.models.db_models import (
Account,
FilterAction,
FilterCondition,
FilterLog,
FilterRule,
LogLevel,
ProcessedMail,
)
from app.schemas.schemas import FilterRuleCreate, FilterRuleResponse, FilterRuleUpdate
logger = logging.getLogger(__name__)
def _reset_processed_for_rule(db: Session, rule_id: int) -> int:
"""Reset processed mails for a specific rule so they get re-evaluated."""
count = (
db.query(ProcessedMail)
.filter(ProcessedMail.rule_id == rule_id)
.delete()
)
if count:
logger.info("Filter geändert: %d verarbeitete Mails für Regel %d zurückgesetzt", count, rule_id)
return count
router = APIRouter(prefix="/api/filters", tags=["filters"])
@router.post("/reset-all-processed")
def reset_all_processed(db: Session = Depends(get_db)):
"""Setzt den 'verarbeitet'-Status ALLER Mails zurück — alle Regeln bewerten
beim nächsten Poll wieder jede Mail im Ordner neu. Sicher, weil schon
verschobene Mails nicht mehr in der Quelle liegen und daher nicht erneut
verschoben werden können. Nutzt DROP+CREATE statt DELETE, damit auch
bei Millionen Zeilen sekundenschnell."""
try:
row_count = db.query(ProcessedMail).count()
db.close() # Sitzung freigeben, damit DDL nicht am Lock hängt
with engine.begin() as conn:
ProcessedMail.__table__.drop(conn, checkfirst=True)
ProcessedMail.__table__.create(conn, checkfirst=True)
conn.execute(text("PRAGMA wal_checkpoint(TRUNCATE)"))
logger.info("ProcessedMail vollständig zurückgesetzt (%d Einträge entfernt)", row_count)
return {"reset": row_count}
except Exception as e:
logger.exception("reset-all-processed fehlgeschlagen")
raise HTTPException(500, f"Fehler beim Zurücksetzen: {e}")
@router.post("/retry-failed")
def retry_failed_fetches(account_id: int | None = None, db: Session = Depends(get_db)):
"""Löscht ProcessedMail-Einträge für UIDs, für die es einen ERROR-Log
mit 'Fehler beim Abrufen' gibt — damit der Scheduler sie beim nächsten
Poll neu versucht (z.B. nach einem Codec-Fix)."""
try:
log_q = db.query(FilterLog.mail_uid, FilterLog.account_id).filter(
FilterLog.level == LogLevel.ERROR,
FilterLog.message.like("Fehler beim Abrufen%"),
FilterLog.mail_uid.isnot(None),
)
if account_id is not None:
log_q = log_q.filter(FilterLog.account_id == account_id)
pairs = log_q.distinct().all()
if not pairs:
return {"reset": 0, "unique_uids": 0}
# Bulk-Delete pro Account: alle UIDs in einem Query
by_account: dict[int | None, set[str]] = {}
for uid, acc_id in pairs:
by_account.setdefault(acc_id, set()).add(uid)
total = 0
for acc_id, uids in by_account.items():
q = db.query(ProcessedMail).filter(ProcessedMail.mail_uid.in_(list(uids)))
if acc_id is not None:
q = q.filter(ProcessedMail.account_id == acc_id)
total += q.delete(synchronize_session=False)
db.commit()
unique_uids = sum(len(v) for v in by_account.values())
logger.info(
"Retry-Failed: %d ProcessedMail-Einträge für %d UIDs entfernt", total, unique_uids
)
return {"reset": total, "unique_uids": unique_uids}
except Exception as e:
logger.exception("retry-failed fehlgeschlagen")
raise HTTPException(500, f"Fehler bei retry-failed: {e}")
@router.get("/account/{account_id}", response_model=list[FilterRuleResponse])
def list_filters(account_id: int, db: Session = Depends(get_db)):
account = db.get(Account, account_id)
if not account:
raise HTTPException(404, "Konto nicht gefunden")
rules = (
db.query(FilterRule)
.filter(FilterRule.account_id == account_id)
.order_by(FilterRule.priority)
.all()
)
return rules
@router.get("/{rule_id}", response_model=FilterRuleResponse)
def get_filter(rule_id: int, db: Session = Depends(get_db)):
rule = db.get(FilterRule, rule_id)
if not rule:
raise HTTPException(404, "Filterregel nicht gefunden")
return rule
@router.post("/", response_model=FilterRuleResponse, status_code=201)
def create_filter(data: FilterRuleCreate, db: Session = Depends(get_db)):
account = db.get(Account, data.account_id)
if not account:
raise HTTPException(404, "Konto nicht gefunden")
if data.priority is None:
# Ans Ende einsortieren: max(priority) + 10 für dieses Konto. Falls noch
# keine Regel existiert, starten wir bei 10.
current_max = db.query(func.max(FilterRule.priority)).filter(
FilterRule.account_id == data.account_id
).scalar()
priority = (current_max + 10) if current_max is not None else 10
else:
priority = data.priority
rule = FilterRule(
account_id=data.account_id,
name=data.name,
priority=priority,
enabled=data.enabled,
stop_processing=data.stop_processing,
source_folder=data.source_folder,
)
db.add(rule)
db.flush()
for cond_data in data.conditions:
cond = FilterCondition(rule_id=rule.id, **cond_data.model_dump())
db.add(cond)
for action_data in data.actions:
action = FilterAction(rule_id=rule.id, **action_data.model_dump())
db.add(action)
# Neue Regel → hat noch keinen processed-Status, wird automatisch alle Mails prüfen
db.commit()
db.refresh(rule)
return rule
@router.put("/{rule_id}", response_model=FilterRuleResponse)
def update_filter(rule_id: int, data: FilterRuleUpdate, db: Session = Depends(get_db)):
rule = db.get(FilterRule, rule_id)
if not rule:
raise HTTPException(404, "Filterregel nicht gefunden")
update_data = data.model_dump(exclude_unset=True)
# Update conditions if provided
if "conditions" in update_data:
for cond in rule.conditions:
db.delete(cond)
for cond_data in data.conditions:
cond = FilterCondition(rule_id=rule.id, **cond_data.model_dump())
db.add(cond)
del update_data["conditions"]
# Update actions if provided
if "actions" in update_data:
for action in rule.actions:
db.delete(action)
for action_data in data.actions:
action = FilterAction(rule_id=rule.id, **action_data.model_dump())
db.add(action)
del update_data["actions"]
for key, value in update_data.items():
setattr(rule, key, value)
# Regel geändert → nur diese Regel zurücksetzen
_reset_processed_for_rule(db, rule.id)
db.commit()
db.refresh(rule)
return rule
@router.delete("/{rule_id}", status_code=204)
def delete_filter(rule_id: int, db: Session = Depends(get_db)):
rule = db.get(FilterRule, rule_id)
if not rule:
raise HTTPException(404, "Filterregel nicht gefunden")
# processed-Einträge werden per CASCADE gelöscht
db.delete(rule)
db.commit()
@router.put("/reorder/{account_id}")
def reorder_filters(account_id: int, rule_ids: list[int], db: Session = Depends(get_db)):
for priority, rule_id in enumerate(rule_ids):
rule = db.get(FilterRule, rule_id)
if rule and rule.account_id == account_id:
rule.priority = priority
db.commit()
return {"message": "Reihenfolge aktualisiert"}