233 lines
8.3 KiB
Python
233 lines
8.3 KiB
Python
"""In-process background queue for outbound user emails.
|
|
|
|
The queue keeps SMTP delivery off the request path so message composition
|
|
returns quickly and the rest of the API remains responsive while the worker
|
|
thread performs the blocking SMTP call.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import logging
|
|
import queue
|
|
import threading
|
|
import time
|
|
import uuid
|
|
from dataclasses import dataclass, field
|
|
from typing import Any
|
|
|
|
from media_library_viewer_api.services.mailer import EmailAttachment, describe_smtp_error, send_email_message
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
|
|
@dataclass(frozen=True)
|
|
class QueuedEmailMessage:
|
|
"""A queued outbound email request."""
|
|
|
|
request_id: str
|
|
settings: Any
|
|
recipients: list[str]
|
|
subject: str
|
|
html_body: str
|
|
text_body: str
|
|
attachments: list[EmailAttachment] = field(default_factory=list)
|
|
created_at: float = field(default_factory=time.time)
|
|
|
|
|
|
class MailQueue:
|
|
"""Single-worker in-process queue for SMTP delivery."""
|
|
|
|
def __init__(self) -> None:
|
|
self._queue: queue.Queue[QueuedEmailMessage | None] = queue.Queue()
|
|
self._thread: threading.Thread | None = None
|
|
self._stop_event = threading.Event()
|
|
self._lock = threading.Lock()
|
|
self._pending_count = 0
|
|
self._active_request_id: str | None = None
|
|
self._last_request_id: str | None = None
|
|
self._last_result: str | None = None
|
|
self._last_error = ""
|
|
self._last_error_at: float | None = None
|
|
self._last_success_at: float | None = None
|
|
self._last_activity_at: float | None = None
|
|
self._sent_count = 0
|
|
self._failed_count = 0
|
|
|
|
def start(self) -> None:
|
|
"""Start the worker thread if it is not already running."""
|
|
with self._lock:
|
|
if self._thread and self._thread.is_alive():
|
|
return
|
|
self._stop_event.clear()
|
|
self._thread = threading.Thread(target=self._run, name="mail-queue-worker", daemon=True)
|
|
self._thread.start()
|
|
logger.info("Mail queue worker started")
|
|
|
|
def stop(self, timeout: float = 5.0) -> None:
|
|
"""Stop the worker thread and wait briefly for shutdown."""
|
|
with self._lock:
|
|
thread = self._thread
|
|
if not thread:
|
|
return
|
|
self._stop_event.set()
|
|
self._queue.put(None)
|
|
thread.join(timeout=timeout)
|
|
if thread.is_alive():
|
|
logger.warning("Mail queue worker did not stop within %.1fs", timeout)
|
|
else:
|
|
logger.info("Mail queue worker stopped")
|
|
with self._lock:
|
|
if self._thread is thread:
|
|
self._thread = None
|
|
|
|
def enqueue(
|
|
self,
|
|
*,
|
|
settings: Any,
|
|
recipients: list[str],
|
|
subject: str,
|
|
html_body: str,
|
|
text_body: str = "",
|
|
attachments: list[EmailAttachment] | None = None,
|
|
) -> str:
|
|
"""Queue an outbound email and return a request identifier."""
|
|
request_id = uuid.uuid4().hex
|
|
message = QueuedEmailMessage(
|
|
request_id=request_id,
|
|
settings=settings,
|
|
recipients=list(recipients),
|
|
subject=subject,
|
|
html_body=html_body,
|
|
text_body=text_body,
|
|
attachments=list(attachments or []),
|
|
)
|
|
with self._lock:
|
|
self._pending_count += 1
|
|
self._last_request_id = request_id
|
|
self._last_result = "queued"
|
|
self._last_activity_at = time.time()
|
|
self._queue.put(message)
|
|
logger.info(
|
|
"Queued email request_id=%s recipients=%s attachments=%s subject=%s",
|
|
request_id,
|
|
len(message.recipients),
|
|
len(message.attachments),
|
|
subject,
|
|
)
|
|
return request_id
|
|
|
|
def status(self) -> dict[str, Any]:
|
|
"""Return a snapshot of the queue state for health/status endpoints."""
|
|
with self._lock:
|
|
worker_running = bool(self._thread and self._thread.is_alive())
|
|
stop_requested = self._stop_event.is_set()
|
|
pending_count = self._pending_count
|
|
active_request_id = self._active_request_id
|
|
last_request_id = self._last_request_id
|
|
last_result = self._last_result
|
|
last_error = self._last_error
|
|
last_error_at = self._last_error_at
|
|
last_success_at = self._last_success_at
|
|
last_activity_at = self._last_activity_at
|
|
sent_count = self._sent_count
|
|
failed_count = self._failed_count
|
|
|
|
if not worker_running:
|
|
state = "stopped" if stop_requested else "error"
|
|
elif active_request_id or pending_count > 0:
|
|
state = "busy"
|
|
elif last_result == "failed" and last_error:
|
|
state = "error"
|
|
else:
|
|
state = "idle"
|
|
|
|
return {
|
|
"state": state,
|
|
"worker_running": worker_running,
|
|
"stop_requested": stop_requested,
|
|
"pending_count": pending_count,
|
|
"active_request_id": active_request_id,
|
|
"last_request_id": last_request_id,
|
|
"last_result": last_result,
|
|
"last_error": last_error,
|
|
"last_error_at": last_error_at,
|
|
"last_success_at": last_success_at,
|
|
"last_activity_at": last_activity_at,
|
|
"sent_count": sent_count,
|
|
"failed_count": failed_count,
|
|
}
|
|
|
|
def _run(self) -> None:
|
|
while not self._stop_event.is_set():
|
|
try:
|
|
message = self._queue.get(timeout=0.5)
|
|
except queue.Empty:
|
|
continue
|
|
|
|
try:
|
|
if message is None:
|
|
continue
|
|
|
|
with self._lock:
|
|
self._pending_count = max(0, self._pending_count - 1)
|
|
self._active_request_id = message.request_id
|
|
self._last_request_id = message.request_id
|
|
self._last_result = "sending"
|
|
self._last_activity_at = time.time()
|
|
|
|
logger.info(
|
|
"Mail queue sending request_id=%s recipients=%s attachments=%s subject=%s",
|
|
message.request_id,
|
|
len(message.recipients),
|
|
len(message.attachments),
|
|
message.subject,
|
|
)
|
|
result: dict[str, Any] = send_email_message(
|
|
message.settings,
|
|
recipients=message.recipients,
|
|
subject=message.subject,
|
|
html_body=message.html_body,
|
|
text_body=message.text_body,
|
|
attachments=message.attachments,
|
|
)
|
|
with self._lock:
|
|
self._active_request_id = None
|
|
self._last_result = "sent"
|
|
self._last_success_at = time.time()
|
|
self._last_activity_at = self._last_success_at
|
|
self._last_error = ""
|
|
self._last_error_at = None
|
|
self._sent_count += 1
|
|
logger.info(
|
|
"Mail queue sent request_id=%s mode=%s auth_user=%s recipient_count=%s attachment_count=%s",
|
|
message.request_id,
|
|
(result.get("selected_mode") or {}).get("label", "<unknown>"),
|
|
result.get("authenticated_as") or "<none>",
|
|
result.get("recipient_count", 0),
|
|
result.get("attachment_count", 0),
|
|
)
|
|
except Exception as exc:
|
|
friendly_error = describe_smtp_error(exc)
|
|
with self._lock:
|
|
self._active_request_id = None
|
|
self._last_result = "failed"
|
|
self._last_error = friendly_error
|
|
self._last_error_at = time.time()
|
|
self._last_activity_at = self._last_error_at
|
|
self._failed_count += 1
|
|
logger.exception(
|
|
"Mail queue delivery failed request_id=%s error=%s",
|
|
getattr(message, "request_id", "unknown"),
|
|
friendly_error,
|
|
)
|
|
finally:
|
|
self._queue.task_done()
|
|
|
|
|
|
_MAIL_QUEUE = MailQueue()
|
|
|
|
|
|
def get_mail_queue() -> MailQueue:
|
|
"""Return the singleton mail queue."""
|
|
return _MAIL_QUEUE
|