Add Django site, Docker packaging, and beta/prod Gitea deploys.
Unignore site/ (was blocked by mkdocs /site rule), add compose/Docker/uv tooling, and split deploys so push to main goes to beta while prod stays manual.
This commit is contained in:
@@ -0,0 +1,575 @@
|
||||
"""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(),
|
||||
}
|
||||
Reference in New Issue
Block a user