# app/workers/mail_send_worker.py from __future__ import annotations import smtplib, ssl, email.utils from email.message import EmailMessage from typing import List, Optional from pydantic import BaseModel, EmailStr from app.workers.base_worker import BaseWorker from app.config import Settings from app.db import SessionLocal # <- AsyncSession factory from app.models import Task, TaskStatus from sqlalchemy import update from app.job_registry import routing_key from app.config import settings class _SendPayload(BaseModel): mailbox_id: str recipients: List[EmailStr] subject: str body_text: Optional[str] = None body_html: Optional[str] = None sender: Optional[EmailStr] = None class MailSendWorker(BaseWorker): async def process(self, payload: dict, **meta) -> dict: """ Erwartet: task.job.payload (JSON) mit Feldern wie in _SendPayload. Setzt Task-Status/Result analog zu eurem BaseWorker-Pattern. """ p = _SendPayload(**payload) sender = str(p.sender or settings.smtp_sender_fallback) msg = EmailMessage() msg["From"] = sender msg["To"] = ", ".join([str(r) for r in p.recipients]) msg["Subject"] = p.subject msg["Date"] = email.utils.formatdate(localtime=True) msg["Message-ID"] = email.utils.make_msgid() if p.body_html and p.body_text: msg.set_content(p.body_text) msg.add_alternative(p.body_html, subtype="html") elif p.body_html: msg.add_alternative(p.body_html, subtype="html") else: msg.set_content(p.body_text or "") # SMTP – synchrones I/O im Worker-Kontext ist i.d.R. ok. context = ssl.create_default_context() with smtplib.SMTP(settings.smtp_host, settings.smtp_port, timeout=30) as server: if settings.smtp_starttls: server.starttls(context=context) if settings.smtp_user and settings.smtp_password: server.login(settings.smtp_user, settings.smtp_password) server.send_message(msg) return {"message_id": msg["Message-ID"]} async def main(): worker =MailSendWorker(routing_key("MailSend", "SendClean"), worker_id="mail-send-clean-1") await worker.run() if __name__ == "__main__": import asyncio as _a _a.run(main())