from __future__ import annotations from contextlib import contextmanager from contextvars import ContextVar from dataclasses import dataclass from email.message import EmailMessage from email.utils import formatdate, make_msgid from threading import get_ident from typing import Any, Iterator from sqlalchemy.orm import Session from govoplan_core.core.mail import NotificationMailDeliveryRequest from govoplan_core.core.modules import ModuleContext from govoplan_mail.backend.config import ImapConfig from govoplan_mail.backend.mail_profiles import ( MailProfileError, _assert_campaign_inherits_profile_credentials, assert_campaign_mail_policy_allows_json, assert_mail_policy_allows_send, campaign_mail_owner_context, campaign_profile_transport_revisions, effective_mail_profile_policy, ensure_mail_profile_allowed_for_campaign, get_mail_server_profile, imap_config_from_profile, mail_profile_id_from_campaign_json, smtp_config_from_profile, ) from govoplan_mail.backend.server_hierarchy import ( MailHierarchyContext, MailServerHierarchyError, resolve_mail_transport, select_mail_transport, ) from govoplan_mail.backend.runtime import configure_runtime from govoplan_mail.backend.recovery import ( MailRecoveryError, begin_provider_effect_recovery, ) from govoplan_mail.backend.sending.imap import ( ImapAppendError, ImapBatchSession, ImapConfigurationError, append_message_to_sent, ) from govoplan_mail.backend.sending.rate_limit import wait_for_rate_limit from govoplan_mail.backend.sending.smtp import ( SmtpBatchSession, SmtpConfigurationError, SmtpSendError, send_email_bytes, ) _ACTIVE_SMTP_BATCH: ContextVar[SmtpBatchSession | None] = ContextVar( "govoplan_mail_active_smtp_batch", default=None, ) @dataclass(frozen=True, slots=True) class CampaignSmtpDeliveryResult: envelope_recipients: list[str] refused_recipients: dict[str, dict[str, int | str]] connection_sequence: int = 1 session_reused: bool = False reconnect_count: int = 0 @property def accepted_count(self) -> int: return len(self.envelope_recipients) - len(self.refused_recipients) @dataclass(frozen=True, slots=True) class CampaignImapAppendResult: folder: str connection_sequence: int = 1 session_reused: bool = False reconnect_count: int = 0 class CampaignImapBatchState: """Lazy transport reuse; authorization and recovery remain per message.""" def __init__(self, *, tenant_id: str, campaign_id: str): self.tenant_id = tenant_id self.campaign_id = campaign_id self._session: ImapBatchSession | None = None self._binding: tuple[Any, ...] | None = None self._previous_connections = 0 self._previous_reconnects = 0 self._closed = False self._owner_thread = get_ident() @property def connection_count(self) -> int: return self._previous_connections + (self._session.connection_count if self._session else 0) @property def reconnect_count(self) -> int: return self._previous_reconnects + (self._session.reconnect_count if self._session else 0) def assert_scope(self, *, tenant_id: str, campaign_id: str) -> None: if ( self._closed or get_ident() != self._owner_thread or tenant_id != self.tenant_id or campaign_id != self.campaign_id ): raise ImapConfigurationError("The IMAP batch does not match this campaign scope") def session_for(self, config: ImapConfig, *, binding: tuple[Any, ...]) -> ImapBatchSession: if self._closed: raise ImapConfigurationError("The IMAP batch is closed") if self._session is not None and ( self._binding != binding or not self._session.matches_config(config) ): self._release_session() if self._session is None: self._session = ImapBatchSession(config) self._binding = binding return self._session def _release_session(self) -> None: if self._session is not None: self._previous_connections += self._session.connection_count self._previous_reconnects += self._session.reconnect_count self._session.close() self._session = None def close(self) -> None: self._closed = True self._release_session() _ACTIVE_IMAP_BATCH: ContextVar[CampaignImapBatchState | None] = ContextVar( "govoplan_mail_active_imap_batch", default=None, ) @contextmanager def campaign_imap_batch(*, tenant_id: str, campaign_id: str) -> Iterator[CampaignImapBatchState]: """Open no connection until an individual append passes its current checks.""" state = CampaignImapBatchState(tenant_id=tenant_id, campaign_id=campaign_id) token = _ACTIVE_IMAP_BATCH.set(state) try: yield state finally: _ACTIVE_IMAP_BATCH.reset(token) state.close() @dataclass(frozen=True, slots=True) class CampaignSmtpBatchState: session: SmtpBatchSession @property def status(self) -> str: return "ready" @property def connection_count(self) -> int: return self.session.connection_count @property def reconnect_count(self) -> int: return self.session.reconnect_count def _sanitized_refusals( refused_recipients: dict[str, tuple[int, bytes | str]], ) -> dict[str, dict[str, int | str]]: sanitized: dict[str, dict[str, int | str]] = {} for recipient, (raw_code, _provider_message) in refused_recipients.items(): try: code = int(raw_code) except (TypeError, ValueError): code = 0 if 400 <= code < 500: classification = "temporary" message = "Temporary recipient rejection" elif 500 <= code < 600: classification = "permanent" message = "Permanent recipient rejection" else: classification = "unknown" message = "Recipient rejected" sanitized[str(recipient)] = { "status_code": code, "classification": classification, "message": message, } return sanitized def _sanitized_smtp_error(exc: SmtpSendError) -> SmtpSendError: if exc.outcome_unknown: message = "Mail delivery outcome is unknown after transmission started." elif exc.temporary: message = "Mail delivery failed temporarily." else: message = "Mail delivery was rejected." return SmtpSendError( message, temporary=exc.temporary, outcome_unknown=exc.outcome_unknown, systemic=exc.systemic, reason_code=exc.reason_code, phase=exc.phase, ) def _sanitized_imap_error(exc: ImapAppendError) -> ImapAppendError: if exc.outcome_unknown: message = "The Sent-folder append outcome is unknown; inspect the mailbox before retrying." elif exc.temporary: message = "Appending the sent message failed temporarily." else: message = "Appending the sent message was rejected." return ImapAppendError( message, temporary=exc.temporary, outcome_unknown=exc.outcome_unknown, ) def _authorized_campaign_profile( session: Session, *, tenant_id: str, campaign_id: str, profile_id: str, selection: dict[str, str | None] | None = None, credential_protocol: str | None = None, ): profile = ensure_mail_profile_allowed_for_campaign( session, tenant_id=tenant_id, campaign_id=campaign_id, profile_id=profile_id, require_active=True, ) policy = effective_mail_profile_policy(session, tenant_id=tenant_id, campaign_id=campaign_id) _assert_campaign_inherits_profile_credentials(profile, policy, selection, protocol=credential_protocol) return profile def _campaign_hierarchy_context( session: Session, *, tenant_id: str, campaign_id: str, ) -> MailHierarchyContext: campaign = campaign_mail_owner_context( session, tenant_id=tenant_id, campaign_id=campaign_id, ) return MailHierarchyContext( tenant_id=tenant_id, user_id=campaign.owner_user_id, group_ids=( frozenset({campaign.owner_group_id}) if campaign.owner_group_id else frozenset() ), target_scope_type="campaign", target_scope_id=campaign.id, ) def _selection_payload( *, profile_id: str, smtp_server_id: str | None = None, smtp_credential_id: str | None = None, imap_server_id: str | None = None, imap_credential_id: str | None = None, ) -> dict[str, str | None]: return { "mail_profile_id": profile_id, "smtp_server_id": smtp_server_id, "smtp_credential_id": smtp_credential_id, "imap_server_id": imap_server_id, "imap_credential_id": imap_credential_id, } def _supports_hierarchy(session: object) -> bool: return callable(getattr(session, "execute", None)) def campaign_profile_delivery_summary( session: Session, *, tenant_id: str, profile_id: str, campaign_id: str | None = None, owner_user_id: str | None = None, owner_group_id: str | None = None, smtp_server_id: str | None = None, smtp_credential_id: str | None = None, imap_server_id: str | None = None, imap_credential_id: str | None = None, ) -> dict[str, Any]: """Return only non-secret capabilities and opaque drift evidence.""" selection = _selection_payload( profile_id=profile_id, smtp_server_id=smtp_server_id, smtp_credential_id=smtp_credential_id, imap_server_id=imap_server_id, imap_credential_id=imap_credential_id, ) if campaign_id: profile = _authorized_campaign_profile( session, tenant_id=tenant_id, campaign_id=campaign_id, profile_id=profile_id, selection=selection, ) else: assert_campaign_mail_policy_allows_json( session, tenant_id=tenant_id, raw_json={"server": {key: value for key, value in selection.items() if value}}, owner_user_id=owner_user_id, owner_group_id=owner_group_id, ) profile = get_mail_server_profile( session, tenant_id=tenant_id, profile_id=profile_id, require_active=True, ) if campaign_id and _supports_hierarchy(session): context = _campaign_hierarchy_context( session, tenant_id=tenant_id, campaign_id=campaign_id, ) try: smtp = select_mail_transport( session, profile=profile, protocol="smtp", context=context, server_id=smtp_server_id, credential_id=smtp_credential_id, ) imap = select_mail_transport( session, profile=profile, protocol="imap", context=context, server_id=imap_server_id, credential_id=imap_credential_id, ) except MailServerHierarchyError as exc: raise MailProfileError(str(exc)) from exc smtp_available = smtp.available imap_available = imap.available smtp_revision = smtp.transport_revision imap_revision = imap.transport_revision if imap.available else None resolved_smtp_server_id = smtp.server.id if smtp.server else None resolved_smtp_credential_id = smtp.credential.id if smtp.credential else None resolved_imap_server_id = imap.server.id if imap.server else None resolved_imap_credential_id = imap.credential.id if imap.credential else None else: revisions = campaign_profile_transport_revisions(profile) smtp_config = profile.smtp_config or {} imap_config = profile.imap_config or {} smtp_available = bool(smtp_config.get("host") and smtp_config.get("port")) imap_available = bool(imap_config.get("host") and imap_config.get("port")) smtp_revision = revisions["smtp"] imap_revision = revisions["imap"] resolved_smtp_server_id = smtp_server_id resolved_smtp_credential_id = smtp_credential_id resolved_imap_server_id = imap_server_id resolved_imap_credential_id = imap_credential_id return { "mail_profile_id": profile_id, "smtp_server_id": resolved_smtp_server_id, "smtp_credential_id": resolved_smtp_credential_id, "imap_server_id": resolved_imap_server_id, "imap_credential_id": resolved_imap_credential_id, "smtp_available": smtp_available, "imap_available": imap_available, "smtp_transport_revision": smtp_revision, "imap_transport_revision": imap_revision, } @contextmanager def campaign_smtp_batch( session: Session, *, tenant_id: str, campaign_id: str, profile_id: str, envelope_from: str, envelope_recipients: list[str], from_header: str | None, expected_smtp_transport_revision: str, smtp_server_id: str | None = None, smtp_credential_id: str | None = None, ) -> Iterator[CampaignSmtpBatchState]: """Preflight and retain one authorized SMTP connection for a batch.""" selection = _selection_payload( profile_id=profile_id, smtp_server_id=smtp_server_id, smtp_credential_id=smtp_credential_id, ) try: profile = _authorized_campaign_profile( session, tenant_id=tenant_id, campaign_id=campaign_id, profile_id=profile_id, selection=selection, credential_protocol="smtp", ) except MailProfileError: raise except Exception: raise SmtpConfigurationError("The selected Mail profile's SMTP configuration is unusable.") from None if _supports_hierarchy(session): context = _campaign_hierarchy_context( session, tenant_id=tenant_id, campaign_id=campaign_id, ) try: selected_smtp = select_mail_transport( session, profile=profile, protocol="smtp", context=context, server_id=smtp_server_id, credential_id=smtp_credential_id, ) except MailServerHierarchyError as exc: raise MailProfileError(str(exc)) from exc current_revision = selected_smtp.transport_revision else: context = None current_revision = campaign_profile_transport_revisions(profile)["smtp"] if current_revision != expected_smtp_transport_revision: raise MailProfileError( "The selected Mail profile's SMTP settings changed after this campaign was built. " "Revalidate and rebuild the campaign before delivery." ) try: smtp = ( resolve_mail_transport( session, profile=profile, protocol="smtp", context=context, server_id=smtp_server_id, credential_id=smtp_credential_id, ).config if context is not None else smtp_config_from_profile(profile) ) except MailProfileError: raise except Exception: raise SmtpConfigurationError("The selected Mail profile's SMTP configuration is unusable.") from None try: assert_mail_policy_allows_send( session, tenant_id=tenant_id, campaign_id=campaign_id, smtp=smtp, imap=None, envelope_sender=envelope_from, from_header=from_header, recipients=envelope_recipients, ) except MailProfileError: raise MailProfileError("Mail delivery is blocked by the effective Mail policy.") from None smtp_session = SmtpBatchSession(smtp) smtp_session.preflight() token = _ACTIVE_SMTP_BATCH.set(smtp_session) try: yield CampaignSmtpBatchState(session=smtp_session) finally: _ACTIVE_SMTP_BATCH.reset(token) smtp_session.close() def send_campaign_email_bytes( session: Session, *, tenant_id: str, campaign_id: str, profile_id: str, message_bytes: bytes, envelope_from: str, envelope_recipients: list[str], from_header: str | None, expected_smtp_transport_revision: str, smtp_server_id: str | None = None, smtp_credential_id: str | None = None, recovery_effect_id: str | None = None, recovery_resource_type: str | None = None, recovery_resource_id: str | None = None, ) -> CampaignSmtpDeliveryResult: selection = _selection_payload( profile_id=profile_id, smtp_server_id=smtp_server_id, smtp_credential_id=smtp_credential_id, ) try: profile = _authorized_campaign_profile( session, tenant_id=tenant_id, campaign_id=campaign_id, profile_id=profile_id, selection=selection, credential_protocol="smtp", ) except MailProfileError: raise except Exception: raise SmtpConfigurationError("The selected Mail profile's SMTP configuration is unusable.") from None resolved_smtp = None if _supports_hierarchy(session): context = _campaign_hierarchy_context( session, tenant_id=tenant_id, campaign_id=campaign_id, ) try: selected_smtp = select_mail_transport( session, profile=profile, protocol="smtp", context=context, server_id=smtp_server_id, credential_id=smtp_credential_id, ) except MailServerHierarchyError as exc: raise MailProfileError(str(exc)) from exc current_smtp_revision = selected_smtp.transport_revision else: current_smtp_revision = campaign_profile_transport_revisions(profile)["smtp"] if current_smtp_revision != expected_smtp_transport_revision: raise MailProfileError( "The selected Mail profile's SMTP settings changed after this campaign was built. " "Revalidate and rebuild the campaign before delivery." ) try: if _supports_hierarchy(session): resolved_smtp = resolve_mail_transport( session, profile=profile, protocol="smtp", context=context, server_id=smtp_server_id, credential_id=smtp_credential_id, ) smtp = resolved_smtp.config else: smtp = smtp_config_from_profile(profile) except MailProfileError: raise except Exception: raise SmtpConfigurationError("The selected Mail profile's SMTP configuration is unusable.") from None try: assert_mail_policy_allows_send( session, tenant_id=tenant_id, campaign_id=campaign_id, smtp=smtp, imap=None, envelope_sender=envelope_from, from_header=from_header, recipients=envelope_recipients, ) except MailProfileError: raise MailProfileError("Mail delivery is blocked by the effective Mail policy.") from None try: recovery = begin_provider_effect_recovery( kind="smtp-delivery", effect_id=recovery_effect_id, tenant_id=tenant_id, profile_id=profile_id, message_bytes=message_bytes, expected_transport_revision=expected_smtp_transport_revision, recipient_count=len(envelope_recipients), resource_type=recovery_resource_type, resource_id=recovery_resource_id, ) except MailRecoveryError as exc: raise SmtpConfigurationError(str(exc)) from None if recovery is not None and recovery.replayed: raise SmtpSendError( "The matching SMTP effect already succeeded; reconcile caller state without resending.", outcome_unknown=True, ) try: result = send_email_bytes( message_bytes, smtp_config=smtp, envelope_from=envelope_from, envelope_recipients=envelope_recipients, batch_session=_ACTIVE_SMTP_BATCH.get(), ) except SmtpSendError as exc: sanitized = _sanitized_smtp_error(exc) if recovery is not None: if sanitized.outcome_unknown: recovery.unknown(code="smtp_outcome_unknown", summary=str(sanitized)) else: recovery.reject(code="smtp_rejected", summary=str(sanitized)) raise sanitized from None except SmtpConfigurationError as exc: if recovery is not None: recovery.reject(code="smtp_configuration", summary=str(exc)) raise SmtpConfigurationError("The selected Mail profile's SMTP configuration is unusable.") from None except Exception: if recovery is not None: recovery.unknown( code="unexpected_provider_error", summary="SMTP outcome is unknown after an unexpected provider failure", ) raise SmtpSendError( "Mail delivery outcome is unknown after the provider effect started.", outcome_unknown=True, ) from None sanitized_result = CampaignSmtpDeliveryResult( envelope_recipients=list(result.envelope_recipients), refused_recipients=_sanitized_refusals(result.refused_recipients), connection_sequence=getattr(result, "connection_sequence", 0), session_reused=getattr(result, "session_reused", False), reconnect_count=getattr(result, "reconnect_count", 0), ) if recovery is not None: try: recovery.succeed_smtp( accepted_count=sanitized_result.accepted_count, refused_recipients=sanitized_result.refused_recipients, ) except Exception: raise SmtpSendError( "SMTP returned an outcome, but durable recovery evidence could not be finalized.", outcome_unknown=True, ) from None return sanitized_result def append_campaign_message_to_sent( session: Session, *, tenant_id: str, campaign_id: str, profile_id: str, message_bytes: bytes, folder: str | None, expected_smtp_transport_revision: str, expected_imap_transport_revision: str | None, smtp_server_id: str | None = None, smtp_credential_id: str | None = None, imap_server_id: str | None = None, imap_credential_id: str | None = None, recovery_effect_id: str | None = None, recovery_resource_type: str | None = None, recovery_resource_id: str | None = None, ) -> CampaignImapAppendResult: batch = _ACTIVE_IMAP_BATCH.get() if batch is not None: batch.assert_scope(tenant_id=tenant_id, campaign_id=campaign_id) selection = _selection_payload( profile_id=profile_id, smtp_server_id=smtp_server_id, smtp_credential_id=smtp_credential_id, imap_server_id=imap_server_id, imap_credential_id=imap_credential_id, ) try: profile = _authorized_campaign_profile( session, tenant_id=tenant_id, campaign_id=campaign_id, profile_id=profile_id, selection=selection, credential_protocol="imap", ) except MailProfileError: raise except Exception: raise ImapConfigurationError("The selected Mail profile's IMAP configuration is unusable.") from None if _supports_hierarchy(session): context = _campaign_hierarchy_context( session, tenant_id=tenant_id, campaign_id=campaign_id, ) try: selected_smtp = select_mail_transport( session, profile=profile, protocol="smtp", context=context, server_id=smtp_server_id, credential_id=smtp_credential_id, ) selected_imap = select_mail_transport( session, profile=profile, protocol="imap", context=context, server_id=imap_server_id, credential_id=imap_credential_id, ) except MailServerHierarchyError as exc: raise MailProfileError(str(exc)) from exc smtp_revision = selected_smtp.transport_revision imap_revision = selected_imap.transport_revision else: revisions = campaign_profile_transport_revisions(profile) smtp_revision = revisions["smtp"] imap_revision = revisions["imap"] if smtp_revision != expected_smtp_transport_revision: raise MailProfileError( "The selected Mail profile's SMTP settings changed after this campaign was built. " "Revalidate and rebuild the campaign before append-to-Sent delivery." ) if imap_revision != expected_imap_transport_revision: raise MailProfileError( "The selected Mail profile's IMAP settings changed after this campaign was built. " "Revalidate and rebuild the campaign before append-to-Sent delivery." ) try: if _supports_hierarchy(session): imap = resolve_mail_transport( session, profile=profile, protocol="imap", context=context, server_id=imap_server_id, credential_id=imap_credential_id, ).config else: imap = imap_config_from_profile(profile) except MailProfileError: raise except Exception: raise ImapConfigurationError("The selected Mail profile's IMAP configuration is unusable.") from None if imap is None: raise ImapConfigurationError("The selected Mail profile has no IMAP configuration") try: assert_mail_policy_allows_send( session, tenant_id=tenant_id, campaign_id=campaign_id, smtp=None, imap=imap, ) except MailProfileError: raise MailProfileError("Appending to Sent is blocked by the effective Mail policy.") from None batch_session = None if batch is not None: batch_session = batch.session_for( imap, binding=( tenant_id, campaign_id, profile_id, smtp_server_id, smtp_credential_id, imap_server_id, imap_credential_id, smtp_revision, imap_revision, folder, ), ) try: recovery = begin_provider_effect_recovery( kind="imap-append", effect_id=recovery_effect_id, tenant_id=tenant_id, profile_id=profile_id, message_bytes=message_bytes, expected_transport_revision=expected_imap_transport_revision, folder=folder, resource_type=recovery_resource_type, resource_id=recovery_resource_id, ) except MailRecoveryError as exc: raise ImapConfigurationError(str(exc)) from None if recovery is not None and recovery.replayed: raise ImapAppendError( "The matching IMAP append already succeeded; reconcile caller state without appending again.", outcome_unknown=True, ) try: if batch_session is None: result = append_message_to_sent(message_bytes, imap_config=imap, folder=folder) else: result = append_message_to_sent( message_bytes, imap_config=imap, folder=folder, batch_session=batch_session, ) except ImapAppendError as exc: sanitized = _sanitized_imap_error(exc) if recovery is not None: if sanitized.outcome_unknown: recovery.unknown(code="imap_outcome_unknown", summary=str(sanitized)) else: recovery.reject(code="imap_rejected", summary=str(sanitized)) raise sanitized from None except ImapConfigurationError as exc: if recovery is not None: recovery.reject(code="imap_configuration", summary=str(exc)) raise ImapConfigurationError("The selected Mail profile's IMAP configuration is unusable.") from None except Exception: if recovery is not None: recovery.unknown( code="unexpected_provider_error", summary="IMAP APPEND outcome is unknown after an unexpected provider failure", ) raise ImapAppendError( "The Sent-folder append outcome is unknown; inspect the mailbox before retrying.", outcome_unknown=True, ) from None if recovery is not None: try: recovery.succeed_imap(folder=result.folder) except Exception: if batch is not None: batch.close() raise ImapAppendError( "IMAP APPEND returned success, but durable recovery evidence could not be finalized.", outcome_unknown=True, ) from None return CampaignImapAppendResult( folder=result.folder, connection_sequence=batch.connection_count if batch else getattr(result, "connection_sequence", 1), session_reused=getattr(result, "session_reused", False), reconnect_count=batch.reconnect_count if batch else getattr(result, "reconnect_count", 0), ) class MailCampaignCapability: MailProfileError = MailProfileError SmtpConfigurationError = SmtpConfigurationError SmtpSendError = SmtpSendError ImapConfigurationError = ImapConfigurationError ImapAppendError = ImapAppendError assert_campaign_mail_policy_allows_json = staticmethod(assert_campaign_mail_policy_allows_json) mail_profile_id_from_campaign_json = staticmethod(mail_profile_id_from_campaign_json) campaign_profile_delivery_summary = staticmethod(campaign_profile_delivery_summary) campaign_smtp_batch = staticmethod(campaign_smtp_batch) campaign_imap_batch = staticmethod(campaign_imap_batch) send_campaign_email_bytes = staticmethod(send_campaign_email_bytes) append_campaign_message_to_sent = staticmethod(append_campaign_message_to_sent) wait_for_rate_limit = staticmethod(wait_for_rate_limit) @staticmethod def submit_delivery_command(session: Session, **kwargs: Any) -> dict[str, object]: from govoplan_mail.backend.delivery_outbox import submit_delivery_command return submit_delivery_command(session, **kwargs) @staticmethod def delivery_command_summary( session: Session, *, tenant_id: str, command_id: str, ) -> dict[str, object]: from govoplan_mail.backend.delivery_outbox import ( delivery_command_summary, get_delivery_command, ) return delivery_command_summary( get_delivery_command( session, tenant_id=tenant_id, command_id=command_id, ) ) @staticmethod def mock_mailbox(): from govoplan_mail.backend.dev import mock_mailbox return mock_mailbox def campaign_capability(context: ModuleContext) -> MailCampaignCapability: configure_runtime(settings=context.settings) return MailCampaignCapability() def delivery_outbox_capability(context: ModuleContext): from govoplan_mail.backend.delivery_outbox import MailDeliveryOutboxCapability configure_runtime(settings=context.settings) return MailDeliveryOutboxCapability() class MailNotificationDeliveryCapability: """Submit notification email to the Mail-owned durable outbox.""" def submit_notification_mail( self, session: Session, request: NotificationMailDeliveryRequest, ) -> dict[str, object]: profile_id = str(request.mail_profile_id or "").strip() from_address = str(request.from_address or "").strip() if not profile_id or not from_address: return { "status": "paused", "provider": "mail.notificationDelivery", "error": ( "Notification email requires a Mail profile and sender " "selected by tenant policy." ), } try: transport = campaign_profile_delivery_summary( session, tenant_id=request.tenant_id, profile_id=profile_id, smtp_server_id=request.smtp_server_id, smtp_credential_id=request.smtp_credential_id, ) except MailProfileError: return { "status": "paused", "provider": "mail.notificationDelivery", "error": ( "The selected notification Mail profile is unavailable or " "blocked by effective policy." ), } if not transport.get("smtp_available"): return { "status": "paused", "provider": "mail.notificationDelivery", "error": "The selected notification Mail profile has no usable SMTP transport.", } message = EmailMessage() message["Date"] = formatdate(localtime=True) message["Message-ID"] = make_msgid() message["Subject"] = request.subject message["From"] = from_address message["To"] = request.recipient message.set_content(request.body_text) if request.body_html: message.add_alternative(request.body_html, subtype="html") from govoplan_mail.backend.delivery_outbox import submit_delivery_command command = submit_delivery_command( session, tenant_id=request.tenant_id, command_type="notification", source_module="notifications", source_resource_type="notification", source_resource_id=request.notification_id, source_version_id=None, idempotency_key=f"notification:{request.notification_id}", profile_id=profile_id, smtp_server_id=( str(transport.get("smtp_server_id")) if transport.get("smtp_server_id") else request.smtp_server_id ), smtp_credential_id=( str(transport.get("smtp_credential_id")) if transport.get("smtp_credential_id") else request.smtp_credential_id ), message_bytes=bytes(message), envelope_from=from_address, envelope_recipients=[request.recipient], from_header=from_address, expected_smtp_transport_revision=str( transport["smtp_transport_revision"] ), ) return { "status": "accepted", "provider": "mail.delivery_outbox", "external_message_id": command["id"], "delivery_status": command["status"], "duplicate": bool(command.get("duplicate")), } def notification_delivery_capability( context: ModuleContext, ) -> MailNotificationDeliveryCapability: configure_runtime(settings=context.settings) return MailNotificationDeliveryCapability()