"""SMTP2GO email/SMS event webhooks → Message + ProviderEvent updates.""" from __future__ import annotations import json import logging from typing import Any from django.http import HttpRequest from contacts.models import Channel, Contact from messaging.models import Message, ProviderEvent from messaging.services import set_channel_consent logger = logging.getLogger(__name__) PROVIDER_EMAIL = "smtp2go_email" PROVIDER_SMS = "smtp2go_sms" PROVIDER_PCM = "pcm" PROVIDER = PROVIDER_EMAIL # backward-compatible alias # Do not move a message backward to a weaker delivery state. _STATUS_RANK = { Message.Status.DRAFT: 0, Message.Status.SCHEDULED: 1, Message.Status.QUEUED: 2, Message.Status.SENT: 3, Message.Status.FAILED: 3, Message.Status.DELIVERED: 4, Message.Status.BOUNCED: 5, Message.Status.SUPPRESSED: 5, } _MONICA_HEADER_KEYS = ( "X-Monica-Message-Id", "x-monica-message-id", "X_Monica_Message_Id", "monica-message-id", ) def parse_webhook_payload(request: HttpRequest) -> dict[str, Any]: """Accept JSON or form-encoded SMTP2GO webhook bodies.""" content_type = (request.content_type or "").lower() if "application/json" in content_type: try: data = json.loads(request.body.decode() or "{}") except json.JSONDecodeError: return {} return data if isinstance(data, dict) else {} # Form-encoded (SMTP2GO default) return {key: request.POST.get(key) for key in request.POST.keys()} def extract_monica_message_id(payload: dict[str, Any]) -> str: """Pull our correlation id from flat keys or a nested headers object.""" for key in _MONICA_HEADER_KEYS: value = payload.get(key) if value: return str(value).strip() headers = payload.get("headers") or payload.get("email_headers") or {} if isinstance(headers, dict): for key in _MONICA_HEADER_KEYS: value = headers.get(key) if value: return str(value).strip() # Case-insensitive scan lower_map = {str(k).lower(): v for k, v in headers.items()} for key in _MONICA_HEADER_KEYS: value = lower_map.get(key.lower()) if value: return str(value).strip() return "" def find_message_for_email_event(payload: dict[str, Any]) -> Message | None: monica_id = extract_monica_message_id(payload) if monica_id: message = ( Message.objects.select_related("contact", "campaign") .filter(pk=monica_id) .first() ) if message: return message email_id = (payload.get("email_id") or payload.get("email-id") or "").strip() if email_id: message = ( Message.objects.select_related("contact", "campaign") .filter(provider_message_id=email_id) .first() ) if message: return message rcpt = (payload.get("rcpt") or "").strip().lower() if not rcpt: recipients = payload.get("recipients") if isinstance(recipients, str) and recipients.strip(): rcpt = recipients.split(",")[0].strip().lower() elif isinstance(recipients, list) and recipients: rcpt = str(recipients[0]).strip().lower() if not rcpt: return None contact = Contact.objects.filter(email__iexact=rcpt).first() if not contact: return None return ( Message.objects.select_related("contact", "campaign") .filter( contact=contact, channel=Channel.EMAIL, status__in=[ Message.Status.QUEUED, Message.Status.SENT, Message.Status.DELIVERED, Message.Status.FAILED, Message.Status.BOUNCED, ], ) .order_by("-sent_at", "-updated_at") .first() ) def _maybe_upgrade_status(message: Message, new_status: str, *, error: str = "") -> None: current_rank = _STATUS_RANK.get(message.status, 0) new_rank = _STATUS_RANK.get(new_status, 0) # Always allow bounce/suppress to overwrite delivered; allow delivered over sent. if new_rank < current_rank and new_status not in { Message.Status.BOUNCED, Message.Status.SUPPRESSED, Message.Status.FAILED, }: return if ( message.status in {Message.Status.BOUNCED, Message.Status.SUPPRESSED} and new_status == Message.Status.DELIVERED ): return fields = ["status", "updated_at"] message.status = new_status if error: message.error = error[:2000] fields.append("error") elif new_status == Message.Status.DELIVERED: message.error = "" fields.append("error") message.save(update_fields=fields) def _apply_email_event(message: Message, event: str, payload: dict[str, Any]) -> None: event = (event or "").strip().lower() bounce_kind = (payload.get("bounce") or "").strip().lower() err = (payload.get("message") or payload.get("context") or "").strip() email_id = (payload.get("email_id") or payload.get("email-id") or "").strip() if email_id and message.provider_message_id != email_id: message.provider_message_id = email_id message.provider = PROVIDER_EMAIL message.save( update_fields=["provider_message_id", "provider", "updated_at"] ) if event == "processed": if message.status in {Message.Status.QUEUED, Message.Status.DRAFT}: _maybe_upgrade_status(message, Message.Status.SENT) return if event == "delivered": _maybe_upgrade_status(message, Message.Status.DELIVERED) return if event == "bounce": status = Message.Status.BOUNCED _maybe_upgrade_status( message, status, error=err or f"{bounce_kind or 'unknown'} bounce", ) if bounce_kind == "hard": set_channel_consent( message.contact, Channel.EMAIL, opted_in=False, reason="smtp2go_hard_bounce", ) return if event == "reject": _maybe_upgrade_status( message, Message.Status.FAILED, error=err or "rejected by provider" ) return if event == "spam": _maybe_upgrade_status( message, Message.Status.SUPPRESSED, error=err or "spam complaint" ) set_channel_consent( message.contact, Channel.EMAIL, opted_in=False, reason="smtp2go_spam", ) return if event == "unsubscribe": _maybe_upgrade_status( message, Message.Status.SUPPRESSED, error="provider unsubscribe" ) set_channel_consent( message.contact, Channel.EMAIL, opted_in=False, reason="smtp2go_unsubscribe", ) return # open / click / resubscribe — event row only (status unchanged) def process_smtp2go_email_webhook(payload: dict[str, Any]) -> ProviderEvent | None: """ Persist ProviderEvent and update Message delivery status when possible. Returns the stored event (even if message could not be matched). """ event = (payload.get("event") or "").strip().lower() if not event: logger.warning("SMTP2GO webhook missing event: %s", payload) return None message = find_message_for_email_event(payload) if message: _apply_email_event(message, event, payload) message.refresh_from_db() else: logger.info( "SMTP2GO webhook unmatched event=%s rcpt=%s email_id=%s", event, payload.get("rcpt"), payload.get("email_id"), ) return ProviderEvent.objects.create( message=message, provider=PROVIDER_EMAIL, event_type=event, payload=payload, ) def normalize_phone(value: str) -> str: return "".join(ch for ch in (value or "") if ch.isdigit()) def find_message_for_sms_event(payload: dict[str, Any]) -> Message | None: provider_id = ( payload.get("message_id") or payload.get("sms_id") or payload.get("id") or "" ) provider_id = str(provider_id).strip() if provider_id: message = ( Message.objects.select_related("contact", "campaign") .filter(channel=Channel.SMS, provider_message_id=provider_id) .first() ) if message: return message raw_phone = ( payload.get("destination_number") or payload.get("to") or payload.get("phone") or payload.get("from") or "" ) digits = normalize_phone(str(raw_phone)) if len(digits) < 7: return None # Match last 10 digits so +1 / formatting differences still hit. tail = digits[-10:] contacts = Contact.objects.exclude(phone="").only("id", "phone") contact = None for row in contacts.iterator(): if normalize_phone(row.phone).endswith(tail): contact = row break if not contact: return None return ( Message.objects.select_related("contact", "campaign") .filter( contact=contact, channel=Channel.SMS, status__in=[ Message.Status.QUEUED, Message.Status.SENT, Message.Status.DELIVERED, Message.Status.FAILED, ], ) .order_by("-sent_at", "-updated_at") .first() ) def _apply_sms_event(message: Message, event: str, payload: dict[str, Any]) -> None: event = (event or "").strip().lower().replace("-", "_") err = ( payload.get("message") or payload.get("status_code") or payload.get("context") or "" ) err = str(err).strip() provider_id = ( payload.get("message_id") or payload.get("sms_id") or "" ) provider_id = str(provider_id).strip() if provider_id and message.provider_message_id != provider_id: message.provider_message_id = provider_id message.provider = PROVIDER_SMS message.save( update_fields=["provider_message_id", "provider", "updated_at"] ) if event in {"sms_sending", "sending", "sms_submitted", "submitted"}: if message.status in {Message.Status.QUEUED, Message.Status.DRAFT}: _maybe_upgrade_status(message, Message.Status.SENT) return if event in {"sms_delivered", "delivered"}: _maybe_upgrade_status(message, Message.Status.DELIVERED) return if event in {"sms_failed", "failed", "sms_rejected", "rejected"}: _maybe_upgrade_status( message, Message.Status.FAILED, error=err or event, ) return if event in {"sms_opt_out", "opt_out", "optout"}: _maybe_upgrade_status( message, Message.Status.SUPPRESSED, error="sms opt-out" ) set_channel_consent( message.contact, Channel.SMS, opted_in=False, reason="smtp2go_sms_opt_out", ) return def process_smtp2go_sms_webhook(payload: dict[str, Any]) -> ProviderEvent | None: """Persist SMS delivery/opt-out ProviderEvent and update Message when matched.""" event = (payload.get("event") or "").strip().lower() if not event: logger.warning("SMTP2GO SMS webhook missing event: %s", payload) return None message = find_message_for_sms_event(payload) if message: _apply_sms_event(message, event, payload) message.refresh_from_db() else: # Opt-out with no matched campaign message still suppresses by phone. if event.replace("-", "_") in {"sms_opt_out", "opt_out", "optout"}: phone = ( payload.get("destination_number") or payload.get("from") or payload.get("source_number") or "" ) if phone: from messaging.services import record_sms_stop record_sms_stop(str(phone)) logger.info( "SMTP2GO SMS webhook unmatched event=%s phone=%s message_id=%s", event, payload.get("destination_number"), payload.get("message_id"), ) return ProviderEvent.objects.create( message=message, provider=PROVIDER_SMS, event_type=event, payload=payload, ) def is_inbound_sms_stop(payload: dict[str, Any]) -> bool: """True for gateway-style inbound reply payloads (STOP / UNSUBSCRIBE).""" if payload.get("event"): return False text = ( payload.get("text") or payload.get("message") or payload.get("message_content") or "" ) text = str(text).strip().upper() return text in {"STOP", "UNSUBSCRIBE", "CANCEL", "END", "QUIT"} def _pcm_event_type(payload: dict[str, Any]) -> str: for key in ("event", "eventType", "event_type", "type", "status"): value = payload.get(key) if value: return str(value).strip() return "unknown" def find_message_for_pcm_event(payload: dict[str, Any]) -> Message | None: """Correlate PCM webhook to Message via extRefNbr or orderID.""" ext = ( payload.get("extRefNbr") or payload.get("ext_ref_nbr") or payload.get("externalReference") or "" ) if not ext and isinstance(payload.get("recipient"), dict): ext = payload["recipient"].get("extRefNbr") or "" ext = str(ext).strip() if ext: message = ( Message.objects.select_related("contact", "campaign") .filter(pk=ext) .first() ) if message: return message order_id = ( payload.get("orderID") or payload.get("orderId") or payload.get("order_id") or "" ) order_id = str(order_id).strip() if order_id: message = ( Message.objects.select_related("contact", "campaign") .filter(provider_message_id=order_id, channel=Channel.POSTCARD) .first() ) if message: return message return None def _apply_pcm_status(message: Message, status: str, payload: dict[str, Any]) -> None: status_norm = (status or "").strip().lower() err = ( payload.get("message") or payload.get("error") or payload.get("reason") or "" ) err = str(err).strip() order_id = ( payload.get("orderID") or payload.get("orderId") or payload.get("order_id") or "" ) if order_id and message.provider_message_id != str(order_id): message.provider_message_id = str(order_id) message.provider = PROVIDER_PCM message.save( update_fields=["provider_message_id", "provider", "updated_at"] ) if status_norm in {"delivered"}: _maybe_upgrade_status(message, Message.Status.DELIVERED) return if status_norm in {"undeliverable", "returned"}: _maybe_upgrade_status( message, Message.Status.BOUNCED, error=err or "undeliverable", ) return if status_norm in {"canceled", "cancelled"}: _maybe_upgrade_status( message, Message.Status.FAILED, error=err or "canceled" ) return if status_norm in {"pending", "processing", "processed", "mailed", "intransit", "in_transit"}: if message.status in { Message.Status.QUEUED, Message.Status.DRAFT, Message.Status.SCHEDULED, }: _maybe_upgrade_status(message, Message.Status.SENT) return def process_pcm_postcard_webhook(payload: dict[str, Any]) -> ProviderEvent | None: """Record a PCM Integrations postcard event and advance Message status.""" if not payload: return None # Nested data wrappers some webhook UIs use. if "data" in payload and isinstance(payload["data"], dict): inner = dict(payload["data"]) for key in ("event", "eventType", "type"): if key in payload and key not in inner: inner[key] = payload[key] payload = inner event_type = _pcm_event_type(payload) message = find_message_for_pcm_event(payload) if message: status_for_apply = ( payload.get("status") or payload.get("orderStatus") or event_type ) _apply_pcm_status(message, str(status_for_apply), payload) return ProviderEvent.objects.create( message=message, provider=PROVIDER_PCM, event_type=event_type[:64], payload=payload, ) def campaign_engagement_stats(campaign) -> dict[str, int]: """Aggregate delivery + open/click counts for the campaign report.""" messages_qs = campaign.messages.all() statuses = list(messages_qs.values_list("status", flat=True)) message_ids = list(messages_qs.values_list("pk", flat=True)) events = ProviderEvent.objects.filter(message_id__in=message_ids) open_message_ids = set( events.filter(event_type__iexact="open").values_list("message_id", flat=True) ) click_message_ids = set( events.filter(event_type__iexact="click").values_list("message_id", flat=True) ) return { "total": len(statuses), "sent": sum(1 for s in statuses if s in {"sent", "delivered"}), "delivered": sum(1 for s in statuses if s == "delivered"), "failed": sum(1 for s in statuses if s in {"failed", "bounced"}), "bounced": sum(1 for s in statuses if s == "bounced"), "suppressed": sum(1 for s in statuses if s == "suppressed"), "opens": len(open_message_ids), "clicks": len(click_message_ids), "open_events": events.filter(event_type__iexact="open").count(), "click_events": events.filter(event_type__iexact="click").count(), }