Primo commit
This commit is contained in:
@@ -0,0 +1,754 @@
|
||||
#!/usr/bin/env python3
|
||||
"""
|
||||
Postfix Content Filter con integrazione Rspamd
|
||||
Filtra i messaggi email attraverso Rspamd via HTTP API
|
||||
"""
|
||||
|
||||
import asyncio
|
||||
import re
|
||||
import argparse
|
||||
import configparser
|
||||
import logging
|
||||
import logging.handlers
|
||||
import sys
|
||||
import signal
|
||||
import aiohttp
|
||||
import email
|
||||
import os
|
||||
import uuid
|
||||
import contextvars
|
||||
|
||||
|
||||
# Il content filter riceve il messaggio gia' re-iniettato da Postfix sulla porta
|
||||
# locale, non via milter: l'unico modo per recuperare il queue id originale e'
|
||||
# leggerlo dall'header Received che Postfix stesso aggiunge nel formato
|
||||
# "by <host> (Postfix) with <proto> id <QUEUEID>". Senza Queue-ID nella
|
||||
# richiesta a Rspamd, task:get_queue_id() risulta vuoto e tutto finisce
|
||||
# etichettato "unknown" (es. nel path di quarantena del metadata_exporter).
|
||||
_POSTFIX_QUEUE_ID_RE = re.compile(r'\(Postfix\)\s+with\s+\S+\s+id\s+([0-9A-Za-z]+)')
|
||||
|
||||
|
||||
def _extract_queue_id(message):
|
||||
"""Recupera il queue id Postfix dal primo header Received del messaggio."""
|
||||
received = message.get_all('Received')
|
||||
if not received:
|
||||
return None
|
||||
match = _POSTFIX_QUEUE_ID_RE.search(received[0])
|
||||
return match.group(1) if match else None
|
||||
|
||||
|
||||
# Id della sessione SMTP corrente, isolato per task asyncio: ogni connessione
|
||||
# gira nel proprio Task (creato da asyncio.start_server), quindi il valore
|
||||
# impostato in handle_smtp_session non "perde" verso altre sessioni concorrenti.
|
||||
# asyncio.gather() copia il contesto corrente nei sotto-task che crea, quindi
|
||||
# anche i reinvii per i vari destinatari ereditano lo stesso id.
|
||||
_session_id_var = contextvars.ContextVar('session_id', default='-')
|
||||
|
||||
|
||||
class _SessionIdLogFilter(logging.Filter):
|
||||
"""Inietta l'id della sessione SMTP corrente in ogni record di log, per
|
||||
poter correlare tutte le righe relative al transito di una singola email."""
|
||||
|
||||
def filter(self, record):
|
||||
record.session_id = _session_id_var.get()
|
||||
return True
|
||||
|
||||
|
||||
class PostfixRspamdFilter:
|
||||
def __init__(self, config_file='filter_config.ini'):
|
||||
self.config = configparser.ConfigParser()
|
||||
self.config.read(config_file)
|
||||
|
||||
# Configurazione di default
|
||||
self.listen_port = int(self.config.get('server', 'listen_port', fallback='10026'))
|
||||
self.listen_host = self.config.get('server', 'listen_host', fallback='127.0.0.1')
|
||||
|
||||
# Configurazione Rspamd
|
||||
self.rspamd_host = self.config.get('rspamd', 'host', fallback='127.0.0.1')
|
||||
self.rspamd_port = int(self.config.get('rspamd', 'port', fallback='11333'))
|
||||
self.rspamd_timeout = int(self.config.get('rspamd', 'timeout', fallback='30'))
|
||||
|
||||
# Configurazione Postfix
|
||||
self.postfix_host = self.config.get('postfix', 'reinject_host', fallback='127.0.0.1')
|
||||
self.postfix_port = int(self.config.get('postfix', 'reinject_port', fallback='10027'))
|
||||
self.use_simple_smtp = self.config.getboolean('postfix', 'use_simple_smtp', fallback=True)
|
||||
|
||||
# Soglie spam (solo per header, non per blocco)
|
||||
self.spam_threshold = float(self.config.get('filtering', 'spam_threshold', fallback='5.0'))
|
||||
|
||||
# Limite connessioni concorrenti
|
||||
self.max_concurrent = int(self.config.get('server', 'max_concurrent', fallback='50'))
|
||||
|
||||
# Logging
|
||||
self.setup_logging()
|
||||
|
||||
# Server
|
||||
self.server = None
|
||||
self.session = None
|
||||
self.semaphore = None
|
||||
self.shutdown_event = None
|
||||
|
||||
def setup_logging(self):
|
||||
"""Configura il logging"""
|
||||
log_level = self.config.get('logging', 'level', fallback='INFO')
|
||||
log_file = self.config.get('logging', 'file', fallback='/var/log/postfix-rspamd-filter.log')
|
||||
log_max_bytes = int(self.config.get('logging', 'max_bytes', fallback=str(10 * 1024 * 1024)))
|
||||
log_backup_count = int(self.config.get('logging', 'backup_count', fallback='5'))
|
||||
|
||||
# Assicurati che la directory del log esista
|
||||
log_dir = os.path.dirname(log_file)
|
||||
if log_dir and not os.path.exists(log_dir):
|
||||
try:
|
||||
os.makedirs(log_dir, mode=0o755)
|
||||
except PermissionError:
|
||||
log_file = './postfix-rspamd-filter.log'
|
||||
print(f"Impossibile creare {log_dir}, usando {log_file}")
|
||||
|
||||
logging.basicConfig(
|
||||
level=getattr(logging, log_level.upper()),
|
||||
format='%(asctime)s - %(name)s - %(levelname)s - [session=%(session_id)s] - %(message)s',
|
||||
handlers=[
|
||||
logging.handlers.RotatingFileHandler(
|
||||
log_file, maxBytes=log_max_bytes, backupCount=log_backup_count
|
||||
),
|
||||
logging.StreamHandler(sys.stdout)
|
||||
]
|
||||
)
|
||||
self.logger = logging.getLogger('postfix-rspamd-filter')
|
||||
self.logger.addFilter(_SessionIdLogFilter())
|
||||
self.logger.info(f"Logging inizializzato - Level: {log_level}, File: {log_file}")
|
||||
|
||||
async def check_with_rspamd(self, message_content, destinatario, queue_id=None):
|
||||
"""Invia il messaggio a Rspamd per la verifica spam"""
|
||||
|
||||
self.logger.debug(f"Rspamd check_with_rspamd destinatario: {destinatario}")
|
||||
url = f"http://{self.rspamd_host}:{self.rspamd_port}/checkv2"
|
||||
|
||||
headers = {
|
||||
'Content-Type': 'message/rfc822',
|
||||
'Rcpt': destinatario,
|
||||
}
|
||||
if queue_id:
|
||||
headers['Queue-ID'] = queue_id
|
||||
|
||||
try:
|
||||
timeout = aiohttp.ClientTimeout(total=self.rspamd_timeout)
|
||||
|
||||
async with self.session.post(url, data=message_content,
|
||||
headers=headers, timeout=timeout) as response:
|
||||
|
||||
self.logger.debug(f"Rspamd response status: {response.status}")
|
||||
|
||||
if response.status == 200:
|
||||
result = await response.json()
|
||||
self.logger.debug(f"Rspamd result: {result}")
|
||||
return result
|
||||
else:
|
||||
response_text = await response.text()
|
||||
self.logger.error(f"Errore Rspamd: HTTP {response.status} - {response_text}")
|
||||
return None
|
||||
|
||||
except asyncio.TimeoutError:
|
||||
self.logger.error(f"Timeout connessione a Rspamd ({self.rspamd_timeout}s)")
|
||||
return None
|
||||
except aiohttp.ClientConnectorError as e:
|
||||
self.logger.error(f"Errore connessione a Rspamd: {e}")
|
||||
return None
|
||||
except Exception as e:
|
||||
self.logger.error(f"Errore comunicazione Rspamd: {e}")
|
||||
return None
|
||||
|
||||
def add_spam_headers(self, message, rspamd_result):
|
||||
"""Aggiunge header spam al messaggio"""
|
||||
if not rspamd_result:
|
||||
# Se Rspamd non è disponibile, aggiungi header di fallback
|
||||
message['X-Spam-Status'] = 'Unknown (Rspamd not available)'
|
||||
message['X-Spam-Checker-Version'] = 'Postfix Content Filter v1.0'
|
||||
return message
|
||||
|
||||
score = rspamd_result.get('score', 0)
|
||||
action = rspamd_result.get('action', 'no action')
|
||||
symbols = rspamd_result.get('symbols', {})
|
||||
required_score = rspamd_result.get('required_score', self.spam_threshold)
|
||||
|
||||
# Header X-Spam-Status (compatibile SpamAssassin). Il confronto usa
|
||||
# required_score (la soglia effettivamente applicata da Rspamd) e non
|
||||
# self.spam_threshold, altrimenti se le due soglie divergono l'header
|
||||
# può risultare contraddittorio, es. "No, score=6.00 required=5.00".
|
||||
is_spam = score >= required_score
|
||||
spam_status = f"{'Yes' if is_spam else 'No'}, score={score:.2f} required={required_score:.2f}"
|
||||
message['X-Spam-Status'] = spam_status
|
||||
|
||||
# Header X-Spam-Score (punteggio numerico)
|
||||
message['X-Spam-Score'] = f"{score:.2f}"
|
||||
|
||||
# Header X-Spam-Level (asterischi per visualizzazione)
|
||||
if score > 0:
|
||||
level_stars = '*' * min(int(score), 20) # Max 20 asterischi
|
||||
message['X-Spam-Level'] = level_stars
|
||||
|
||||
# Header X-Spam-Action (azione suggerita da Rspamd)
|
||||
message['X-Spam-Action'] = action
|
||||
|
||||
# Header con simboli e punteggi dettagliati
|
||||
if symbols:
|
||||
symbol_list = []
|
||||
for symbol, info in symbols.items():
|
||||
symbol_score = info.get('score', 0)
|
||||
symbol_list.append(f"{symbol}({symbol_score:.2f})")
|
||||
|
||||
# Limita la lunghezza dell'header per evitare problemi
|
||||
symbols_text = ', '.join(symbol_list)
|
||||
if len(symbols_text) > 1000:
|
||||
symbols_text = symbols_text[:997] + '...'
|
||||
|
||||
message['X-Spam-Symbols'] = symbols_text
|
||||
|
||||
# Header con informazioni aggiuntive da Rspamd
|
||||
if 'message-id' in rspamd_result:
|
||||
message['X-Rspamd-Queue-Id'] = rspamd_result['message-id']
|
||||
|
||||
if 'time_real' in rspamd_result:
|
||||
scan_time = rspamd_result['time_real']
|
||||
message['X-Rspamd-Scan-Time'] = f"{scan_time:.3f}s"
|
||||
|
||||
# Header identificativo del filter
|
||||
message['X-Spam-Checker-Version'] = 'Postfix-Rspamd Content Filter v1.0'
|
||||
|
||||
# Se è greylist, aggiungi header specifico
|
||||
if action == 'greylist':
|
||||
message['X-Spam-Greylist'] = 'Deferred by Rspamd (passed through by filter)'
|
||||
|
||||
self.logger.debug(f"Aggiunti header spam: score={score}, action={action}, symbols={len(symbols)}")
|
||||
|
||||
return message
|
||||
|
||||
async def reinject_message_smtp(self, sender, recipient, message_data):
|
||||
"""Reinvia il messaggio processato a Postfix tramite SMTP"""
|
||||
try:
|
||||
import aiosmtplib
|
||||
|
||||
self.logger.debug(f"Reinviando messaggio da {sender} a {recipient}")
|
||||
|
||||
# Connessione SMTP senza TLS per connessioni locali
|
||||
smtp = aiosmtplib.SMTP(
|
||||
hostname=self.postfix_host,
|
||||
port=self.postfix_port,
|
||||
use_tls=False, # Disabilita TLS
|
||||
start_tls=False # Disabilita STARTTLS
|
||||
)
|
||||
|
||||
await smtp.connect()
|
||||
|
||||
# Non tentare STARTTLS per connessioni locali
|
||||
await smtp.sendmail(sender, [recipient], message_data)
|
||||
await smtp.quit()
|
||||
|
||||
self.logger.info(f"Messaggio reinviato: {sender} -> {recipient}")
|
||||
|
||||
except Exception as e:
|
||||
self.logger.error(f"Errore reinvio messaggio: {e}")
|
||||
# Prova con una connessione più semplice se fallisce
|
||||
try:
|
||||
await self.reinject_message_simple(sender, recipient, message_data)
|
||||
except Exception as e2:
|
||||
self.logger.error(f"Fallito anche reinvio semplice: {e2}")
|
||||
raise e
|
||||
|
||||
@staticmethod
|
||||
def _smtp_quote_data(data):
|
||||
"""Normalizza i fine riga a CRLF e raddoppia i punti a inizio riga
|
||||
(RFC 5321 4.5.2 - trasparenza). Senza questo passaggio, una riga del
|
||||
corpo che inizia per '.' verrebbe interpretata dal server come fine
|
||||
del messaggio (se la riga è solo '.') oppure troncata di un carattere,
|
||||
corrompendo silenziosamente il contenuto."""
|
||||
data = re.sub(rb'\r\n|\r|\n', b'\r\n', data)
|
||||
return re.sub(rb'(?m)^\.', b'..', data)
|
||||
|
||||
async def _read_smtp_response(self, reader):
|
||||
"""Legge una risposta SMTP (anche multilinea) e restituisce (codice, testo)"""
|
||||
lines = []
|
||||
while True:
|
||||
raw_line = await asyncio.wait_for(reader.readline(), timeout=60)
|
||||
if not raw_line:
|
||||
raise ConnectionError("Connessione chiusa inattesa durante la lettura della risposta SMTP")
|
||||
line = raw_line.decode('ascii', errors='replace').rstrip('\r\n')
|
||||
lines.append(line)
|
||||
if len(line) >= 4 and line[3] == '-':
|
||||
continue # riga di continuazione (es. "250-...")
|
||||
break
|
||||
text = '\n'.join(lines)
|
||||
code = int(lines[-1][:3]) if lines[-1][:3].isdigit() else 0
|
||||
return code, text
|
||||
|
||||
async def reinject_message_simple(self, sender, recipient, message_data):
|
||||
"""Reinvio semplice tramite socket TCP diretto"""
|
||||
writer = None
|
||||
try:
|
||||
self.logger.debug("Tentando reinvio con socket diretto")
|
||||
|
||||
reader, writer = await asyncio.open_connection(
|
||||
self.postfix_host, self.postfix_port
|
||||
)
|
||||
|
||||
# Leggi banner
|
||||
code, text = await self._read_smtp_response(reader)
|
||||
self.logger.debug(f"Banner SMTP: {text}")
|
||||
if code != 220:
|
||||
raise RuntimeError(f"Banner SMTP inatteso: {text}")
|
||||
|
||||
# HELO
|
||||
writer.write(b"HELO localhost\r\n")
|
||||
await writer.drain()
|
||||
code, text = await self._read_smtp_response(reader)
|
||||
self.logger.debug(f"HELO response: {text}")
|
||||
if code != 250:
|
||||
raise RuntimeError(f"HELO rifiutato: {text}")
|
||||
|
||||
# MAIL FROM
|
||||
writer.write(f"MAIL FROM:<{sender}>\r\n".encode())
|
||||
await writer.drain()
|
||||
code, text = await self._read_smtp_response(reader)
|
||||
self.logger.debug(f"MAIL FROM response: {text}")
|
||||
if code != 250:
|
||||
raise RuntimeError(f"MAIL FROM rifiutato: {text}")
|
||||
|
||||
# RCPT TO
|
||||
writer.write(f"RCPT TO:<{recipient}>\r\n".encode())
|
||||
await writer.drain()
|
||||
code, text = await self._read_smtp_response(reader)
|
||||
self.logger.debug(f"RCPT TO response: {text}")
|
||||
if code not in (250, 251):
|
||||
raise RuntimeError(f"RCPT TO rifiutato: {text}")
|
||||
|
||||
# DATA
|
||||
writer.write(b"DATA\r\n")
|
||||
await writer.drain()
|
||||
code, text = await self._read_smtp_response(reader)
|
||||
self.logger.debug(f"DATA response: {text}")
|
||||
if code != 354:
|
||||
raise RuntimeError(f"DATA rifiutato: {text}")
|
||||
|
||||
# Invia il messaggio (già in bytes: nessuna re-encode che corromperebbe
|
||||
# contenuto 8bit non-UTF8), applicando il dot-stuffing richiesto dal
|
||||
# protocollo SMTP prima del terminatore di fine messaggio
|
||||
writer.write(self._smtp_quote_data(message_data))
|
||||
writer.write(b"\r\n.\r\n")
|
||||
await writer.drain()
|
||||
code, text = await self._read_smtp_response(reader)
|
||||
self.logger.debug(f"Message response: {text}")
|
||||
if code != 250:
|
||||
raise RuntimeError(f"Messaggio rifiutato dal server: {text}")
|
||||
|
||||
# QUIT
|
||||
writer.write(b"QUIT\r\n")
|
||||
await writer.drain()
|
||||
try:
|
||||
code, text = await self._read_smtp_response(reader)
|
||||
self.logger.debug(f"QUIT response: {text}")
|
||||
except Exception:
|
||||
pass # il messaggio è già stato accettato, un QUIT non confermato non è un errore
|
||||
|
||||
self.logger.info(f"Messaggio reinviato con socket diretto: {sender} -> {recipient}")
|
||||
|
||||
except Exception as e:
|
||||
self.logger.error(f"Errore reinvio con socket diretto: {e}")
|
||||
raise
|
||||
finally:
|
||||
if writer is not None:
|
||||
try:
|
||||
writer.close()
|
||||
await writer.wait_closed()
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
async def handle_smtp_session(self, reader, writer):
|
||||
"""Gestisce una sessione SMTP per il content filter"""
|
||||
_session_id_var.set(uuid.uuid4().hex[:8])
|
||||
async with self.semaphore:
|
||||
await self._handle_smtp_session(reader, writer)
|
||||
|
||||
async def _handle_smtp_session(self, reader, writer):
|
||||
client_addr = writer.get_extra_info('peername')
|
||||
self.logger.info(f"Nuova connessione da {client_addr}")
|
||||
|
||||
try:
|
||||
# Invia il banner SMTP
|
||||
writer.write(b"220 Content-Filter Ready\r\n")
|
||||
await writer.drain()
|
||||
|
||||
sender = None
|
||||
recipient = None
|
||||
message_data = []
|
||||
in_data_mode = False
|
||||
|
||||
while True:
|
||||
try:
|
||||
raw_line = await asyncio.wait_for(reader.readline(), timeout=60)
|
||||
if not raw_line:
|
||||
self.logger.info("Connessione chiusa dal client")
|
||||
break
|
||||
|
||||
if in_data_mode:
|
||||
# Corpo del messaggio: lavoriamo sui byte grezzi, senza
|
||||
# decodificarli come UTF-8. Il corpo può legittimamente
|
||||
# contenere contenuto 8bit in charset diversi (es. ISO-8859-1):
|
||||
# un decode/encode UTF-8 con errors='ignore' lo corromperebbe
|
||||
# silenziosamente.
|
||||
content = raw_line.rstrip(b'\r\n')
|
||||
if content == b'.':
|
||||
# Fine del messaggio
|
||||
in_data_mode = False
|
||||
|
||||
message_bytes = b"".join(message_data)
|
||||
self.logger.info(f"Messaggio ricevuto completamente ({len(message_bytes)} bytes)")
|
||||
|
||||
# Processa il messaggio
|
||||
status = await self.process_message(sender, recipient, message_bytes)
|
||||
|
||||
if status == 'rejected':
|
||||
# 250 e non 5xx/4xx: il messaggio è già stato accettato da
|
||||
# Postfix a monte ed è già in quarantena lato Rspamd; la nota
|
||||
# tra parentesi finisce nel log di Postfix (status=sent (...))
|
||||
# rendendo visibile il reject anche lì.
|
||||
writer.write(b"250 OK Message processed (rspamd rejected)\r\n")
|
||||
self.logger.warning("Messaggio rifiutato da Rspamd, non reiniettato")
|
||||
elif status == 'accepted':
|
||||
writer.write(b"250 OK Message processed\r\n")
|
||||
self.logger.info("Messaggio processato con successo")
|
||||
else:
|
||||
writer.write(b"450 Temporary failure\r\n")
|
||||
self.logger.error("Errore nel processamento del messaggio")
|
||||
await writer.drain()
|
||||
|
||||
# Reset per il prossimo messaggio
|
||||
sender = None
|
||||
recipient = None
|
||||
message_data = []
|
||||
|
||||
else:
|
||||
# Accumula i dati del messaggio
|
||||
if content.startswith(b'.'):
|
||||
content = content[1:] # Rimuovi il punto di escape
|
||||
message_data.append(content + b"\r\n")
|
||||
else:
|
||||
# Comandi SMTP: l'envelope è sempre testo ASCII per specifica
|
||||
line = raw_line.decode('utf-8', errors='ignore').rstrip('\r\n')
|
||||
if line: # Log solo linee non vuote
|
||||
self.logger.debug(f"Ricevuto da {client_addr}: {line}")
|
||||
|
||||
if not line.strip():
|
||||
continue
|
||||
|
||||
parts = line.strip().split()
|
||||
if not parts:
|
||||
writer.write(b"500 Empty command\r\n")
|
||||
await writer.drain()
|
||||
continue
|
||||
|
||||
cmd = parts[0].upper()
|
||||
|
||||
if cmd == "HELO" or cmd == "EHLO":
|
||||
writer.write(b"250 OK\r\n")
|
||||
|
||||
elif cmd == "MAIL":
|
||||
# MAIL FROM:<sender> [parametri ESMTP, es. SIZE=, BODY=8BITMIME]
|
||||
try:
|
||||
line_upper = line.upper()
|
||||
if "FROM:" in line_upper:
|
||||
from_index = line_upper.find("FROM:")
|
||||
from_part = line[from_index + 5:].strip()
|
||||
# Estrae solo l'indirizzo tra < >, ignorando eventuali
|
||||
# parametri ESMTP successivi (es. "<a@b> SIZE=123")
|
||||
match = re.match(r'<([^>]*)>', from_part)
|
||||
sender = match.group(1) if match else from_part.split()[0] if from_part.split() else ''
|
||||
self.logger.debug(f"Sender impostato: {sender}")
|
||||
else:
|
||||
self.logger.warning(f"MAIL FROM senza FROM: {line}")
|
||||
writer.write(b"250 OK\r\n")
|
||||
except Exception as e:
|
||||
self.logger.error(f"Errore parsing MAIL FROM '{line}': {e}")
|
||||
writer.write(b"500 Syntax error in MAIL FROM\r\n")
|
||||
|
||||
elif cmd == "RCPT":
|
||||
# RCPT TO:<recipient> [parametri ESMTP, es. NOTIFY=, ORCPT=]
|
||||
# Il filtro è configurato per ricevere un solo destinatario per
|
||||
# transazione (consegna Postfix già frazionata a monte): una
|
||||
# seconda RCPT viene rifiutata invece di essere accettata in
|
||||
# silenzio, per non verificare/reiniettare con Rspamd solo il
|
||||
# primo indirizzo ignorando gli altri.
|
||||
if recipient is not None:
|
||||
self.logger.warning(f"RCPT TO aggiuntivo rifiutato (atteso un solo destinatario): {line}")
|
||||
writer.write(b"452 4.5.3 Too many recipients\r\n")
|
||||
else:
|
||||
try:
|
||||
line_upper = line.upper()
|
||||
if "TO:" in line_upper:
|
||||
to_index = line_upper.find("TO:")
|
||||
to_part = line[to_index + 3:].strip()
|
||||
# Estrae solo l'indirizzo tra < >, ignorando eventuali
|
||||
# parametri ESMTP successivi
|
||||
match = re.match(r'<([^>]*)>', to_part)
|
||||
recipient = match.group(1) if match else to_part.split()[0] if to_part.split() else ''
|
||||
self.logger.debug(f"Recipient impostato: {recipient}")
|
||||
else:
|
||||
self.logger.warning(f"RCPT TO senza TO: {line}")
|
||||
writer.write(b"250 OK\r\n")
|
||||
except Exception as e:
|
||||
self.logger.error(f"Errore parsing RCPT TO '{line}': {e}")
|
||||
writer.write(b"500 Syntax error in RCPT TO\r\n")
|
||||
|
||||
elif cmd == "DATA":
|
||||
if sender is None or recipient is None:
|
||||
writer.write(b"503 Bad sequence of commands\r\n")
|
||||
else:
|
||||
writer.write(b"354 Start mail input; end with <CRLF>.<CRLF>\r\n")
|
||||
in_data_mode = True
|
||||
|
||||
elif cmd == "QUIT":
|
||||
writer.write(b"221 Bye\r\n")
|
||||
break
|
||||
|
||||
elif cmd == "RSET":
|
||||
sender = None
|
||||
recipient = None
|
||||
message_data = []
|
||||
in_data_mode = False
|
||||
writer.write(b"250 OK\r\n")
|
||||
|
||||
elif cmd == "NOOP":
|
||||
writer.write(b"250 OK\r\n")
|
||||
|
||||
else:
|
||||
self.logger.warning(f"Comando non riconosciuto: {cmd}")
|
||||
writer.write(b"500 Command not recognized\r\n")
|
||||
|
||||
await writer.drain()
|
||||
|
||||
except asyncio.TimeoutError:
|
||||
self.logger.warning("Timeout connessione SMTP")
|
||||
break
|
||||
|
||||
except Exception as e:
|
||||
self.logger.error(f"Errore sessione SMTP: {e}")
|
||||
|
||||
finally:
|
||||
try:
|
||||
writer.close()
|
||||
await writer.wait_closed()
|
||||
except Exception:
|
||||
pass
|
||||
self.logger.info(f"Connessione chiusa per {client_addr}")
|
||||
|
||||
async def process_message(self, sender, recipient, message_data):
|
||||
"""Processa il messaggio con Rspamd e lo reinvia (modalità passiva - nessun blocco,
|
||||
salvo il caso action=reject: il messaggio è già in quarantena lato Rspamd tramite
|
||||
il modulo metadata_exporter, quindi qui viene solo confermato a Postfix senza
|
||||
rimetterlo in coda).
|
||||
|
||||
Restituisce 'accepted', 'rejected' oppure 'error': il chiamante usa questo stato
|
||||
per scegliere il testo della risposta SMTP, che finisce anche nel log di Postfix."""
|
||||
try:
|
||||
self.logger.info(f"Processando messaggio da {sender} per {recipient}")
|
||||
|
||||
# Parsing su bytes: preserva corpi 8bit non-UTF8 senza passare per una
|
||||
# decodifica/encodifica testuale lossy. Fatto qui, prima della verifica
|
||||
# Rspamd, per poter estrarre il queue id Postfix dall'header Received.
|
||||
message = email.message_from_bytes(message_data)
|
||||
queue_id = _extract_queue_id(message)
|
||||
if queue_id:
|
||||
self.logger.debug(f"Queue id Postfix rilevato: {queue_id}")
|
||||
else:
|
||||
self.logger.warning("Queue id Postfix non rilevato nell'header Received")
|
||||
|
||||
# Verifica con Rspamd
|
||||
rspamd_result = await self.check_with_rspamd(message_data, recipient, queue_id)
|
||||
|
||||
if rspamd_result:
|
||||
score = rspamd_result.get('score', 0)
|
||||
action = rspamd_result.get('action', 'no action')
|
||||
|
||||
self.logger.info(f"Rspamd score: {score}, action: {action}")
|
||||
|
||||
if action == 'reject':
|
||||
# Rspamd ha già salvato il messaggio in quarantena (metadata_exporter):
|
||||
# confermiamo a Postfix senza reiniettarlo, per non consegnarlo due volte.
|
||||
self.logger.warning(
|
||||
f"Messaggio da {sender} per {recipient} non reiniettato "
|
||||
f"(Rspamd action=reject, score={score}): già in quarantena"
|
||||
)
|
||||
return 'rejected'
|
||||
else:
|
||||
self.logger.warning("Rspamd non disponibile, procedo senza analisi")
|
||||
|
||||
# Aggiungi header spam (anche se Rspamd non risponde)
|
||||
message = self.add_spam_headers(message, rspamd_result)
|
||||
processed_message = message.as_bytes()
|
||||
|
||||
# Reinvia il messaggio processato al destinatario
|
||||
reinject = self.reinject_message_simple if self.use_simple_smtp else self.reinject_message_smtp
|
||||
await reinject(sender, recipient, processed_message)
|
||||
|
||||
self.logger.info(f"Messaggio processato e reinviato con successo")
|
||||
return 'accepted'
|
||||
|
||||
except Exception as e:
|
||||
self.logger.error(f"Errore processamento messaggio: {e}")
|
||||
|
||||
# Anche in caso di errore, tenta di reinviare il messaggio originale
|
||||
try:
|
||||
self.logger.info("Tentando reinvio del messaggio originale senza modifiche")
|
||||
reinject = self.reinject_message_simple if self.use_simple_smtp else self.reinject_message_smtp
|
||||
await reinject(sender, recipient, message_data)
|
||||
self.logger.info("Messaggio originale reinviato con successo")
|
||||
return 'accepted'
|
||||
except Exception as e2:
|
||||
self.logger.error(f"Errore fatale nel reinvio del messaggio originale: {e2}")
|
||||
return 'error'
|
||||
|
||||
async def test_rspamd_connection(self):
|
||||
"""Testa la connessione a Rspamd"""
|
||||
self.logger.info("Testando connessione a Rspamd...")
|
||||
|
||||
try:
|
||||
url = f"http://{self.rspamd_host}:{self.rspamd_port}/ping"
|
||||
timeout = aiohttp.ClientTimeout(total=5)
|
||||
|
||||
async with self.session.get(url, timeout=timeout) as response:
|
||||
if response.status == 200:
|
||||
self.logger.info("✓ Connessione a Rspamd OK")
|
||||
return True
|
||||
else:
|
||||
self.logger.warning(f"Rspamd risponde con status {response.status}")
|
||||
return False
|
||||
|
||||
except Exception as e:
|
||||
self.logger.error(f"✗ Impossibile connettersi a Rspamd: {e}")
|
||||
self.logger.error("Verifica che Rspamd sia in esecuzione e configurato correttamente")
|
||||
return False
|
||||
|
||||
async def start_server(self):
|
||||
"""Avvia il server"""
|
||||
# Sessione HTTP condivisa con connection pool verso Rspamd
|
||||
connector = aiohttp.TCPConnector(ssl=False, limit=100)
|
||||
self.session = aiohttp.ClientSession(connector=connector)
|
||||
|
||||
# Semaforo per limitare le connessioni SMTP concorrenti
|
||||
self.semaphore = asyncio.Semaphore(self.max_concurrent)
|
||||
|
||||
# Evento di shutdown: serve_forever() non si sblocca con server.close(),
|
||||
# serve un segnale esplicito per uscire dal blocco "async with" sotto
|
||||
self.shutdown_event = asyncio.Event()
|
||||
|
||||
# Testa la connessione a Rspamd
|
||||
await self.test_rspamd_connection()
|
||||
|
||||
self.server = await asyncio.start_server(
|
||||
self.handle_smtp_session,
|
||||
self.listen_host,
|
||||
self.listen_port
|
||||
)
|
||||
|
||||
addr = self.server.sockets[0].getsockname()
|
||||
self.logger.info(f"Server avviato su {addr[0]}:{addr[1]}")
|
||||
self.logger.info(f"Configurato per Rspamd su {self.rspamd_host}:{self.rspamd_port}")
|
||||
self.logger.info(f"Reinvio messaggi a {self.postfix_host}:{self.postfix_port}")
|
||||
|
||||
async with self.server:
|
||||
await self.shutdown_event.wait()
|
||||
|
||||
async def stop_server(self):
|
||||
"""Ferma il server"""
|
||||
if self.shutdown_event and not self.shutdown_event.is_set():
|
||||
self.shutdown_event.set()
|
||||
|
||||
if self.server:
|
||||
self.server.close()
|
||||
await self.server.wait_closed()
|
||||
|
||||
if self.session and not self.session.closed:
|
||||
await self.session.close()
|
||||
|
||||
self.logger.info("Server fermato")
|
||||
|
||||
def create_default_config(self, config_file):
|
||||
"""Crea un file di configurazione di default"""
|
||||
config = configparser.ConfigParser()
|
||||
|
||||
config['server'] = {
|
||||
'listen_host': '127.0.0.1',
|
||||
'listen_port': '10026',
|
||||
'max_concurrent': '50',
|
||||
}
|
||||
|
||||
config['rspamd'] = {
|
||||
'host': '127.0.0.1',
|
||||
'port': '11333',
|
||||
'timeout': '30',
|
||||
}
|
||||
|
||||
config['postfix'] = {
|
||||
'reinject_host': '127.0.0.1',
|
||||
'reinject_port': '10027',
|
||||
'use_simple_smtp': 'True',
|
||||
}
|
||||
|
||||
config['filtering'] = {
|
||||
'spam_threshold': '5.0',
|
||||
'passthrough_mode': 'True',
|
||||
}
|
||||
|
||||
config['logging'] = {
|
||||
'level': 'INFO',
|
||||
'file': '/var/log/postfix-rspamd-filter.log',
|
||||
'max_bytes': str(10 * 1024 * 1024),
|
||||
'backup_count': '5',
|
||||
}
|
||||
|
||||
with open(config_file, 'w') as f:
|
||||
config.write(f)
|
||||
|
||||
print(f"File di configurazione creato: {config_file}")
|
||||
|
||||
|
||||
async def main():
|
||||
parser = argparse.ArgumentParser(description='Postfix Content Filter con Rspamd')
|
||||
parser.add_argument('-c', '--config', default='filter_config.ini',
|
||||
help='File di configurazione')
|
||||
parser.add_argument('--create-config', action='store_true',
|
||||
help='Crea file di configurazione di default')
|
||||
args = parser.parse_args()
|
||||
|
||||
filter_instance = PostfixRspamdFilter(args.config)
|
||||
|
||||
if args.create_config:
|
||||
filter_instance.create_default_config(args.config)
|
||||
return
|
||||
|
||||
# Gestione segnali per shutdown pulito
|
||||
def signal_handler():
|
||||
asyncio.create_task(filter_instance.stop_server())
|
||||
|
||||
for sig in (signal.SIGTERM, signal.SIGINT):
|
||||
asyncio.get_event_loop().add_signal_handler(sig, signal_handler)
|
||||
|
||||
try:
|
||||
await filter_instance.start_server()
|
||||
except KeyboardInterrupt:
|
||||
filter_instance.logger.info("Shutdown richiesto dall'utente")
|
||||
except Exception as e:
|
||||
filter_instance.logger.error(f"Errore fatale: {e}")
|
||||
finally:
|
||||
await filter_instance.stop_server()
|
||||
|
||||
|
||||
if __name__ == '__main__':
|
||||
# Installa dipendenze se necessario
|
||||
try:
|
||||
import aiohttp
|
||||
import aiosmtplib
|
||||
except ImportError:
|
||||
print("Installare dipendenze: pip install aiohttp aiosmtplib")
|
||||
sys.exit(1)
|
||||
|
||||
asyncio.run(main())
|
||||
Reference in New Issue
Block a user