# app/workers/mail_store_worker.py from __future__ import annotations import base64, email.utils from email import policy from email.parser import BytesParser from datetime import datetime from uuid import UUID from sqlalchemy import select, update from app.workers.base_worker import BaseWorker from app.db import SessionLocal from app.models import Message, Task, TaskStatus class MailStoreWorker(BaseWorker): routing_key = "mail.store" async def process(self, payload: dict, **meta) -> dict: raw = base64.b64decode(payload["raw_eml_b64"]) mailbox_id = UUID(payload["mailbox_id"]) msg = BytesParser(policy=policy.default).parsebytes(raw) message_id = msg.get("Message-ID") subject = msg.get("Subject") from_addr = msg.get("From") to_addr = msg.get("To") dt = None try: dt = email.utils.parsedate_to_datetime(msg.get("Date")) if msg.get("Date") else None except Exception: pass async with SessionLocal() as session: # Idempotenz: bei Message-ID prüfen if message_id: q = select(Message).where(Message.message_id == message_id) existing = (await session.execute(q)).scalar_one_or_none() if existing: result = {"id": str(existing.id), "message_id": message_id, "dedup": True} await session.execute( update(Task) .where(Task.id == task.id) .values(status=TaskStatus.FINISHED, result=result, worker_id=self.worker_id) ) await session.commit() return result m = Message( mailbox_id=mailbox_id, message_id=message_id, subject=subject, from_addr=from_addr, to_addr=to_addr, date_header=dt, received_at=datetime.utcnow(), raw_eml=raw, ) session.add(m) await session.commit() # id verfügbar, da PG/UUID serverseitig result = {"id": str(m.id), "message_id": message_id} await session.execute( update(Task) .where(Task.id == task.id) .values(status=TaskStatus.FINISHED, result=result, worker_id=self.worker_id) ) await session.commit() return result async def main(): worker =MailStoreWorker(routing_key("MailSend", "ImportMessage"), worker_id="mail-import-1") await worker.run() if __name__ == "__main__": import asyncio as _a _a.run(main())