#!/usr/bin/env python3 """ connettore_metadataexporter HTTP receiver for Rspamd metadata_exporter, compatibile con il database rqwatch. Endpoints: POST /api/metadata_importer_multipart — formatter = "multipart" (raccomandata) POST /api/metadata_importer — formatter = "default" + meta_headers """ import email as _email import email.header import email.policy import json import logging import os import re import secrets import uuid from datetime import date from pathlib import Path from typing import Optional import pymysql from dotenv import load_dotenv from fastapi import Depends, FastAPI, File, Form, HTTPException, Request, UploadFile from fastapi.responses import PlainTextResponse from fastapi.security import HTTPBasic, HTTPBasicCredentials from pymysql.cursors import DictCursor load_dotenv() # ─── Configurazione ─────────────────────────────────────────────────────────── API_USER: str = os.getenv("RSPAMD_API_USER", "rspamd") API_PASS: str = os.getenv("RSPAMD_API_PASS", "") API_ACL: set[str] = {ip.strip() for ip in os.getenv("RSPAMD_API_ACL", "127.0.0.1").split(",")} QUARANTINE_DIR: str = os.getenv("QUARANTINE_DIR", "/quarantine") SERVER_ALIAS: str = os.getenv("MY_API_SERVER_ALIAS", "mx1") DB_HOST: str = os.getenv("DB_HOST", "127.0.0.1") DB_PORT: int = int(os.getenv("DB_PORT", "3306")) DB_NAME: str = os.getenv("DB_NAME", "rqwatch") DB_USER: str = os.getenv("DB_USER", "rqwatch") DB_PASS: str = os.getenv("DB_PASS", "") MAILLOGS_TABLE: str = os.getenv("MAILLOGS_TABLE", "mail_logs") MAIL_RECIPIENTS_TABLE: str = os.getenv("MAIL_RECIPIENTS_TABLE", "mail_log_recipients") _store_flags: dict[str, str] = { "no action": os.getenv("STORE_NO_ACTION", "false"), "add header": os.getenv("STORE_ADD_HEADER", "true"), "rewrite subject": os.getenv("STORE_REWRITE_SUBJECT", "true"), "greylist": os.getenv("STORE_GREYLIST", "false"), "discard": os.getenv("STORE_DISCARD", "true"), "reject": os.getenv("STORE_REJECT", "true"), } STORE_ACTIONS: set[str] = { k for k, v in _store_flags.items() if v.strip().lower() in ("1", "true", "yes") } # Limiti di campo compatibili con rqwatch MailLog::FIELD_LIMITS FIELD_LIMITS: dict[str, int] = { "qid": 30, "server": 10, "subject": 1024, "action": 20, "ip": 50, "mail_from": 255, "mime_from": 255, "rcpt_to": 1024, "mime_to": 1024, "mail_location": 255, "message_id": 1024, } # ─── Logging ────────────────────────────────────────────────────────────────── logging.basicConfig( level=logging.INFO, format="%(asctime)s %(levelname)s [connettore] %(message)s", ) log = logging.getLogger("connettore") # ─── App ────────────────────────────────────────────────────────────────────── app = FastAPI(docs_url=None, redoc_url=None) _security = HTTPBasic() # ─── Helper DB ──────────────────────────────────────────────────────────────── def _db() -> pymysql.Connection: return pymysql.connect( host=DB_HOST, port=DB_PORT, user=DB_USER, password=DB_PASS, database=DB_NAME, charset="utf8mb4", cursorclass=DictCursor, autocommit=False, ) def _insert(data: dict, recipients: list[str]) -> int: conn = _db() try: with conn.cursor() as cur: cols = ", ".join(f"`{k}`" for k in data) ph = ", ".join(["%s"] * len(data)) cur.execute( f"INSERT INTO `{MAILLOGS_TABLE}` ({cols}) VALUES ({ph})", list(data.values()), ) db_id: int = cur.lastrowid if db_id and recipients: unique = list({r.lower().strip() for r in recipients if r.strip()}) cur.executemany( f"INSERT INTO `{MAIL_RECIPIENTS_TABLE}` " "(mail_log_id, recipient_email) VALUES (%s, %s)", [(db_id, r) for r in unique], ) conn.commit() return db_id except Exception: conn.rollback() raise finally: conn.close() # ─── Helper email ───────────────────────────────────────────────────────────── def _decode_header(value: str) -> str: try: parts = _email.header.decode_header(value) out = [] for raw, charset in parts: if isinstance(raw, bytes): out.append(raw.decode(charset or "utf-8", errors="replace")) else: out.append(raw) return "".join(out) except Exception: return value def _parse_mime(raw: bytes) -> dict: try: msg = _email.message_from_bytes(raw, policy=_email.policy.compat32) def hdr(name: str) -> str: v = msg.get(name, "") return _decode_header(str(v)) if v else "" raw_lines = [f"{k}: {v}" for k, v in msg.items()] headers_raw = "\r\n".join(raw_lines).encode("utf-8", "ignore").decode("utf-8") return { "mime_from": hdr("From"), "mime_to": hdr("To")[:1024], "mime_subject": hdr("Subject"), "message_id": hdr("Message-ID"), "headers_raw": headers_raw, } except Exception as exc: log.warning(f"MIME parse error: {exc}") return { "mime_from": "", "mime_to": "", "mime_subject": "", "message_id": "", "headers_raw": "", } # ─── Helper quarantena ──────────────────────────────────────────────────────── def _store_email(qid: str, raw: bytes) -> Optional[str]: q = Path(QUARANTINE_DIR) if not q.is_dir() or not os.access(str(q), os.W_OK): log.error(f"Quarantine dir non accessibile: {QUARANTINE_DIR}") return None today = date.today().isoformat() safe = qid if re.match(r"^[a-zA-Z0-9]+$", qid) else "" subdir = safe if safe and safe != "unknown" else f"unknown/{uuid.uuid4().hex}" mail_dir = q / today / subdir mail_dir.mkdir(parents=True, exist_ok=True) dest = mail_dir / "mail.eml" try: dest.write_bytes(raw) return str(dest) except Exception as exc: log.error(f"Scrittura fallita {dest}: {exc}") return None def _has_virus(symbols) -> bool: if isinstance(symbols, dict): return any( isinstance(v, dict) and v.get("group") == "antivirus" for v in symbols.values() ) if isinstance(symbols, list): return any( isinstance(s, dict) and s.get("group") == "antivirus" for s in symbols ) return False def _trim_fields(data: dict) -> dict: for field, limit in FIELD_LIMITS.items(): v = data.get(field) if isinstance(v, str) and len(v) > limit: log.warning(f"Campo '{field}' troncato a {limit} caratteri") data[field] = v[:limit] return data def _sanitize_server(s: str) -> str: return re.sub(r"[^a-zA-Z0-9.\-]", "", s)[:10] # ─── Logica principale ──────────────────────────────────────────────────────── def _process( *, qid: str, server: str, subject: str, score: float, action: str, symbols_json: str, fuzzy_json: str, ip: str, mail_from: str, rcpt_list: list[str], size: int, raw_email: bytes, ) -> int: try: symbols = json.loads(symbols_json) except Exception: symbols = {} virus = _has_virus(symbols) mail_stored, mail_location = 0, None if action in STORE_ACTIONS or virus: mail_location = _store_email(qid, raw_email) if mail_location: mail_stored = 1 log.info(f"{qid} in quarantena: {mail_location}") else: log.error(f"{qid} salvataggio quarantena fallito") if not mail_from: mail_from = "empty-mail-from@localhost" mime = _parse_mime(raw_email) data = { "qid": qid, "server": _sanitize_server(server), "subject": mime["mime_subject"] or subject or "", "score": round(score, 2), "action": action, "symbols": symbols_json or "[]", "has_virus": 1 if virus else 0, "fuzzy_hashes": fuzzy_json or "[]", "ip": ip or "", "mail_from": (mail_from or "").lower(), "mime_from": mime["mime_from"], "rcpt_to": "unknown" if not rcpt_list else ", ".join(r.lower() for r in rcpt_list), "mime_to": mime["mime_to"], "mail_stored": mail_stored, "mail_location": mail_location, "size": size, "headers": mime["headers_raw"], "message_id": mime["message_id"], } data = _trim_fields(data) return _insert(data, rcpt_list) # ─── Dipendenza: ACL + autenticazione ───────────────────────────────────────── async def _auth( request: Request, credentials: HTTPBasicCredentials = Depends(_security), ) -> None: client_ip = request.client.host if request.client else "" if client_ip not in API_ACL: log.warning(f"Richiesta da {client_ip} rifiutata (non in RSPAMD_API_ACL)") raise HTTPException(status_code=403, detail="Forbidden") ok = ( secrets.compare_digest(credentials.username.encode(), API_USER.encode()) and secrets.compare_digest(credentials.password.encode(), API_PASS.encode()) ) if not ok: raise HTTPException( status_code=401, detail="Unauthorized", headers={"WWW-Authenticate": 'Basic realm="rqwatch-api"'}, ) # ─── Endpoint: multipart/form-data (formatter = "multipart") ───────────────── @app.post("/api/metadata_importer_multipart", response_class=PlainTextResponse) async def metadata_importer_multipart( request: Request, metadata: str = Form(...), message: UploadFile = File(...), _: None = Depends(_auth), ) -> str: try: meta: dict = json.loads(metadata) except json.JSONDecodeError as exc: raise HTTPException(status_code=400, detail=f"metadata JSON non valido: {exc}") raw_email = await message.read() if not raw_email: raise HTTPException(status_code=400, detail="File message vuoto") qid = str(meta.get("qid") or "unknown") if not re.match(r"^[a-zA-Z0-9]+$", qid): qid = "unknown" score = float(meta.get("score") or 0.0) action = str(meta.get("action") or "") server = request.query_params.get("server", SERVER_ALIAS) if not qid and not score and not action: raise HTTPException(status_code=400, detail="qid, score e action mancanti") rcpt = meta.get("rcpt", []) if isinstance(rcpt, str) and rcpt not in ("", "unknown"): rcpt = [rcpt] elif not isinstance(rcpt, list): rcpt = [] rcpt = [r.lower().strip() for r in rcpt if isinstance(r, str) and r.strip()] fuzzy = meta.get("fuzzy") if isinstance(fuzzy, list): fuzzy_json = json.dumps(fuzzy, ensure_ascii=False) elif fuzzy in (None, "", "unknown"): fuzzy_json = "[]" else: fuzzy_json = str(fuzzy) symbols = meta.get("symbols", {}) if isinstance(symbols, (dict, list)): symbols_json = json.dumps(symbols, ensure_ascii=False) else: symbols_json = str(symbols) if symbols else "[]" try: db_id = _process( qid=qid, server=server, subject=str(meta.get("subject") or ""), score=score, action=action, symbols_json=symbols_json, fuzzy_json=fuzzy_json, ip=str(meta.get("ip") or ""), mail_from=str(meta.get("from") or ""), rcpt_list=rcpt, size=int(meta.get("size") or 0), raw_email=raw_email, ) except Exception as exc: log.error(f"{qid} errore DB: {exc}") raise HTTPException(status_code=500, detail="Errore database") log.info(f"{qid} score:{score:.2f} action:'{action}' salvato [id:{db_id}]") return "Message saved" # ─── Endpoint: raw body + X-Rspamd-* headers (formatter = "default") ───────── @app.post("/api/metadata_importer", response_class=PlainTextResponse) async def metadata_importer( request: Request, _: None = Depends(_auth), ) -> str: raw_email = await request.body() if not raw_email: raise HTTPException(status_code=400, detail="Body vuoto") h = request.headers qid = h.get("x-rspamd-qid", "unknown") if qid and not re.match(r"^[a-zA-Z0-9]+$", qid): qid = "unknown" action = h.get("x-rspamd-action", "") server = request.query_params.get("server", SERVER_ALIAS) try: score = float(h.get("x-rspamd-score") or "0") except ValueError: score = 0.0 try: size = int(h.get("x-rspamd-size") or "0") except ValueError: size = 0 symbols_raw = h.get("x-rspamd-symbols", "[]") fuzzy_raw = h.get("x-rspamd-fuzzy", "[]") fuzzy_json = "[]" if fuzzy_raw in ("unknown", "", None) else fuzzy_raw rcpt_raw = h.get("x-rspamd-rcpt", "[]") try: rcpt = json.loads(rcpt_raw) if rcpt_raw not in ("", "unknown") else [] if not isinstance(rcpt, list): rcpt = [str(rcpt)] if rcpt else [] except Exception: rcpt = [] rcpt = [r.lower().strip() for r in rcpt if isinstance(r, str) and r.strip()] try: db_id = _process( qid=qid, server=server, subject=h.get("x-rspamd-subject", ""), score=score, action=action, symbols_json=symbols_raw or "[]", fuzzy_json=fuzzy_json, ip=h.get("x-rspamd-ip", ""), mail_from=h.get("x-rspamd-from", ""), rcpt_list=rcpt, size=size, raw_email=raw_email, ) except Exception as exc: log.error(f"{qid} errore DB: {exc}") raise HTTPException(status_code=500, detail="Errore database") log.info(f"{qid} score:{score:.2f} action:'{action}' salvato [id:{db_id}]") return "Message saved" # ─── Entry point ────────────────────────────────────────────────────────────── if __name__ == "__main__": import uvicorn uvicorn.run( "main:app", host=os.getenv("LISTEN_HOST", "127.0.0.1"), port=int(os.getenv("LISTEN_PORT", "8080")), reload=False, )