Files
govoplan-campaign/src/govoplan_campaign/backend/sending/jobs.py
T

4229 lines
139 KiB
Python

from __future__ import annotations
import hashlib
import json
from collections import Counter
from dataclasses import asdict, dataclass
from datetime import datetime, timezone
from email import policy
from email.parser import BytesParser
from pathlib import Path
from typing import Any
from uuid import uuid4
from sqlalchemy import func
from sqlalchemy.exc import IntegrityError
from sqlalchemy.orm import Session
from govoplan_core.core.notifications import (
NotificationDispatchRequest,
notification_dispatch_provider,
)
from govoplan_core.core.object_storage import (
StorageBackendError,
configured_storage_backend,
)
from govoplan_core.audit.logging import audit_event
from govoplan_core.security.redaction import redact_secret_values
from govoplan_core.settings import settings as core_settings
from govoplan_campaign.backend.db.models import (
Campaign,
CampaignJob,
CampaignMessageAction,
CampaignMessageActionAttempt,
CampaignStatus,
CampaignVersion,
CampaignVersionWorkflowState,
JobBuildStatus,
JobImapStatus,
JobPostboxStatus,
JobQueueStatus,
JobSendStatus,
JobValidationStatus,
ImapAppendAttempt,
PostboxDeliveryAttempt,
SendAttempt,
)
from govoplan_campaign.backend.campaign.models import DeliveryChannelPolicy
from govoplan_campaign.backend.approval_gate import (
CampaignApprovalGateError,
assert_campaign_approval,
)
from govoplan_campaign.backend.delivery_policy import (
CampaignDeliveryPolicyError,
SynchronousSendPolicy,
effective_synchronous_send_policy,
)
from govoplan_campaign.backend.sending.execution import (
ExecutionSnapshot,
ExecutionSnapshotError,
ensure_execution_snapshot,
profile_delivery_summary,
)
from govoplan_campaign.backend.sending.postbox_delivery import (
PostboxChannelOutcome,
deliver_campaign_job_to_postboxes,
)
from govoplan_campaign.backend.runtime import get_registry, get_settings
from govoplan_campaign.backend.integrations import (
ImapAppendError,
ImapConfigurationError,
MailProfileError,
SmtpConfigurationError,
SmtpSendError,
files_integration,
mail_integration,
)
class QueueingError(RuntimeError):
pass
class SendJobError(RuntimeError):
pass
class SynchronousSendRejected(QueueingError):
def __init__(
self,
message: str,
*,
reason: str,
eligible_count: int | None = None,
policy: SynchronousSendPolicy | None = None,
) -> None:
super().__init__(message)
self.reason = reason
self.eligible_count = eligible_count
self.policy = policy
def audit_details(self) -> dict[str, Any]:
return {
"delivery_mode": "synchronous",
"rejection_reason": self.reason,
"eligible_recipient_job_count": self.eligible_count,
"synchronous_send_policy": self.policy.as_dict()
if self.policy is not None
else None,
}
@dataclass(frozen=True, slots=True)
class QueueCampaignResult:
campaign_id: str
version_id: str
queued_count: int
skipped_count: int
blocked_count: int
enqueued_count: int
delivery_mode: str = "worker_queue"
worker_queue_available: bool = False
dry_run: bool = False
def as_dict(self) -> dict[str, Any]:
return {
"campaign_id": self.campaign_id,
"version_id": self.version_id,
"queued_count": self.queued_count,
"skipped_count": self.skipped_count,
"blocked_count": self.blocked_count,
"enqueued_count": self.enqueued_count,
"delivery_mode": self.delivery_mode,
"worker_queue_available": self.worker_queue_available,
"dry_run": self.dry_run,
}
@dataclass(frozen=True, slots=True)
class SendCampaignNowResult:
campaign_id: str
version_id: str
attempted_count: int
sent_count: int
failed_count: int
outcome_unknown_count: int
skipped_count: int
preflight_count: int = 0
synchronous_send_policy: dict[str, Any] | None = None
dry_run: bool = False
results: list[dict[str, Any]] | None = None
def as_dict(self) -> dict[str, Any]:
return {
"campaign_id": self.campaign_id,
"version_id": self.version_id,
"attempted_count": self.attempted_count,
"sent_count": self.sent_count,
"failed_count": self.failed_count,
"outcome_unknown_count": self.outcome_unknown_count,
"skipped_count": self.skipped_count,
"preflight_count": self.preflight_count,
"delivery_mode": "synchronous",
"synchronous_send_policy": self.synchronous_send_policy or {},
"dry_run": self.dry_run,
"results": self.results or [],
}
@dataclass(frozen=True, slots=True)
class SendJobResult:
job_id: str
status: str
attempt_number: int
dry_run: bool = False
message: str | None = None
def as_dict(self) -> dict[str, Any]:
return {
"job_id": self.job_id,
"status": self.status,
"attempt_number": self.attempt_number,
"dry_run": self.dry_run,
"message": self.message,
}
@dataclass(frozen=True, slots=True)
class AppendSentResult:
job_id: str
status: str
attempt_number: int
dry_run: bool = False
folder: str | None = None
message: str | None = None
def as_dict(self) -> dict[str, Any]:
return {
"job_id": self.job_id,
"status": self.status,
"attempt_number": self.attempt_number,
"dry_run": self.dry_run,
"folder": self.folder,
"message": self.message,
}
@dataclass(frozen=True, slots=True)
class _CampaignDeliveryCounts:
accepted: int
partial: int
unknown: int
failed: int
active: int
cancelled: int
not_started: int
@dataclass(frozen=True, slots=True)
class _SendJobDeliveryContext:
version: CampaignVersion
snapshot: ExecutionSnapshot
message_bytes: bytes
envelope_from: str | None
envelope_recipients: list[str]
@dataclass(frozen=True, slots=True)
class _MailChannelOutcome:
accepted: bool = False
rejected_temporary: bool = False
rejected_permanent: bool = False
outcome_unknown: bool = False
message: str | None = None
@property
def rejected_before_acceptance(self) -> bool:
return (
not self.accepted
and not self.outcome_unknown
and (self.rejected_temporary or self.rejected_permanent)
)
@dataclass(frozen=True, slots=True)
class _DeliveryOutcomeSummary:
mail_accepted: bool
postbox_accepted: int
postbox_rejected: int
outcome_unknown: bool
temporary_rejection: bool
@property
def accepted_count(self) -> int:
return int(self.mail_accepted) + self.postbox_accepted
@dataclass(frozen=True, slots=True)
class _ImapAppendContext:
snapshot: ExecutionSnapshot
message_bytes: bytes
folder: str
@dataclass(frozen=True, slots=True)
class _ClaimedImapAppend:
job: CampaignJob
attempt: ImapAppendAttempt
claim_token: str
QUEUEABLE_VALIDATION_STATUSES = {
JobValidationStatus.READY.value,
JobValidationStatus.WARNING.value,
}
SMTP_ACCEPTED_STATUSES = {JobSendStatus.SMTP_ACCEPTED.value, JobSendStatus.SENT.value}
DELIVERY_ACCEPTED_STATUSES = SMTP_ACCEPTED_STATUSES | {
JobSendStatus.POSTBOX_ACCEPTED.value,
JobSendStatus.DELIVERED.value,
JobSendStatus.PARTIALLY_ACCEPTED.value,
}
FULLY_ACCEPTED_STATUSES = SMTP_ACCEPTED_STATUSES | {
JobSendStatus.POSTBOX_ACCEPTED.value,
JobSendStatus.DELIVERED.value,
}
DELIVERY_MODE_SYNCHRONOUS = "synchronous"
DELIVERY_MODE_WORKER_QUEUE = "worker_queue"
DELIVERY_MODE_DATABASE_QUEUE = "database_queue"
DELIVERY_MODES = {
DELIVERY_MODE_SYNCHRONOUS,
DELIVERY_MODE_WORKER_QUEUE,
DELIVERY_MODE_DATABASE_QUEUE,
}
AUTOMATICALLY_SENDABLE_STATUSES = {JobSendStatus.QUEUED.value}
EXPLICIT_RETRY_STATUSES = {
JobSendStatus.FAILED_TEMPORARY.value,
JobSendStatus.FAILED_PERMANENT.value,
}
INITIAL_QUEUE_SKIPPED_SEND_STATUSES = DELIVERY_ACCEPTED_STATUSES | {
JobSendStatus.SKIPPED.value,
JobSendStatus.CLAIMED.value,
JobSendStatus.SENDING.value,
JobSendStatus.OUTCOME_UNKNOWN.value,
JobSendStatus.FAILED_TEMPORARY.value,
JobSendStatus.FAILED_PERMANENT.value,
JobSendStatus.CANCELLED.value,
}
INITIAL_QUEUE_SKIPPED_QUEUE_STATUSES = {
JobQueueStatus.CANCELLED.value,
JobQueueStatus.SENDING.value,
JobQueueStatus.PAUSED.value,
}
CAMPAIGN_STATUS_NOTIFICATION_EVENTS = {
CampaignStatus.QUEUED.value: ("campaign.queued", "Campaign queued", 1),
CampaignStatus.SENDING.value: ("campaign.sending", "Campaign sending started", 1),
CampaignStatus.SENT.value: ("campaign.sent", "Campaign sent", 1),
CampaignStatus.PARTIALLY_COMPLETED.value: (
"campaign.partially_completed",
"Campaign partially completed",
4,
),
CampaignStatus.OUTCOME_UNKNOWN.value: (
"campaign.outcome_unknown",
"Campaign requires review",
5,
),
CampaignStatus.FAILED.value: ("campaign.failed", "Campaign failed", 5),
CampaignStatus.CANCELLED.value: ("campaign.cancelled", "Campaign cancelled", 2),
}
def _version_user_lock_state(version: CampaignVersion) -> str | None:
state = getattr(version, "user_lock_state", None)
if state in {"temporary", "permanent"}:
return state
return "permanent" if version.published_at else None
def _version_is_user_locked(version: CampaignVersion) -> bool:
return _version_user_lock_state(version) is not None
def _version_is_validated_and_locked(version: CampaignVersion) -> bool:
validation = (
version.validation_summary
if isinstance(version.validation_summary, dict)
else {}
)
return bool(
version.locked_at
and validation.get("ok") is True
and not _version_is_user_locked(version)
)
def _ensure_version_validated_and_locked(version: CampaignVersion) -> None:
state = _version_user_lock_state(version)
if state == "temporary":
raise QueueingError(
"This version has a temporary user lock. Unlock it before queueing, dry-run or sending."
)
if state == "permanent":
raise QueueingError(
"This version is permanently user-locked. Create an editable copy instead."
)
if not _version_is_validated_and_locked(version):
raise QueueingError(
"Campaign version must be validated and locked before building, queueing, dry-run or sending."
)
def _reviewed_needs_review_keys(version: CampaignVersion) -> set[str]:
build_summary = (
version.build_summary if isinstance(version.build_summary, dict) else {}
)
build_token = str(
build_summary.get("build_token") or build_summary.get("built_at") or ""
)
editor_state = (
version.editor_state if isinstance(version.editor_state, dict) else {}
)
review_state = (
editor_state.get("review_send")
if isinstance(editor_state.get("review_send"), dict)
else {}
)
if not build_token or str(review_state.get("build_token") or "") != build_token:
return set()
if review_state.get("inspection_complete") is not True:
return set()
return {
str(value)
for value in (review_state.get("reviewed_message_keys") or [])
if str(value).strip()
}
def _job_review_key(job: CampaignJob) -> str:
return str(job.entry_id or job.entry_index)
def _utcnow() -> datetime:
return datetime.now(timezone.utc)
def _set_version_delivery_mode(version: CampaignVersion, mode: str) -> None:
if mode not in DELIVERY_MODES:
raise QueueingError(f"Unsupported Campaign delivery mode: {mode}")
version.delivery_mode = mode
version.delivery_mode_selected_at = _utcnow()
def _get_campaign_for_tenant(
session: Session, *, campaign_id: str, tenant_id: str
) -> Campaign:
campaign = (
session.query(Campaign)
.filter(Campaign.id == campaign_id, Campaign.tenant_id == tenant_id)
.one_or_none()
)
if not campaign:
raise QueueingError(f"Campaign not found or not accessible: {campaign_id}")
return campaign
def _get_version_for_campaign(
session: Session,
campaign: Campaign,
version_id: str | None = None,
) -> CampaignVersion:
wanted = version_id or campaign.current_version_id
if not wanted:
raise QueueingError("Campaign has no selected version")
version = session.get(CampaignVersion, wanted)
if not version or version.campaign_id != campaign.id:
raise QueueingError(
f"Campaign version not found or not part of campaign: {wanted}"
)
return version
def _get_current_version(
session: Session, campaign: Campaign, version_id: str | None = None
) -> CampaignVersion:
version = _get_version_for_campaign(session, campaign, version_id)
if campaign.current_version_id != version.id:
raise QueueingError(
"Historical campaign versions cannot start a new execution. Use the explicit retry or reconciliation actions for an existing execution."
)
return version
def _celery_enabled() -> bool:
return bool(core_settings.celery_enabled)
def _should_enqueue_celery(enqueue_celery: bool) -> bool:
return bool(enqueue_celery and _celery_enabled())
def _asynchronous_delivery_mode(enqueue_celery: bool) -> str:
return (
DELIVERY_MODE_WORKER_QUEUE
if _should_enqueue_celery(enqueue_celery)
else DELIVERY_MODE_DATABASE_QUEUE
)
def _campaign_notification_body(campaign: Campaign, status: str) -> str:
return {
CampaignStatus.QUEUED.value: f"{campaign.name} has been queued for delivery.",
CampaignStatus.SENDING.value: f"{campaign.name} started sending.",
CampaignStatus.SENT.value: f"{campaign.name} finished sending.",
CampaignStatus.PARTIALLY_COMPLETED.value: f"{campaign.name} finished with partial delivery results. Review the campaign report.",
CampaignStatus.OUTCOME_UNKNOWN.value: f"{campaign.name} has unresolved delivery outcomes and needs operator review.",
CampaignStatus.FAILED.value: f"{campaign.name} failed during delivery. Review the failed jobs.",
CampaignStatus.CANCELLED.value: f"{campaign.name} was cancelled.",
}.get(status, f"{campaign.name} changed status to {status}.")
def _emit_campaign_status_notification(
session: Session,
*,
campaign: Campaign,
status: str,
previous_status: str | None = None,
version_id: str | None = None,
message: str | None = None,
) -> None:
provider = notification_dispatch_provider(get_registry())
if provider is None:
return
event = CAMPAIGN_STATUS_NOTIFICATION_EVENTS.get(status)
if event is None:
return
event_kind, title, priority = event
try:
provider.enqueue_notification(
session,
NotificationDispatchRequest(
tenant_id=campaign.tenant_id,
source_module="campaigns",
source_resource_type="campaign",
source_resource_id=campaign.id,
event_kind=event_kind,
channel="inbox",
subject=f"{title}: {campaign.name}",
body_text=message or _campaign_notification_body(campaign, status),
action_url=f"/campaigns/{campaign.id}",
priority=priority,
payload={
"campaign_id": campaign.id,
"version_id": version_id or campaign.current_version_id,
"status": status,
"previous_status": previous_status,
},
metadata={
"campaign_id": campaign.id,
"version_id": version_id or campaign.current_version_id,
},
),
enqueue_delivery=False,
)
except Exception:
return
def _celery_enqueue_send_job(job_id: str) -> None:
from govoplan_core.celery_app import celery
celery.send_task("govoplan.campaigns.send_email", args=[job_id], queue="send_email")
def _celery_enqueue_append_sent_job(job_id: str) -> None:
from govoplan_core.celery_app import celery
celery.send_task(
"govoplan.campaigns.append_sent", args=[job_id], queue="append_sent"
)
def _ensure_campaign_execution_snapshot(
session: Session,
version: CampaignVersion,
) -> None:
try:
ensure_execution_snapshot(session, version)
except ExecutionSnapshotError as exc:
raise QueueingError(str(exc)) from exc
def _ensure_campaign_approval_gate(
session: Session,
*,
tenant_id: str,
version: CampaignVersion,
) -> None:
try:
assert_campaign_approval(
session,
tenant_id=tenant_id,
version=version,
)
except CampaignApprovalGateError as exc:
raise QueueingError(str(exc)) from exc
def _queue_validation_statuses(*, include_warnings: bool) -> set[str]:
allowed_validation = {JobValidationStatus.READY.value}
if include_warnings:
allowed_validation.add(JobValidationStatus.WARNING.value)
return allowed_validation
def _campaign_jobs_for_queue(
session: Session,
*,
tenant_id: str,
version_id: str,
) -> list[CampaignJob]:
jobs = (
session.query(CampaignJob)
.filter(
CampaignJob.tenant_id == tenant_id,
CampaignJob.campaign_version_id == version_id,
)
.order_by(CampaignJob.entry_index.asc())
.all()
)
if not jobs:
raise QueueingError(
"Campaign version has no jobs. Build messages before queueing."
)
return jobs
def _initial_queue_disposition(
job: CampaignJob,
*,
allowed_validation: set[str],
reviewed_needs_review_keys: set[str],
) -> str:
if job.send_status in INITIAL_QUEUE_SKIPPED_SEND_STATUSES:
# Initial queueing never doubles as retry/reconciliation. Those states
# require the dedicated explicit actions below.
return "skipped"
if job.queue_status in INITIAL_QUEUE_SKIPPED_QUEUE_STATUSES:
return "skipped"
validation_allowed = job.validation_status in allowed_validation or (
job.validation_status == JobValidationStatus.NEEDS_REVIEW.value
and _job_review_key(job) in reviewed_needs_review_keys
)
if job.build_status != JobBuildStatus.BUILT.value or not validation_allowed:
return "blocked"
if not job.eml_local_path and not job.eml_storage_key:
job.last_error = (
"Job has no generated EML path/storage key. Rebuild with write_eml "
"enabled before queueing."
)
return "blocked"
return "queued"
def _mark_campaign_job_queued(session: Session, job: CampaignJob) -> None:
job.queue_status = JobQueueStatus.QUEUED.value
job.send_status = JobSendStatus.QUEUED.value
job.queued_at = _utcnow()
job.claimed_at = None
job.claim_token = None
job.smtp_started_at = None
job.outcome_unknown_at = None
job.last_error = None
session.add(job)
def _select_campaign_jobs_for_queue(
session: Session,
*,
jobs: list[CampaignJob],
allowed_validation: set[str],
reviewed_needs_review_keys: set[str],
dry_run: bool,
) -> tuple[list[CampaignJob], int, int]:
queued: list[CampaignJob] = []
skipped_count = 0
blocked_count = 0
for job in jobs:
disposition = _initial_queue_disposition(
job,
allowed_validation=allowed_validation,
reviewed_needs_review_keys=reviewed_needs_review_keys,
)
if disposition == "skipped":
skipped_count += 1
continue
if disposition == "blocked":
blocked_count += 1
continue
queued.append(job)
if not dry_run:
_mark_campaign_job_queued(session, job)
return queued, skipped_count, blocked_count
def synchronous_send_candidate_jobs(
version: CampaignVersion,
jobs: list[CampaignJob],
*,
include_warnings: bool = True,
) -> list[CampaignJob]:
"""Return the exact persisted job set an immediate send may attempt.
This projection is deliberately side-effect free so the delivery-options
API and adaptive documentation can explain the current built execution
without changing queue state. The final command repeats the projection
after queueing and before any SMTP effect.
"""
allowed_validation = _queue_validation_statuses(include_warnings=include_warnings)
reviewed_needs_review_keys = _reviewed_needs_review_keys(version)
candidates: list[CampaignJob] = []
for job in jobs:
if (
job.queue_status == JobQueueStatus.QUEUED.value
and job.send_status == JobSendStatus.QUEUED.value
):
candidates.append(job)
continue
if job.send_status in INITIAL_QUEUE_SKIPPED_SEND_STATUSES:
continue
if job.queue_status in INITIAL_QUEUE_SKIPPED_QUEUE_STATUSES:
continue
validation_allowed = job.validation_status in allowed_validation or (
job.validation_status == JobValidationStatus.NEEDS_REVIEW.value
and _job_review_key(job) in reviewed_needs_review_keys
)
if job.build_status != JobBuildStatus.BUILT.value or not validation_allowed:
continue
if not job.eml_local_path and not job.eml_storage_key:
continue
candidates.append(job)
return candidates
def synchronous_send_options(
session: Session,
*,
tenant_id: str,
campaign_id: str,
version_id: str | None = None,
include_warnings: bool = True,
) -> dict[str, Any]:
"""Describe the actor-independent effective delivery-mode constraints."""
campaign = _get_campaign_for_tenant(
session, campaign_id=campaign_id, tenant_id=tenant_id
)
version = _get_version_for_campaign(session, campaign, version_id=version_id)
jobs = _campaign_jobs_for_version(
session, tenant_id=tenant_id, version_id=version.id
)
candidates = synchronous_send_candidate_jobs(
version, jobs, include_warnings=include_warnings
)
try:
policy = effective_synchronous_send_policy(session, tenant_id=tenant_id)
except CampaignDeliveryPolicyError as exc:
return {
"campaign_id": campaign.id,
"version_id": version.id,
"worker_queue_available": bool(_celery_enabled()),
"synchronous_send": {
"allowed": False,
"reason": "policy_configuration_invalid",
"message": str(exc),
"eligible_recipient_job_count": len(candidates),
"policy": {},
},
}
ready = _version_is_validated_and_locked(version) and bool(version.build_summary)
eligible_count = len(candidates)
if not ready:
reason = "version_not_ready"
elif eligible_count == 0:
reason = "no_eligible_recipient_jobs"
elif eligible_count > policy.max_recipient_jobs:
reason = "recipient_limit_exceeded"
else:
reason = None
return {
"campaign_id": campaign.id,
"version_id": version.id,
"worker_queue_available": bool(_celery_enabled()),
"synchronous_send": {
"allowed": reason is None,
"reason": reason,
"eligible_recipient_job_count": eligible_count,
"policy": policy.as_dict(),
},
}
def _campaign_jobs_for_version(
session: Session,
*,
tenant_id: str,
version_id: str,
) -> list[CampaignJob]:
return (
session.query(CampaignJob)
.filter(
CampaignJob.tenant_id == tenant_id,
CampaignJob.campaign_version_id == version_id,
)
.order_by(CampaignJob.entry_index.asc())
.all()
)
def _persist_campaign_queue(
session: Session,
*,
campaign: Campaign,
version: CampaignVersion,
queued: list[CampaignJob],
delivery_mode: str,
commit: bool = True,
) -> None:
if queued:
previous_status = campaign.status
campaign.status = CampaignStatus.QUEUED.value
version.workflow_state = CampaignVersionWorkflowState.QUEUED.value
_set_version_delivery_mode(version, delivery_mode)
if version.locked_at is None:
version.locked_at = _utcnow()
session.add(version)
if previous_status != campaign.status:
_emit_campaign_status_notification(
session,
campaign=campaign,
status=campaign.status,
previous_status=previous_status,
version_id=version.id,
)
session.add(campaign)
if commit:
session.commit()
def _enqueue_campaign_jobs(queued: list[CampaignJob], *, enabled: bool) -> int:
if not enabled:
return 0
for job in queued:
_celery_enqueue_send_job(job.id)
return len(queued)
def queue_campaign_jobs(
session: Session,
*,
tenant_id: str,
campaign_id: str,
version_id: str | None = None,
enqueue_celery: bool = True,
include_warnings: bool = True,
dry_run: bool = False,
delivery_mode: str | None = None,
commit_queue: bool = True,
) -> QueueCampaignResult:
"""Move queueable DB jobs to QUEUED and optionally enqueue Celery tasks."""
selected_delivery_mode = delivery_mode or _asynchronous_delivery_mode(
enqueue_celery
)
if selected_delivery_mode not in DELIVERY_MODES:
raise QueueingError(
f"Unsupported Campaign delivery mode: {selected_delivery_mode}"
)
campaign = _get_campaign_for_tenant(
session, campaign_id=campaign_id, tenant_id=tenant_id
)
version = _get_current_version(session, campaign, version_id=version_id)
_ensure_version_validated_and_locked(version)
_ensure_campaign_execution_snapshot(session, version)
if not dry_run:
_ensure_campaign_approval_gate(session, tenant_id=tenant_id, version=version)
allowed_validation = _queue_validation_statuses(include_warnings=include_warnings)
jobs = _campaign_jobs_for_queue(
session,
tenant_id=tenant_id,
version_id=version.id,
)
reviewed_needs_review_keys = _reviewed_needs_review_keys(version)
queued, skipped_count, blocked_count = _select_campaign_jobs_for_queue(
session,
jobs=jobs,
allowed_validation=allowed_validation,
reviewed_needs_review_keys=reviewed_needs_review_keys,
dry_run=dry_run,
)
if not dry_run:
_persist_campaign_queue(
session,
campaign=campaign,
version=version,
queued=queued,
delivery_mode=selected_delivery_mode,
commit=commit_queue,
)
enqueued_count = _enqueue_campaign_jobs(
queued,
enabled=_should_enqueue_celery(enqueue_celery) and not dry_run,
)
return QueueCampaignResult(
campaign_id=campaign.id,
version_id=version.id,
queued_count=len(queued),
skipped_count=skipped_count,
blocked_count=blocked_count,
enqueued_count=enqueued_count,
delivery_mode=selected_delivery_mode,
worker_queue_available=_celery_enabled(),
dry_run=dry_run,
)
def send_campaign_now(
session: Session,
*,
tenant_id: str,
campaign_id: str,
version_id: str | None = None,
include_warnings: bool = True,
dry_run: bool = False,
use_rate_limit: bool = True,
enqueue_imap_task: bool = False,
) -> SendCampaignNowResult:
"""Queue and send all eligible jobs for one campaign synchronously.
The effective deployment/tenant limit is checked against the immutable,
persisted job set before queue state changes. After queueing, every message
and transport revision is preflighted before the first SMTP effect starts.
"""
campaign = _get_campaign_for_tenant(
session, campaign_id=campaign_id, tenant_id=tenant_id
)
version = _get_current_version(session, campaign, version_id=version_id)
_ensure_version_validated_and_locked(version)
_ensure_campaign_execution_snapshot(session, version)
try:
synchronous_policy = effective_synchronous_send_policy(
session, tenant_id=tenant_id
)
except CampaignDeliveryPolicyError as exc:
raise SynchronousSendRejected(
f"Synchronous Campaign delivery is disabled by an invalid policy configuration: {exc}",
reason="policy_configuration_invalid",
) from exc
initial_jobs = _campaign_jobs_for_queue(
session,
tenant_id=tenant_id,
version_id=version.id,
)
initial_candidates = synchronous_send_candidate_jobs(
version,
initial_jobs,
include_warnings=include_warnings,
)
_ensure_synchronous_send_count_allowed(
len(initial_candidates),
policy=synchronous_policy,
)
queue_result = queue_campaign_jobs(
session,
tenant_id=tenant_id,
campaign_id=campaign.id,
version_id=version.id,
include_warnings=include_warnings,
enqueue_celery=False,
dry_run=dry_run,
delivery_mode=DELIVERY_MODE_SYNCHRONOUS,
commit_queue=False,
)
if dry_run:
return SendCampaignNowResult(
campaign_id=campaign.id,
version_id=version.id,
attempted_count=0,
sent_count=0,
failed_count=0,
outcome_unknown_count=0,
skipped_count=queue_result.skipped_count + queue_result.blocked_count,
preflight_count=len(initial_candidates),
synchronous_send_policy=synchronous_policy.as_dict(),
dry_run=True,
results=[queue_result.as_dict()],
)
jobs = [
job
for job in _campaign_jobs_for_version(
session,
tenant_id=tenant_id,
version_id=version.id,
)
if job.queue_status == JobQueueStatus.QUEUED.value
and job.send_status == JobSendStatus.QUEUED.value
]
# Repeat the hard bound against the post-queue set. This closes the window
# where a concurrent queue operation could otherwise enlarge an immediate
# run between the initial decision and the first provider effect.
_ensure_synchronous_send_count_allowed(len(jobs), policy=synchronous_policy)
delivery_contexts = _preflight_synchronous_send_batch(
session,
version=version,
jobs=jobs,
policy=synchronous_policy,
)
# Queue state and its inbox notification become durable only after every
# message and the selected transport revision have passed preflight. This
# preserves late-ack recovery without leaving rejected work eligible for a
# background worker.
session.commit()
results: list[dict[str, Any]] = []
sent_count = 0
failed_count = 0
outcome_unknown_count = 0
skipped_after_queue = 0
for job in jobs:
try:
claimed = _claimed_campaign_job_for_delivery(session, job)
if isinstance(claimed, SendJobResult):
result = claimed
else:
claimed_job, claim_token = claimed
result = _send_claimed_campaign_job(
session,
job=claimed_job,
claim_token=claim_token,
context=delivery_contexts[job.id],
use_rate_limit=use_rate_limit,
enqueue_imap_task=enqueue_imap_task,
)
result_dict = result.as_dict()
results.append(result_dict)
if result.status in DELIVERY_ACCEPTED_STATUSES | {"already_accepted"}:
sent_count += 1
elif result.status == JobSendStatus.OUTCOME_UNKNOWN.value:
outcome_unknown_count += 1
else:
skipped_after_queue += 1
except Exception as exc: # keep sending other jobs and return per-job details
failed_count += 1
results.append({"job_id": job.id, "status": "failed", "message": str(exc)})
return SendCampaignNowResult(
campaign_id=campaign.id,
version_id=version.id,
attempted_count=len(jobs),
sent_count=sent_count,
failed_count=failed_count,
outcome_unknown_count=outcome_unknown_count,
skipped_count=queue_result.skipped_count
+ queue_result.blocked_count
+ skipped_after_queue,
preflight_count=len(delivery_contexts),
synchronous_send_policy=synchronous_policy.as_dict(),
dry_run=False,
results=results,
)
def _ensure_synchronous_send_count_allowed(
eligible_count: int,
*,
policy: SynchronousSendPolicy,
) -> None:
if eligible_count <= 0:
raise SynchronousSendRejected(
"No eligible built recipient messages are available for synchronous delivery.",
reason="no_eligible_recipient_jobs",
eligible_count=eligible_count,
policy=policy,
)
if eligible_count > policy.max_recipient_jobs:
raise SynchronousSendRejected(
(
f"This built run contains {eligible_count} eligible recipient jobs, above the effective "
f"synchronous limit of {policy.max_recipient_jobs}. Queue it for background workers instead."
),
reason="recipient_limit_exceeded",
eligible_count=eligible_count,
policy=policy,
)
def _preflight_synchronous_send_batch(
session: Session,
*,
version: CampaignVersion,
jobs: list[CampaignJob],
policy: SynchronousSendPolicy,
) -> dict[str, _SendJobDeliveryContext]:
"""Freeze and validate every local input before the first SMTP effect."""
contexts: dict[str, _SendJobDeliveryContext] = {}
try:
for job in jobs:
state_result = _preflight_send_campaign_job(session, job, dry_run=True)
if state_result is not None:
raise SendJobError(
f"Job {job.id} changed to {state_result.status} during synchronous preflight"
)
contexts[job.id] = _send_job_delivery_context(session, job)
mail_contexts = [
context
for context in contexts.values()
if getattr(context.snapshot, "uses_mail", True)
]
if mail_contexts:
profile_summary = profile_delivery_summary(session, version)
expected_revision = mail_contexts[0].snapshot.smtp_transport_revision
if profile_summary.get("smtp_transport_revision") != expected_revision:
raise SendJobError(
"The selected Mail profile changed after this execution was built. Revalidate and rebuild before delivery."
)
except (
ExecutionSnapshotError,
MailProfileError,
SmtpConfigurationError,
SendJobError,
OSError,
) as exc:
raise SynchronousSendRejected(
(
"Synchronous preflight stopped before contacting SMTP because one or more frozen message "
"inputs or the selected Mail profile no longer match the reviewed build. Rebuild and review "
"the version before trying again."
),
reason="batch_preflight_failed",
eligible_count=len(jobs),
policy=policy,
) from exc
return contexts
def enqueue_existing_queued_jobs(
session: Session, *, tenant_id: str, campaign_id: str
) -> int:
if not _celery_enabled():
return 0
campaign = _get_campaign_for_tenant(
session, campaign_id=campaign_id, tenant_id=tenant_id
)
jobs = (
session.query(CampaignJob)
.filter(
CampaignJob.tenant_id == tenant_id,
CampaignJob.campaign_id == campaign.id,
CampaignJob.queue_status == JobQueueStatus.QUEUED.value,
CampaignJob.send_status.in_(
[JobSendStatus.QUEUED.value, JobSendStatus.FAILED_TEMPORARY.value]
),
)
.order_by(CampaignJob.entry_index.asc())
.all()
)
for job in jobs:
_celery_enqueue_send_job(job.id)
return len(jobs)
def pause_campaign_jobs(
session: Session, *, tenant_id: str, campaign_id: str
) -> dict[str, Any]:
campaign = _get_campaign_for_tenant(
session, campaign_id=campaign_id, tenant_id=tenant_id
)
changed = (
session.query(CampaignJob)
.filter(
CampaignJob.tenant_id == tenant_id,
CampaignJob.campaign_id == campaign.id,
CampaignJob.queue_status == JobQueueStatus.QUEUED.value,
)
.update(
{CampaignJob.queue_status: JobQueueStatus.PAUSED.value},
synchronize_session=False,
)
)
if changed:
campaign.status = CampaignStatus.READY_TO_QUEUE.value
session.add(campaign)
session.commit()
return {"campaign_id": campaign.id, "paused_count": int(changed)}
def resume_campaign_jobs(
session: Session, *, tenant_id: str, campaign_id: str, enqueue_celery: bool = True
) -> dict[str, Any]:
campaign = _get_campaign_for_tenant(
session, campaign_id=campaign_id, tenant_id=tenant_id
)
jobs = (
session.query(CampaignJob)
.filter(
CampaignJob.tenant_id == tenant_id,
CampaignJob.campaign_id == campaign.id,
CampaignJob.queue_status == JobQueueStatus.PAUSED.value,
)
.order_by(CampaignJob.entry_index.asc())
.all()
)
for job in jobs:
job.queue_status = JobQueueStatus.QUEUED.value
job.send_status = JobSendStatus.QUEUED.value
session.add(job)
if jobs:
delivery_mode = _asynchronous_delivery_mode(enqueue_celery)
for version_id in {job.campaign_version_id for job in jobs}:
version = session.get(CampaignVersion, version_id)
if version is not None and version.campaign_id == campaign.id:
_set_version_delivery_mode(version, delivery_mode)
session.add(version)
previous_status = campaign.status
campaign.status = CampaignStatus.QUEUED.value
session.add(campaign)
if previous_status != campaign.status:
_emit_campaign_status_notification(
session,
campaign=campaign,
status=campaign.status,
previous_status=previous_status,
)
session.commit()
enqueued_count = 0
if _should_enqueue_celery(enqueue_celery):
for job in jobs:
_celery_enqueue_send_job(job.id)
enqueued_count += 1
return {
"campaign_id": campaign.id,
"resumed_count": len(jobs),
"enqueued_count": enqueued_count,
}
def cancel_campaign_jobs(
session: Session, *, tenant_id: str, campaign_id: str
) -> dict[str, Any]:
campaign = _get_campaign_for_tenant(
session, campaign_id=campaign_id, tenant_id=tenant_id
)
jobs = (
session.query(CampaignJob)
.filter(
CampaignJob.tenant_id == tenant_id, CampaignJob.campaign_id == campaign.id
)
.all()
)
cancelled_count = 0
protected_count = 0
skipped_count = 0
version_ids: set[str] = set()
for job in jobs:
version_ids.add(job.campaign_version_id)
if job.send_status == JobSendStatus.SKIPPED.value:
skipped_count += 1
continue
if job.send_status in DELIVERY_ACCEPTED_STATUSES | {
JobSendStatus.OUTCOME_UNKNOWN.value,
JobSendStatus.CLAIMED.value,
JobSendStatus.SENDING.value,
}:
protected_count += 1
continue
if job.send_status != JobSendStatus.CANCELLED.value:
job.queue_status = JobQueueStatus.CANCELLED.value
job.send_status = JobSendStatus.CANCELLED.value
job.claim_token = None
session.add(job)
cancelled_count += 1
for version_id in version_ids:
_update_campaign_after_job(session, campaign.id, version_id)
session.commit()
return {
"campaign_id": campaign.id,
"cancelled_count": cancelled_count,
"protected_count": protected_count,
"skipped_count": skipped_count,
"campaign_status": campaign.status,
}
def _eligible_jobs_for_explicit_action(
session: Session,
*,
tenant_id: str,
campaign: Campaign,
version: CampaignVersion,
job_ids: list[str] | None,
) -> list[CampaignJob]:
query = session.query(CampaignJob).filter(
CampaignJob.tenant_id == tenant_id,
CampaignJob.campaign_id == campaign.id,
CampaignJob.campaign_version_id == version.id,
)
if job_ids:
query = query.filter(CampaignJob.id.in_(job_ids))
return query.order_by(CampaignJob.entry_index.asc()).all()
def queue_failed_jobs_for_retry(
session: Session,
*,
tenant_id: str,
campaign_id: str,
version_id: str | None = None,
job_ids: list[str] | None = None,
include_permanent: bool = False,
force_max_attempts: bool = False,
enqueue_celery: bool = True,
dry_run: bool = False,
) -> dict[str, Any]:
"""Queue known failures and incomplete multi-channel deliveries.
Accepted SMTP attempts and accepted Postbox targets remain immutable and
are skipped by the channel orchestrator. Unknown outcomes are deliberately
excluded until an operator reconciles them.
"""
campaign = _get_campaign_for_tenant(
session, campaign_id=campaign_id, tenant_id=tenant_id
)
version = _get_version_for_campaign(session, campaign, version_id=version_id)
_ensure_version_validated_and_locked(version)
snapshot = ensure_execution_snapshot(session, version)
allowed = {
JobSendStatus.FAILED_TEMPORARY.value,
JobSendStatus.PARTIALLY_ACCEPTED.value,
}
if include_permanent:
allowed.add(JobSendStatus.FAILED_PERMANENT.value)
selected: list[CampaignJob] = []
skipped: list[dict[str, str]] = []
for job in _eligible_jobs_for_explicit_action(
session,
tenant_id=tenant_id,
campaign=campaign,
version=version,
job_ids=job_ids,
):
if job.send_status not in allowed:
skipped.append(
{
"job_id": job.id,
"reason": f"status {job.send_status} is not an explicit retry candidate",
}
)
continue
delivery_rounds = max(
job.attempt_count,
int(
session.query(func.max(PostboxDeliveryAttempt.attempt_number))
.filter(PostboxDeliveryAttempt.job_id == job.id)
.scalar()
or 0
),
)
if (
not force_max_attempts
and delivery_rounds >= snapshot.delivery.retry.max_attempts
):
skipped.append(
{"job_id": job.id, "reason": "configured maximum attempts reached"}
)
continue
selected.append(job)
if not dry_run:
job.queue_status = JobQueueStatus.QUEUED.value
job.send_status = JobSendStatus.QUEUED.value
job.queued_at = _utcnow()
job.claimed_at = None
job.claim_token = None
job.smtp_started_at = None
job.outcome_unknown_at = None
session.add(job)
if not dry_run:
if selected:
campaign.status = CampaignStatus.QUEUED.value
version.workflow_state = CampaignVersionWorkflowState.QUEUED.value
_set_version_delivery_mode(
version,
_asynchronous_delivery_mode(enqueue_celery),
)
session.add(campaign)
session.add(version)
session.commit()
enqueued = 0
if _should_enqueue_celery(enqueue_celery) and not dry_run:
for job in selected:
_celery_enqueue_send_job(job.id)
enqueued += 1
return {
"campaign_id": campaign.id,
"version_id": version.id,
"action": "retry_failed",
"selected_count": len(selected),
"enqueued_count": enqueued,
"skipped": skipped,
"dry_run": dry_run,
}
def queue_unattempted_jobs(
session: Session,
*,
tenant_id: str,
campaign_id: str,
version_id: str | None = None,
job_ids: list[str] | None = None,
enqueue_celery: bool = True,
dry_run: bool = False,
) -> dict[str, Any]:
"""Explicitly queue built jobs that have never started an SMTP attempt."""
campaign = _get_campaign_for_tenant(
session, campaign_id=campaign_id, tenant_id=tenant_id
)
version = _get_version_for_campaign(session, campaign, version_id=version_id)
_ensure_version_validated_and_locked(version)
ensure_execution_snapshot(session, version)
if not dry_run:
_ensure_campaign_approval_gate(session, tenant_id=tenant_id, version=version)
selected: list[CampaignJob] = []
skipped: list[dict[str, str]] = []
for job in _eligible_jobs_for_explicit_action(
session,
tenant_id=tenant_id,
campaign=campaign,
version=version,
job_ids=job_ids,
):
eligible = (
job.attempt_count == 0
and job.postbox_attempt_count == 0
and job.send_status
in {JobSendStatus.NOT_QUEUED.value, JobSendStatus.CANCELLED.value}
and job.build_status == JobBuildStatus.BUILT.value
and job.validation_status in QUEUEABLE_VALIDATION_STATUSES
)
if not eligible:
skipped.append(
{
"job_id": job.id,
"reason": "job is not an unattempted queueable message",
}
)
continue
selected.append(job)
if not dry_run:
job.queue_status = JobQueueStatus.QUEUED.value
job.send_status = JobSendStatus.QUEUED.value
job.queued_at = _utcnow()
job.claimed_at = None
job.claim_token = None
job.smtp_started_at = None
job.outcome_unknown_at = None
job.last_error = None
session.add(job)
if not dry_run:
if selected:
campaign.status = CampaignStatus.QUEUED.value
version.workflow_state = CampaignVersionWorkflowState.QUEUED.value
_set_version_delivery_mode(
version,
_asynchronous_delivery_mode(enqueue_celery),
)
session.add(campaign)
session.add(version)
session.commit()
enqueued = 0
if _should_enqueue_celery(enqueue_celery) and not dry_run:
for job in selected:
_celery_enqueue_send_job(job.id)
enqueued += 1
return {
"campaign_id": campaign.id,
"version_id": version.id,
"action": "send_unattempted",
"selected_count": len(selected),
"enqueued_count": enqueued,
"skipped": skipped,
"dry_run": dry_run,
}
def send_single_campaign_job(
session: Session,
*,
tenant_id: str,
campaign_id: str,
job_id: str,
kind: str,
idempotency_key: str,
actor_user_id: str | None,
actor_api_key_id: str | None = None,
reason: str | None = None,
action_context: dict[str, Any] | None = None,
include_warnings: bool = True,
use_rate_limit: bool = True,
enqueue_imap_task: bool = False,
) -> dict[str, Any]:
"""Perform one explicit, idempotent action for one immutable built message."""
campaign = _get_campaign_for_tenant(
session, campaign_id=campaign_id, tenant_id=tenant_id
)
job = session.get(CampaignJob, job_id)
if not job or job.tenant_id != tenant_id or job.campaign_id != campaign.id:
raise QueueingError("Campaign job not found or not accessible")
normalized_kind = str(kind or "").strip()
if normalized_kind not in {"test", "single_send", "single_resend"}:
raise QueueingError(
"Single-message action kind must be test, single_send, or single_resend"
)
clean_key = str(idempotency_key or "").strip()
if not clean_key or len(clean_key) > 200:
raise QueueingError("A bounded idempotency key is required")
clean_reason = " ".join(str(reason or "").split()) or None
if normalized_kind == "single_resend" and not clean_reason:
raise QueueingError("Single-message resend requires a reason")
if clean_reason and len(clean_reason) > 2000:
raise QueueingError("Single-message action reason is too long")
version = _get_version_for_campaign(
session,
campaign,
version_id=job.campaign_version_id,
)
if normalized_kind == "single_send" and campaign.current_version_id != version.id:
raise QueueingError(
"An initial official send is only available for the current campaign version."
)
if normalized_kind != "test":
_ensure_campaign_approval_gate(session, tenant_id=tenant_id, version=version)
recipients = _recipients_from_job(job)
recipient_manifest_sha256 = hashlib.sha256(
json.dumps(
recipients,
ensure_ascii=True,
separators=(",", ":"),
).encode("utf-8")
).hexdigest()
canonical_request_hash = hashlib.sha256(
json.dumps(
{
"tenant_id": tenant_id,
"campaign_id": campaign.id,
"version_id": version.id,
"job_id": job.id,
"kind": normalized_kind,
"message_sha256": job.eml_sha256,
"recipient_manifest_sha256": recipient_manifest_sha256,
"reason": clean_reason,
"context": _bounded_single_action_context(action_context),
},
sort_keys=True,
separators=(",", ":"),
ensure_ascii=True,
).encode("utf-8")
).hexdigest()
action, duplicate = _create_single_message_action(
session,
tenant_id=tenant_id,
campaign=campaign,
version=version,
job=job,
kind=normalized_kind,
idempotency_key=clean_key,
canonical_request_hash=canonical_request_hash,
reason=clean_reason,
action_context=_bounded_single_action_context(action_context),
actor_user_id=actor_user_id,
actor_api_key_id=actor_api_key_id,
recipient_manifest_sha256=recipient_manifest_sha256,
recipient_count=len(recipients),
)
if duplicate:
return _single_message_action_response(action, duplicate=True)
try:
_validate_single_message_action(
version=version,
job=job,
kind=normalized_kind,
include_warnings=include_warnings,
)
delivery_context = _send_job_delivery_context(session, job)
except Exception as exc:
_finish_single_message_action(
session,
action=action,
status="initiation_failed",
error_type=exc.__class__.__name__,
error_message=str(exc),
final_send_status=job.send_status,
)
raise
if normalized_kind in {"test", "single_resend"}:
return _send_single_message_direct(
session,
campaign=campaign,
version=version,
job=job,
action=action,
delivery_context=delivery_context,
use_rate_limit=use_rate_limit,
enqueue_imap_task=enqueue_imap_task,
)
try:
queue_action = _prepare_single_job_for_send(
session,
version=version,
job=job,
include_warnings=include_warnings,
dry_run=False,
)
except Exception as exc:
_finish_single_message_action(
session,
action=action,
status="initiation_failed",
error_type=exc.__class__.__name__,
error_message=str(exc),
final_send_status=job.send_status,
)
raise
if queue_action == "queued":
previous_status = campaign.status
campaign.status = CampaignStatus.QUEUED.value
version.workflow_state = CampaignVersionWorkflowState.QUEUED.value
_set_version_delivery_mode(version, DELIVERY_MODE_SYNCHRONOUS)
session.add(campaign)
session.add(version)
session.commit()
if previous_status != campaign.status:
_emit_campaign_status_notification(
session,
campaign=campaign,
status=campaign.status,
previous_status=previous_status,
version_id=version.id,
)
action_attempt = _start_single_message_action_attempt(session, action)
try:
result = send_campaign_job(
session,
job_id=job.id,
dry_run=False,
use_rate_limit=use_rate_limit,
enqueue_imap_task=enqueue_imap_task,
)
except Exception as exc:
latest = (
session.query(SendAttempt)
.filter(SendAttempt.job_id == job.id)
.order_by(SendAttempt.attempt_number.desc())
.first()
)
if latest is not None and latest.started_at is not None:
action.linked_send_attempt_id = latest.id
status = _single_action_status_from_job(job)
_finish_single_message_action(
session,
action=action,
attempt=action_attempt,
status=status,
error_type=exc.__class__.__name__,
error_message=str(exc),
final_send_status=job.send_status,
)
raise
latest = (
session.query(SendAttempt)
.filter(SendAttempt.job_id == job.id)
.order_by(SendAttempt.attempt_number.desc())
.first()
)
if latest is not None:
action.linked_send_attempt_id = latest.id
status = _single_action_status_from_result(result.status)
_finish_single_message_action(
session,
action=action,
attempt=action_attempt,
status=status,
final_send_status=job.send_status,
)
response = _single_message_action_response(action)
response["queue_action"] = queue_action
response["result"] = result.as_dict()
return response
def _bounded_single_action_context(
value: dict[str, Any] | None,
) -> dict[str, Any]:
if not isinstance(value, dict):
return {}
result: dict[str, Any] = {}
for raw_key, raw_value in list(value.items())[:20]:
key = str(raw_key).strip()[:80]
if not key:
continue
if isinstance(raw_value, (str, int, float, bool)) or raw_value is None:
result[key] = (
str(raw_value)[:500] if isinstance(raw_value, str) else raw_value
)
redacted = redact_secret_values(result)
return redacted if isinstance(redacted, dict) else {}
def _create_single_message_action(
session: Session,
*,
tenant_id: str,
campaign: Campaign,
version: CampaignVersion,
job: CampaignJob,
kind: str,
idempotency_key: str,
canonical_request_hash: str,
reason: str | None,
action_context: dict[str, Any],
actor_user_id: str | None,
actor_api_key_id: str | None,
recipient_manifest_sha256: str,
recipient_count: int,
) -> tuple[CampaignMessageAction, bool]:
existing = (
session.query(CampaignMessageAction)
.filter(
CampaignMessageAction.tenant_id == tenant_id,
CampaignMessageAction.idempotency_key == idempotency_key,
)
.first()
)
if existing is not None:
if existing.canonical_request_hash != canonical_request_hash:
raise QueueingError(
"The idempotency key is already bound to another single-message action"
)
_freeze_replayed_in_progress_action(session, existing)
return existing, True
action = CampaignMessageAction(
tenant_id=tenant_id,
campaign_id=campaign.id,
campaign_version_id=version.id,
job_id=job.id,
kind=kind,
idempotency_key=idempotency_key,
canonical_request_hash=canonical_request_hash,
reason=reason,
context=action_context,
actor_user_id=actor_user_id,
actor_api_key_id=actor_api_key_id,
message_sha256=str(job.eml_sha256 or ""),
message_size_bytes=job.eml_size_bytes,
recipient_manifest_sha256=recipient_manifest_sha256,
recipient_count=recipient_count,
prior_send_status=job.send_status,
prior_attempt_count=job.attempt_count,
status="initiated",
)
try:
with session.begin_nested():
session.add(action)
session.flush()
except IntegrityError:
existing = (
session.query(CampaignMessageAction)
.filter(
CampaignMessageAction.tenant_id == tenant_id,
CampaignMessageAction.idempotency_key == idempotency_key,
)
.first()
)
if (
existing is None
or existing.canonical_request_hash != canonical_request_hash
):
raise QueueingError(
"The idempotency key is already bound to another single-message action"
) from None
_freeze_replayed_in_progress_action(session, existing)
return existing, True
session.commit()
audit_event(
session,
tenant_id=tenant_id,
user_id=actor_user_id,
api_key_id=actor_api_key_id,
action="campaign.message_action_initiated",
object_type="campaign_message_action",
object_id=action.id,
details={
"campaign_id": campaign.id,
"campaign_version_id": version.id,
"job_id": job.id,
"kind": kind,
"message_sha256": action.message_sha256,
"recipient_manifest_sha256": recipient_manifest_sha256,
"recipient_count": recipient_count,
"prior_send_status": job.send_status,
"reason": reason,
},
)
session.commit()
return action, False
def _freeze_replayed_in_progress_action(
session: Session,
action: CampaignMessageAction,
) -> None:
if action.status != "effect_in_progress":
return
now = _utcnow()
action.status = "outcome_unknown"
action.completed_at = now
action.error_type = "OutcomeUnknown"
action.error_message = (
"A repeated command found an unfinished provider effect. "
"Automatic resend was stopped."
)
attempt = (
session.query(CampaignMessageActionAttempt)
.filter(CampaignMessageActionAttempt.action_id == action.id)
.order_by(CampaignMessageActionAttempt.attempt_number.desc())
.first()
)
if attempt is not None:
attempt.status = "outcome_unknown"
attempt.completed_at = now
attempt.outcome_code = "replayed_in_progress"
attempt.diagnostic_summary = action.error_message
session.add(attempt)
session.add(action)
session.commit()
def _validate_single_message_action(
*,
version: CampaignVersion,
job: CampaignJob,
kind: str,
include_warnings: bool,
) -> None:
_ensure_version_validated_and_locked(version)
if job.build_status != JobBuildStatus.BUILT.value:
raise QueueingError("This message has not been built yet.")
if not _single_job_validation_allowed(
version,
job,
include_warnings=include_warnings,
):
raise QueueingError(
f"This message cannot be sent while validation status is {job.validation_status}."
)
if not job.eml_sha256 or (not job.eml_local_path and not job.eml_storage_key):
raise QueueingError(
"This message has no immutable generated EML evidence. Rebuild the campaign before sending."
)
if DeliveryChannelPolicy(job.delivery_channel_policy) != DeliveryChannelPolicy.MAIL:
raise QueueingError(
"Single-message SMTP actions currently require a Mail-only delivery policy."
)
if kind == "single_send":
if job.send_status in {
JobSendStatus.CLAIMED.value,
JobSendStatus.SENDING.value,
JobSendStatus.OUTCOME_UNKNOWN.value,
}:
raise QueueingError(
f"This message is in delivery state {job.send_status}; reconcile or wait before sending it."
)
if job.attempt_count > 0 or job.send_status not in {
JobSendStatus.NOT_QUEUED.value,
JobSendStatus.CANCELLED.value,
JobSendStatus.QUEUED.value,
}:
raise QueueingError(
"Single-send is only available for a message without an official delivery attempt."
)
elif kind == "single_resend" and job.send_status in {
JobSendStatus.CLAIMED.value,
JobSendStatus.SENDING.value,
JobSendStatus.OUTCOME_UNKNOWN.value,
}:
raise QueueingError(
f"This message is in delivery state {job.send_status}; reconcile it before resending."
)
def _start_single_message_action_attempt(
session: Session,
action: CampaignMessageAction,
) -> CampaignMessageActionAttempt:
now = _utcnow()
attempt = CampaignMessageActionAttempt(
action_id=action.id,
attempt_number=1,
status="effect_in_progress",
started_at=now,
effect_started_at=now,
)
action.status = "effect_in_progress"
action.effect_started_at = now
session.add(action)
session.add(attempt)
session.commit()
return attempt
def _finish_single_message_action(
session: Session,
*,
action: CampaignMessageAction,
status: str,
attempt: CampaignMessageActionAttempt | None = None,
accepted_count: int = 0,
refused_recipients: dict[str, dict[str, Any]] | None = None,
error_type: str | None = None,
error_message: str | None = None,
final_send_status: str | None = None,
) -> None:
now = _utcnow()
refusals = refused_recipients or {}
action.status = status
action.accepted_count = accepted_count
action.refused_count = len(refusals)
action.refusal_summary = dict(
Counter(
str(item.get("classification") or "unknown") for item in refusals.values()
)
)
action.error_type = str(error_type)[:120] if error_type else None
action.error_message = (
" ".join(str(error_message).split())[:500] if error_message else None
)
action.final_send_status = final_send_status
action.completed_at = now
session.add(action)
if attempt is not None:
attempt.status = status
attempt.completed_at = now
attempt.accepted_count = accepted_count
attempt.refused_count = len(refusals)
attempt.outcome_code = action.error_type or status
attempt.diagnostic_summary = action.error_message
session.add(attempt)
session.commit()
audit_event(
session,
tenant_id=action.tenant_id,
user_id=action.actor_user_id,
api_key_id=action.actor_api_key_id,
action="campaign.message_action_completed",
object_type="campaign_message_action",
object_id=action.id,
details={
"campaign_id": action.campaign_id,
"campaign_version_id": action.campaign_version_id,
"job_id": action.job_id,
"kind": action.kind,
"status": status,
"message_sha256": action.message_sha256,
"recipient_manifest_sha256": action.recipient_manifest_sha256,
"recipient_count": action.recipient_count,
"accepted_count": accepted_count,
"refused_count": len(refusals),
"refusal_summary": action.refusal_summary,
"prior_send_status": action.prior_send_status,
"final_send_status": final_send_status,
"linked_send_attempt_id": action.linked_send_attempt_id,
"reason": action.reason,
},
)
session.commit()
def _single_action_status_from_result(status: str) -> str:
if status in DELIVERY_ACCEPTED_STATUSES or status == "already_accepted":
return "accepted"
if status == JobSendStatus.OUTCOME_UNKNOWN.value:
return "outcome_unknown"
if status == JobSendStatus.FAILED_TEMPORARY.value:
return "failed_temporary"
if status == JobSendStatus.FAILED_PERMANENT.value:
return "failed_permanent"
return status
def _single_action_status_from_job(job: CampaignJob) -> str:
return _single_action_status_from_result(job.send_status)
def _single_message_action_response(
action: CampaignMessageAction,
*,
duplicate: bool = False,
) -> dict[str, Any]:
return {
"campaign_id": action.campaign_id,
"version_id": action.campaign_version_id,
"job_id": action.job_id,
"action_id": action.id,
"action": action.kind,
"kind": action.kind,
"status": action.status,
"message_sha256": action.message_sha256,
"recipient_manifest_sha256": action.recipient_manifest_sha256,
"recipient_count": action.recipient_count,
"prior_send_status": action.prior_send_status,
"final_send_status": action.final_send_status,
"accepted_count": action.accepted_count,
"refused_count": action.refused_count,
"refusal_summary": dict(action.refusal_summary or {}),
"reason": action.reason,
"linked_send_attempt_id": action.linked_send_attempt_id,
"created_at": action.created_at,
"effect_started_at": action.effect_started_at,
"completed_at": action.completed_at,
"duplicate": duplicate,
"result": SendJobResult(
job_id=action.job_id,
status=action.status,
attempt_number=action.prior_attempt_count,
message=action.error_message,
).as_dict(),
}
def _send_single_message_direct(
session: Session,
*,
campaign: Campaign,
version: CampaignVersion,
job: CampaignJob,
action: CampaignMessageAction,
delivery_context: _SendJobDeliveryContext,
use_rate_limit: bool,
enqueue_imap_task: bool,
) -> dict[str, Any]:
if (
delivery_context.envelope_from is None
or not delivery_context.envelope_recipients
):
_finish_single_message_action(
session,
action=action,
status="initiation_failed",
error_type="SmtpConfigurationError",
error_message="Mail delivery has no frozen envelope sender or recipients.",
final_send_status=job.send_status,
)
raise SmtpConfigurationError(
"Mail delivery has no frozen envelope sender or recipients."
)
mail_integration().wait_for_rate_limit(
key=f"tenant:{job.tenant_id}:campaign:{job.campaign_id}:single-action",
messages_per_minute=delivery_context.snapshot.delivery.rate_limit.messages_per_minute,
enabled=use_rate_limit,
)
attempt = _start_single_message_action_attempt(session, action)
try:
result = mail_integration().send_campaign_email_bytes(
session,
tenant_id=job.tenant_id,
campaign_id=job.campaign_id,
profile_id=delivery_context.snapshot.mail_profile_id,
message_bytes=delivery_context.message_bytes,
envelope_from=delivery_context.envelope_from,
envelope_recipients=delivery_context.envelope_recipients,
from_header=_from_header_from_job(job),
expected_smtp_transport_revision=(
delivery_context.snapshot.smtp_transport_revision or ""
),
smtp_server_id=delivery_context.snapshot.smtp_server_id,
smtp_credential_id=delivery_context.snapshot.smtp_credential_id,
)
except SmtpSendError as exc:
if exc.outcome_unknown:
status = "outcome_unknown"
elif exc.temporary:
status = "failed_temporary"
else:
status = "failed_permanent"
_finish_single_message_action(
session,
action=action,
attempt=attempt,
status=status,
error_type=exc.__class__.__name__,
error_message=str(exc),
final_send_status=job.send_status,
)
return _single_message_action_response(action)
except (MailProfileError, SmtpConfigurationError, SendJobError, OSError) as exc:
_finish_single_message_action(
session,
action=action,
attempt=attempt,
status="failed_permanent",
error_type=exc.__class__.__name__,
error_message=str(exc),
final_send_status=job.send_status,
)
return _single_message_action_response(action)
except Exception:
_finish_single_message_action(
session,
action=action,
attempt=attempt,
status="outcome_unknown",
error_type="OutcomeUnknown",
error_message=(
"The provider effect started, but the delivery outcome could "
"not be established."
),
final_send_status=job.send_status,
)
return _single_message_action_response(action)
refusals = dict(result.refused_recipients)
accepted_count = result.accepted_count
if accepted_count <= 0:
classifications = {
str(item.get("classification") or "unknown") for item in refusals.values()
}
if classifications and classifications <= {"temporary"}:
status = "failed_temporary"
elif "unknown" in classifications:
status = "outcome_unknown"
else:
status = "failed_permanent"
else:
status = "accepted_with_refusals" if refusals else "accepted"
if action.kind == "single_resend" and accepted_count > 0:
if job.send_status not in FULLY_ACCEPTED_STATUSES:
job.queue_status = JobQueueStatus.DRAFT.value
job.send_status = JobSendStatus.SMTP_ACCEPTED.value
job.sent_at = _utcnow()
job.outcome_unknown_at = None
job.last_error = (
"Some SMTP recipients were refused during explicit resend."
if refusals
else None
)
job.imap_status = (
JobImapStatus.PENDING.value
if delivery_context.snapshot.delivery.imap_append_sent.enabled
else JobImapStatus.NOT_REQUESTED.value
)
files_integration().mark_job_attachment_uses_sent(session, job)
session.add(job)
if campaign.current_version_id == version.id:
_update_campaign_after_job(
session,
job.campaign_id,
job.campaign_version_id,
)
session.commit()
if (
enqueue_imap_task
and _celery_enabled()
and job.imap_status == JobImapStatus.PENDING.value
):
try:
_celery_enqueue_append_sent_job(job.id)
except Exception:
pass
_finish_single_message_action(
session,
action=action,
attempt=attempt,
status=status,
accepted_count=accepted_count,
refused_recipients=refusals,
error_type=(
None if accepted_count > 0 else f"Smtp{status.title().replace('_', '')}"
),
error_message=(
None if accepted_count > 0 else "SMTP did not accept an envelope recipient."
),
final_send_status=job.send_status,
)
return _single_message_action_response(action)
def _prepare_single_job_for_send(
session: Session,
*,
version: CampaignVersion,
job: CampaignJob,
include_warnings: bool,
dry_run: bool,
) -> str:
if (
job.queue_status == JobQueueStatus.QUEUED.value
and job.send_status == JobSendStatus.QUEUED.value
):
return "already_queued"
if job.send_status in DELIVERY_ACCEPTED_STATUSES:
raise QueueingError(
"This message has already been accepted by a delivery channel and "
"cannot be sent again from preview."
)
if job.send_status in {
JobSendStatus.CLAIMED.value,
JobSendStatus.SENDING.value,
JobSendStatus.OUTCOME_UNKNOWN.value,
}:
raise QueueingError(
f"This message is in delivery state {job.send_status}; reconcile or wait before sending it again."
)
if job.send_status in {
JobSendStatus.FAILED_TEMPORARY.value,
JobSendStatus.FAILED_PERMANENT.value,
}:
raise QueueingError(
"This message has failed before. Use the explicit retry action so the retry is visible in the delivery protocol."
)
if job.attempt_count > 0 or job.postbox_attempt_count > 0:
raise QueueingError(
"This message already has delivery attempts. Use retry or "
"reconciliation instead of preview send."
)
if job.build_status != JobBuildStatus.BUILT.value:
raise QueueingError("This message has not been built yet.")
if not _single_job_validation_allowed(
version, job, include_warnings=include_warnings
):
raise QueueingError(
f"This message cannot be sent while validation status is {job.validation_status}."
)
if not job.eml_local_path and not job.eml_storage_key:
raise QueueingError(
"This message has no generated EML evidence. Rebuild the campaign before sending."
)
if job.queue_status not in {
JobQueueStatus.DRAFT.value,
JobQueueStatus.CANCELLED.value,
}:
raise QueueingError(
f"This message cannot be sent from queue state {job.queue_status}."
)
if dry_run:
return "would_queue"
job.queue_status = JobQueueStatus.QUEUED.value
job.send_status = JobSendStatus.QUEUED.value
job.queued_at = _utcnow()
job.claimed_at = None
job.claim_token = None
job.smtp_started_at = None
job.outcome_unknown_at = None
job.last_error = None
session.add(job)
return "queued"
def _single_job_validation_allowed(
version: CampaignVersion,
job: CampaignJob,
*,
include_warnings: bool,
) -> bool:
allowed = {JobValidationStatus.READY.value}
if include_warnings:
allowed.add(JobValidationStatus.WARNING.value)
if job.validation_status in allowed:
return True
return (
job.validation_status == JobValidationStatus.NEEDS_REVIEW.value
and _job_review_key(job) in _reviewed_needs_review_keys(version)
)
def reconcile_job_outcome(
session: Session,
*,
tenant_id: str,
campaign_id: str,
job_id: str,
decision: str,
note: str | None = None,
attempt_id: str | None = None,
commit: bool = True,
) -> dict[str, Any]:
"""Record an operator decision for an uncertain/incomplete worker state.
``smtp_accepted`` protects the message from retry and enables IMAP append.
``not_sent`` converts it into a failed-temporary candidate; retry remains a
separate explicit action. ``imap_appended`` completes an uncertain append;
only the evidence-backed ``imap_not_appended`` decision makes it retryable.
"""
campaign = _get_campaign_for_tenant(
session, campaign_id=campaign_id, tenant_id=tenant_id
)
job = session.get(CampaignJob, job_id)
if not job or job.tenant_id != tenant_id or job.campaign_id != campaign.id:
raise QueueingError("Campaign job not found or not accessible")
if decision in {"imap_appended", "imap_not_appended"}:
return _reconcile_imap_append_outcome(
session,
campaign=campaign,
job=job,
decision=decision,
note=note,
commit=commit,
)
if decision in {"postbox_accepted", "postbox_not_accepted"}:
return _reconcile_postbox_outcome(
session,
campaign=campaign,
job=job,
decision=decision,
note=note,
attempt_id=attempt_id,
commit=commit,
)
evidence_note = (note or "").strip()
if not evidence_note:
raise QueueingError("SMTP reconciliation requires an evidence note")
if job.send_status != JobSendStatus.OUTCOME_UNKNOWN.value:
raise QueueingError(
f"Job status {job.send_status} does not require reconciliation"
)
version = _get_version_for_campaign(
session, campaign, version_id=job.campaign_version_id
)
snapshot = ensure_execution_snapshot(session, version)
now = _utcnow()
attempt = _unfinished_attempt(session, job)
if decision == "smtp_accepted":
job.send_status = JobSendStatus.SMTP_ACCEPTED.value
job.queue_status = JobQueueStatus.DRAFT.value
job.sent_at = job.sent_at or now
job.outcome_unknown_at = None
job.claim_token = None
job.last_error = evidence_note
job.imap_status = (
JobImapStatus.PENDING.value
if snapshot.delivery.imap_append_sent.enabled
else JobImapStatus.NOT_REQUESTED.value
)
files_integration().mark_job_attachment_uses_sent(session, job)
attempt_status = "reconciled_smtp_accepted"
elif decision == "not_sent":
job.send_status = JobSendStatus.FAILED_TEMPORARY.value
job.queue_status = JobQueueStatus.DRAFT.value
job.outcome_unknown_at = None
job.claim_token = None
job.last_error = evidence_note
attempt_status = "reconciled_not_sent"
else:
raise QueueingError("decision must be 'smtp_accepted' or 'not_sent'")
if attempt:
attempt.status = attempt_status
attempt.finished_at = now
attempt.error_type = "OperatorReconciliation"
attempt.error_message = job.last_error
session.add(attempt)
session.add(job)
_update_campaign_after_job(session, campaign.id, version.id)
if commit:
session.commit()
else:
session.flush()
return {
"campaign_id": campaign.id,
"version_id": version.id,
"job_id": job.id,
"decision": decision,
"send_status": job.send_status,
"imap_status": job.imap_status,
"note": job.last_error,
}
def _postbox_outcome_from_attempts(
session: Session,
job: CampaignJob,
) -> PostboxChannelOutcome:
attempts = (
session.query(PostboxDeliveryAttempt)
.filter(PostboxDeliveryAttempt.job_id == job.id)
.order_by(
PostboxDeliveryAttempt.target_index.asc(),
PostboxDeliveryAttempt.attempt_number.desc(),
)
.all()
)
latest_by_target: dict[str, PostboxDeliveryAttempt] = {}
for attempt in attempts:
latest_by_target.setdefault(attempt.target_key, attempt)
outcome = PostboxChannelOutcome()
for attempt in latest_by_target.values():
if attempt.status == JobPostboxStatus.ACCEPTED.value:
outcome.accepted += 1
elif attempt.status == JobPostboxStatus.ACCEPTED_VACANT.value:
outcome.accepted_vacant += 1
elif attempt.status == JobPostboxStatus.REJECTED_TEMPORARY.value:
outcome.rejected_temporary += 1
elif attempt.status == JobPostboxStatus.REJECTED_PERMANENT.value:
outcome.rejected_permanent += 1
elif attempt.status in {
JobPostboxStatus.OUTCOME_UNKNOWN.value,
JobPostboxStatus.DELIVERING.value,
}:
outcome.outcome_unknown += 1
return outcome
def _reconcile_postbox_outcome(
session: Session,
*,
campaign: Campaign,
job: CampaignJob,
decision: str,
note: str | None,
attempt_id: str | None,
commit: bool,
) -> dict[str, Any]:
evidence_note = (note or "").strip()
if not evidence_note:
raise QueueingError("Postbox reconciliation requires an evidence note")
unknown_query = session.query(PostboxDeliveryAttempt).filter(
PostboxDeliveryAttempt.job_id == job.id,
PostboxDeliveryAttempt.status.in_(
[
JobPostboxStatus.OUTCOME_UNKNOWN.value,
JobPostboxStatus.DELIVERING.value,
]
),
)
if attempt_id:
unknown_query = unknown_query.filter(PostboxDeliveryAttempt.id == attempt_id)
unknown_attempts = unknown_query.order_by(
PostboxDeliveryAttempt.target_index.asc(),
PostboxDeliveryAttempt.attempt_number.desc(),
).all()
if not unknown_attempts:
raise QueueingError(
"No unresolved Postbox attempt matches this reconciliation."
)
if not attempt_id and len(unknown_attempts) > 1:
raise QueueingError(
"More than one Postbox attempt is unresolved; select a specific "
"attempt before reconciling."
)
attempt = unknown_attempts[0]
evidence = dict(attempt.evidence or {})
evidence["operator_reconciliation"] = {
"decision": decision,
"note": evidence_note,
"at": _utcnow().isoformat(),
}
attempt.evidence = evidence
attempt.error_type = "OperatorReconciliation"
attempt.error_message = evidence_note
attempt.finished_at = attempt.finished_at or _utcnow()
attempt.status = (
JobPostboxStatus.ACCEPTED.value
if decision == "postbox_accepted"
else JobPostboxStatus.REJECTED_TEMPORARY.value
)
session.add(attempt)
session.flush()
postbox_outcome = _postbox_outcome_from_attempts(session, job)
postbox_outcome.messages.append(evidence_note)
job.postbox_status = postbox_outcome.status
mail_outcome = (
_MailChannelOutcome(accepted=True)
if _mail_was_accepted(
session,
job.id,
send_status=job.send_status,
)
else None
)
result = _finalize_multichannel_job(
session,
job_id=job.id,
channel_policy=DeliveryChannelPolicy(job.delivery_channel_policy),
mail=mail_outcome,
postbox=postbox_outcome,
commit=commit,
)
return {
"campaign_id": campaign.id,
"version_id": job.campaign_version_id,
"job_id": job.id,
"channel": "postbox",
"attempt_id": attempt.id,
"decision": decision,
"send_status": result.status,
"postbox_status": job.postbox_status,
"note": evidence_note,
}
def _reconcile_imap_append_outcome(
session: Session,
*,
campaign: Campaign,
job: CampaignJob,
decision: str,
note: str | None,
commit: bool,
) -> dict[str, Any]:
"""Resolve an uncertain append while retaining its immutable attempt row."""
evidence_note = (note or "").strip()
if not evidence_note:
raise QueueingError("IMAP reconciliation requires an evidence note")
if not _mail_was_accepted(
session,
job.id,
send_status=job.send_status,
):
raise QueueingError(
"Only an SMTP-accepted job can reconcile a Sent-folder append"
)
if job.imap_status != JobImapStatus.OUTCOME_UNKNOWN.value:
raise QueueingError(
f"IMAP status {job.imap_status} does not require reconciliation"
)
attempt = (
session.query(ImapAppendAttempt)
.filter(ImapAppendAttempt.job_id == job.id)
.order_by(ImapAppendAttempt.attempt_number.desc())
.first()
)
if decision == "imap_appended":
job.imap_status = JobImapStatus.APPENDED.value
attempt_status = "reconciled_imap_appended"
files_integration().mark_job_attachment_uses_sent(session, job)
elif decision == "imap_not_appended":
# FAILED is deliberately retryable only after this explicit,
# evidence-backed operator decision.
job.imap_status = JobImapStatus.FAILED.value
attempt_status = "reconciled_imap_not_appended"
else: # Kept defensive for non-HTTP callers.
raise QueueingError(
"IMAP decision must be 'imap_appended' or 'imap_not_appended'"
)
job.imap_claimed_at = None
job.imap_claim_token = None
job.last_error = evidence_note
if attempt is not None:
attempt.status = attempt_status
attempt.error_message = evidence_note
session.add(attempt)
session.add(job)
if commit:
session.commit()
else:
session.flush()
return {
"campaign_id": campaign.id,
"version_id": job.campaign_version_id,
"job_id": job.id,
"channel": "imap",
"decision": decision,
"send_status": job.send_status,
"imap_status": job.imap_status,
"note": evidence_note,
}
def _verify_eml_evidence(job: CampaignJob, payload: bytes) -> None:
if job.eml_size_bytes is not None and len(payload) != job.eml_size_bytes:
raise SendJobError(
f"Generated EML size mismatch for job {job.id}: expected {job.eml_size_bytes}, got {len(payload)}"
)
if job.eml_sha256:
actual = hashlib.sha256(payload).hexdigest()
if actual != job.eml_sha256:
raise SendJobError(
f"Generated EML SHA-256 mismatch for job {job.id}; rebuild the campaign version before delivery"
)
if job.message_id_header:
parsed = BytesParser(policy=policy.default).parsebytes(payload)
message_id = parsed.get("Message-ID")
parsed_message_id = str(message_id) if message_id else None
if parsed_message_id != job.message_id_header:
raise SendJobError(
f"Generated EML Message-ID mismatch for job {job.id}; rebuild the campaign version before delivery"
)
def _load_eml_bytes_for_job(job: CampaignJob) -> bytes:
if job.eml_storage_key:
try:
payload = configured_storage_backend(
get_settings() or core_settings
).get_bytes(job.eml_storage_key)
except StorageBackendError as exc:
raise SendJobError(
f"Generated EML object could not be read for job {job.id}: {exc}"
) from exc
_verify_eml_evidence(job, payload)
return payload
if job.eml_local_path:
path = Path(job.eml_local_path)
if not path.exists():
raise SendJobError(f"Generated EML file does not exist: {path}")
payload = path.read_bytes()
_verify_eml_evidence(job, payload)
return payload
raise SendJobError("Generated EML evidence is not available for this job")
def _load_eml_for_job(job: CampaignJob):
return BytesParser(policy=policy.default).parsebytes(_load_eml_bytes_for_job(job))
def _addresses_from_job(job: CampaignJob, field: str) -> list[str]:
data = job.resolved_recipients or {}
values = data.get(field) or []
return [
item.get("email")
for item in values
if isinstance(item, dict) and item.get("email")
]
def _sender_from_job(job: CampaignJob) -> str:
data = job.resolved_recipients or {}
bounce_to = _addresses_from_job(job, "bounce_to")
if bounce_to:
return bounce_to[0]
from_data = data.get("from") if isinstance(data, dict) else None
if isinstance(from_data, dict) and from_data.get("email"):
return from_data["email"]
raise SmtpConfigurationError(
"No envelope sender could be determined from Campaign-owned recipient data; configure recipients.from or bounce_to before building"
)
def _from_header_from_job(job: CampaignJob) -> str | None:
data = job.resolved_recipients or {}
from_data = data.get("from") if isinstance(data, dict) else None
if isinstance(from_data, dict) and from_data.get("email"):
return str(from_data["email"])
return None
def _recipients_from_job(job: CampaignJob) -> list[str]:
recipients: list[str] = []
for field in ["to", "cc", "bcc"]:
recipients.extend(_addresses_from_job(job, field))
# Preserve order while de-duplicating.
return list(dict.fromkeys(recipients))
def _claim_job_for_sending(session: Session, job: CampaignJob) -> str | None:
"""Atomically claim a queued job and return the claim token.
A duplicate task can observe CLAIMED/SENDING but cannot acquire a second
claim. CLAIMED is deliberately not timed out automatically: an operator can
safely release it after checking worker state, while an automatic timeout
could race a slow worker.
"""
claim_token = str(uuid4())
changed = (
session.query(CampaignJob)
.filter(
CampaignJob.id == job.id,
CampaignJob.queue_status == JobQueueStatus.QUEUED.value,
CampaignJob.send_status == JobSendStatus.QUEUED.value,
)
.update(
{
CampaignJob.queue_status: JobQueueStatus.SENDING.value,
CampaignJob.send_status: JobSendStatus.CLAIMED.value,
CampaignJob.claimed_at: _utcnow(),
CampaignJob.claim_token: claim_token,
CampaignJob.last_error: None,
},
synchronize_session=False,
)
)
session.commit()
session.expire_all()
return claim_token if changed == 1 else None
def _record_attempt_start(
session: Session, job: CampaignJob, claim_token: str
) -> SendAttempt:
now = _utcnow()
attempt = SendAttempt(
job_id=job.id,
attempt_number=job.attempt_count + 1,
status="smtp_in_progress",
claim_token=claim_token,
started_at=now,
)
job.attempt_count += 1
job.queue_status = JobQueueStatus.SENDING.value
job.send_status = JobSendStatus.SENDING.value
job.smtp_started_at = now
job.outcome_unknown_at = None
job.last_error = None
session.add(attempt)
session.add(job)
session.commit()
return attempt
def _unfinished_attempt(session: Session, job: CampaignJob) -> SendAttempt | None:
return (
session.query(SendAttempt)
.filter(SendAttempt.job_id == job.id, SendAttempt.finished_at.is_(None))
.order_by(SendAttempt.attempt_number.desc())
.first()
)
def mark_job_outcome_unknown(
session: Session, job: CampaignJob, *, reason: str
) -> SendJobResult:
"""Freeze an uncertain SMTP result and prohibit automatic redelivery."""
now = _utcnow()
attempt = _unfinished_attempt(session, job)
if attempt:
attempt.status = "outcome_unknown"
attempt.finished_at = now
attempt.error_type = "OutcomeUnknown"
attempt.error_message = reason
session.add(attempt)
job.queue_status = JobQueueStatus.DRAFT.value
job.send_status = JobSendStatus.OUTCOME_UNKNOWN.value
job.outcome_unknown_at = now
job.last_error = reason
session.add(job)
_update_campaign_after_job(session, job.campaign_id, job.campaign_version_id)
session.commit()
return SendJobResult(
job_id=job.id,
status=JobSendStatus.OUTCOME_UNKNOWN.value,
attempt_number=job.attempt_count,
message=reason,
)
def _update_campaign_after_job(
session: Session, campaign_id: str, version_id: str | None = None
) -> None:
session.flush()
campaign = session.get(Campaign, campaign_id)
if not campaign:
return
previous_status = campaign.status
version = session.get(CampaignVersion, version_id) if version_id else None
counts = _campaign_delivery_counts(
session, campaign_id=campaign_id, version_id=version_id
)
_apply_campaign_delivery_status(campaign, version, counts)
if version and (counts.accepted or counts.unknown):
if version.locked_at is None:
version.locked_at = _utcnow()
session.add(version)
session.add(campaign)
if campaign.status != previous_status:
_emit_campaign_status_notification(
session,
campaign=campaign,
status=campaign.status,
previous_status=previous_status,
version_id=version_id,
)
def _campaign_delivery_counts(
session: Session, *, campaign_id: str, version_id: str | None
) -> _CampaignDeliveryCounts:
base_filters = _campaign_delivery_base_filters(
campaign_id=campaign_id, version_id=version_id
)
counts = {
status: session.query(CampaignJob)
.filter(*base_filters, CampaignJob.send_status == status)
.count()
for status in {item.value for item in JobSendStatus}
}
return _CampaignDeliveryCounts(
accepted=sum(counts.get(status, 0) for status in DELIVERY_ACCEPTED_STATUSES),
partial=counts.get(JobSendStatus.PARTIALLY_ACCEPTED.value, 0),
unknown=counts.get(JobSendStatus.OUTCOME_UNKNOWN.value, 0),
failed=counts.get(JobSendStatus.FAILED_TEMPORARY.value, 0)
+ counts.get(JobSendStatus.FAILED_PERMANENT.value, 0),
active=(
counts.get(JobSendStatus.QUEUED.value, 0)
+ counts.get(JobSendStatus.CLAIMED.value, 0)
+ counts.get(JobSendStatus.SENDING.value, 0)
),
cancelled=counts.get(JobSendStatus.CANCELLED.value, 0),
not_started=_queueable_not_started_job_count(session, base_filters),
)
def _campaign_delivery_base_filters(
*, campaign_id: str, version_id: str | None
) -> list[object]:
base_filters: list[object] = [CampaignJob.campaign_id == campaign_id]
if version_id:
base_filters.append(CampaignJob.campaign_version_id == version_id)
return base_filters
def _queueable_not_started_job_count(
session: Session, base_filters: list[object]
) -> int:
return (
session.query(CampaignJob)
.filter(
*base_filters,
CampaignJob.send_status == JobSendStatus.NOT_QUEUED.value,
CampaignJob.build_status == JobBuildStatus.BUILT.value,
CampaignJob.validation_status.in_(list(QUEUEABLE_VALIDATION_STATUSES)),
)
.count()
)
def _apply_campaign_delivery_status(
campaign: Campaign,
version: CampaignVersion | None,
counts: _CampaignDeliveryCounts,
) -> None:
status_pair = _campaign_delivery_status_pair(counts)
if status_pair is None:
return
campaign_status, workflow_state = status_pair
campaign.status = campaign_status
if version:
version.workflow_state = workflow_state
def _campaign_delivery_status_pair(
counts: _CampaignDeliveryCounts,
) -> tuple[str, str] | None:
if counts.active:
return CampaignStatus.SENDING.value, CampaignVersionWorkflowState.SENDING.value
if counts.accepted and (
counts.partial
or counts.failed
or counts.unknown
or counts.cancelled
or counts.not_started
):
return (
CampaignStatus.PARTIALLY_COMPLETED.value,
CampaignVersionWorkflowState.PARTIALLY_COMPLETED.value,
)
if counts.unknown:
return (
CampaignStatus.OUTCOME_UNKNOWN.value,
CampaignVersionWorkflowState.OUTCOME_UNKNOWN.value,
)
if counts.failed:
return CampaignStatus.FAILED.value, CampaignVersionWorkflowState.FAILED.value
if counts.accepted:
return CampaignStatus.SENT.value, CampaignVersionWorkflowState.COMPLETED.value
if counts.cancelled:
return (
CampaignStatus.CANCELLED.value,
CampaignVersionWorkflowState.CANCELLED.value,
)
return None
def send_campaign_job(
session: Session,
*,
job_id: str,
dry_run: bool = False,
use_rate_limit: bool = True,
enqueue_imap_task: bool = False,
) -> SendJobResult:
job = session.get(CampaignJob, job_id)
if not job:
raise SendJobError(f"Job not found: {job_id}")
preflight_result = _preflight_send_campaign_job(session, job, dry_run=dry_run)
if preflight_result is not None:
return preflight_result
context = _send_job_delivery_context(session, job)
if dry_run:
policy = DeliveryChannelPolicy(job.delivery_channel_policy)
descriptions: list[str] = []
if policy.uses_mail:
descriptions.append(
f"Mail to {len(context.envelope_recipients)} recipient(s)"
)
if policy.uses_postbox:
descriptions.append(
f"Postbox to {len(job.resolved_postbox_targets or [])} target(s)"
)
return SendJobResult(
job_id=job.id,
status="dry_run",
attempt_number=job.attempt_count,
dry_run=True,
message=f"Would deliver via {'; '.join(descriptions)}",
)
claimed = _claimed_campaign_job_for_delivery(session, job)
if isinstance(claimed, SendJobResult):
return claimed
claimed_job, claim_token = claimed
return _send_claimed_campaign_job(
session,
job=claimed_job,
claim_token=claim_token,
context=context,
use_rate_limit=use_rate_limit,
enqueue_imap_task=enqueue_imap_task,
)
def _preflight_send_campaign_job(
session: Session,
job: CampaignJob,
*,
dry_run: bool,
) -> SendJobResult | None:
if (
job.queue_status == JobQueueStatus.CANCELLED.value
or job.send_status == JobSendStatus.CANCELLED.value
):
return SendJobResult(
job_id=job.id,
status="cancelled",
attempt_number=job.attempt_count,
dry_run=dry_run,
)
if job.queue_status == JobQueueStatus.PAUSED.value:
return SendJobResult(
job_id=job.id,
status="paused",
attempt_number=job.attempt_count,
dry_run=dry_run,
)
if job.send_status in FULLY_ACCEPTED_STATUSES:
return SendJobResult(
job_id=job.id,
status="already_accepted",
attempt_number=job.attempt_count,
dry_run=dry_run,
)
if job.send_status == JobSendStatus.OUTCOME_UNKNOWN.value:
return SendJobResult(
job_id=job.id,
status=JobSendStatus.OUTCOME_UNKNOWN.value,
attempt_number=job.attempt_count,
dry_run=dry_run,
message="A delivery outcome is unresolved; reconcile it before any retry.",
)
if job.send_status == JobSendStatus.SENDING.value:
return mark_job_outcome_unknown(
session,
job,
reason="A delivery task resumed while the previous channel attempt was still marked in progress. Automatic redelivery was stopped.",
)
if job.send_status == JobSendStatus.CLAIMED.value:
return SendJobResult(
job_id=job.id,
status="already_claimed",
attempt_number=job.attempt_count,
dry_run=dry_run,
message="Another worker has claimed this job, or a stale pre-SMTP claim requires operator review.",
)
if (
job.queue_status != JobQueueStatus.QUEUED.value
or job.send_status != JobSendStatus.QUEUED.value
):
raise SendJobError(
f"Job is not explicitly queued for delivery: queue={job.queue_status}, send={job.send_status}"
)
return None
def _send_job_delivery_context(
session: Session, job: CampaignJob
) -> _SendJobDeliveryContext:
version = session.get(CampaignVersion, job.campaign_version_id)
if not version:
raise SendJobError("Campaign version not found")
try:
snapshot = ensure_execution_snapshot(session, version, effect_job=job)
except ExecutionSnapshotError as exc:
raise SendJobError(str(exc)) from exc
message_bytes = _load_eml_bytes_for_job(job)
channel_policy = DeliveryChannelPolicy(
getattr(job, "delivery_channel_policy", DeliveryChannelPolicy.MAIL.value)
)
envelope_from: str | None = None
envelope_recipients: list[str] = []
if channel_policy.uses_mail:
envelope_from = _sender_from_job(job)
envelope_recipients = _recipients_from_job(job)
if not envelope_recipients:
raise SmtpConfigurationError("No envelope recipients could be determined")
return _SendJobDeliveryContext(
version=version,
snapshot=snapshot,
message_bytes=message_bytes,
envelope_from=envelope_from,
envelope_recipients=envelope_recipients,
)
def _claimed_campaign_job_for_delivery(
session: Session,
job: CampaignJob,
) -> tuple[CampaignJob, str] | SendJobResult:
claim_token = _claim_job_for_sending(session, job)
if claim_token is None:
return _not_claimed_send_job_result(session, job)
job = session.get(CampaignJob, job.id)
if job is None:
raise SendJobError("Claimed campaign job disappeared before send.")
return job, claim_token
def _not_claimed_send_job_result(session: Session, job: CampaignJob) -> SendJobResult:
current = session.get(CampaignJob, job.id)
if not current:
raise SendJobError(f"Job disappeared while claiming: {job.id}")
if current.send_status == JobSendStatus.SENDING.value:
return mark_job_outcome_unknown(
session,
current,
reason="A duplicate/redelivered task found an unfinished SMTP attempt. Automatic resend was stopped.",
)
return SendJobResult(
job_id=current.id,
status="not_claimed",
attempt_number=current.attempt_count,
message=f"Job is no longer queueable: {current.send_status}",
)
def _send_claimed_campaign_job(
session: Session,
*,
job: CampaignJob,
claim_token: str,
context: _SendJobDeliveryContext,
use_rate_limit: bool,
enqueue_imap_task: bool,
) -> SendJobResult:
channel_policy = DeliveryChannelPolicy(
getattr(job, "delivery_channel_policy", DeliveryChannelPolicy.MAIL.value)
)
if channel_policy == DeliveryChannelPolicy.MAIL:
return _send_claimed_mail_only_job(
session,
job=job,
claim_token=claim_token,
context=context,
use_rate_limit=use_rate_limit,
enqueue_imap_task=enqueue_imap_task,
)
return _send_claimed_multichannel_job(
session,
job=job,
claim_token=claim_token,
context=context,
channel_policy=channel_policy,
use_rate_limit=use_rate_limit,
enqueue_imap_task=enqueue_imap_task,
)
def _send_claimed_mail_only_job(
session: Session,
*,
job: CampaignJob,
claim_token: str,
context: _SendJobDeliveryContext,
use_rate_limit: bool,
enqueue_imap_task: bool,
) -> SendJobResult:
if context.envelope_from is None or not context.envelope_recipients:
raise SmtpConfigurationError(
"Mail delivery has no frozen envelope sender or recipients."
)
mail_integration().wait_for_rate_limit(
key=f"tenant:{job.tenant_id}:campaign:{job.campaign_id}",
messages_per_minute=context.snapshot.delivery.rate_limit.messages_per_minute,
enabled=use_rate_limit,
)
attempt = _record_attempt_start(session, job, claim_token)
try:
result = mail_integration().send_campaign_email_bytes(
session,
tenant_id=job.tenant_id,
campaign_id=job.campaign_id,
profile_id=context.snapshot.mail_profile_id,
message_bytes=context.message_bytes,
envelope_from=context.envelope_from,
envelope_recipients=context.envelope_recipients,
from_header=_from_header_from_job(job),
expected_smtp_transport_revision=context.snapshot.smtp_transport_revision
or "",
smtp_server_id=context.snapshot.smtp_server_id,
smtp_credential_id=context.snapshot.smtp_credential_id,
)
if result.accepted_count <= 0:
raise SmtpSendError(
"SMTP did not accept any envelope recipients", temporary=False
)
except SmtpSendError as exc:
outcome_unknown = _record_smtp_send_error(
session, job=job, attempt=attempt, exc=exc
)
if outcome_unknown is not None:
return outcome_unknown
raise
except (MailProfileError, SmtpConfigurationError, SendJobError, OSError) as exc:
_record_permanent_send_error(session, job=job, attempt=attempt, exc=exc)
raise
try:
success = _record_smtp_send_success(
session,
job=job,
attempt=attempt,
snapshot=context.snapshot,
result=result,
)
except Exception:
session.rollback()
current = session.get(CampaignJob, job.id)
if current is None:
raise SendJobError(
"SMTP accepted the message, but its durable Campaign job record is unavailable."
) from None
return mark_job_outcome_unknown(
session,
current,
reason=(
"SMTP accepted the provider effect, but persisting the accepted outcome failed. "
"Automatic retry is stopped until an operator reconciles the job."
),
)
if (
enqueue_imap_task
and _celery_enabled()
and job.imap_status == JobImapStatus.PENDING.value
):
try:
_celery_enqueue_append_sent_job(job.id)
except Exception:
return SendJobResult(
job_id=success.job_id,
status=success.status,
attempt_number=success.attempt_number,
message=(
f"{success.message + ' ' if success.message else ''}"
"SMTP acceptance was persisted, but the Sent-folder append could not be enqueued; it remains pending."
),
)
return success
def _mail_was_accepted(
session: Session,
job_id: str,
*,
send_status: str | None = None,
) -> bool:
if send_status in {
JobSendStatus.SMTP_ACCEPTED.value,
JobSendStatus.SENT.value,
JobSendStatus.DELIVERED.value,
}:
return True
return (
session.query(SendAttempt.id)
.filter(
SendAttempt.job_id == job_id,
SendAttempt.status.in_(
[
JobSendStatus.SMTP_ACCEPTED.value,
"smtp_accepted_with_refusals",
"reconciled_smtp_accepted",
]
),
)
.first()
is not None
)
def _deliver_mail_channel(
session: Session,
*,
job: CampaignJob,
claim_token: str,
context: _SendJobDeliveryContext,
use_rate_limit: bool,
enqueue_imap_task: bool,
) -> _MailChannelOutcome:
if _mail_was_accepted(
session,
job.id,
send_status=job.send_status,
):
return _MailChannelOutcome(accepted=True)
try:
result = _send_claimed_mail_only_job(
session,
job=job,
claim_token=claim_token,
context=context,
use_rate_limit=use_rate_limit,
enqueue_imap_task=enqueue_imap_task,
)
except Exception as exc:
current = session.get(CampaignJob, job.id)
if current is None:
raise SendJobError(
"Campaign job disappeared while recording Mail delivery."
) from exc
if current.send_status == JobSendStatus.OUTCOME_UNKNOWN.value:
return _MailChannelOutcome(
outcome_unknown=True,
message=current.last_error or str(exc),
)
if current.send_status == JobSendStatus.FAILED_TEMPORARY.value:
return _MailChannelOutcome(
rejected_temporary=True,
message=current.last_error or str(exc),
)
if current.send_status == JobSendStatus.FAILED_PERMANENT.value:
return _MailChannelOutcome(
rejected_permanent=True,
message=current.last_error or str(exc),
)
unknown = mark_job_outcome_unknown(
session,
current,
reason=(
"Mail delivery raised an unexpected error after its durable "
"attempt started. The outcome is unknown and no fallback was "
f"started: {exc}"
),
)
return _MailChannelOutcome(
outcome_unknown=True,
message=unknown.message,
)
return _MailChannelOutcome(
accepted=result.status
in {
JobSendStatus.SMTP_ACCEPTED.value,
"already_accepted",
},
outcome_unknown=result.status == JobSendStatus.OUTCOME_UNKNOWN.value,
message=result.message,
)
def _empty_postbox_outcome() -> PostboxChannelOutcome:
return PostboxChannelOutcome()
def _final_multichannel_status(
*,
channel_policy: DeliveryChannelPolicy,
mail: _MailChannelOutcome | None,
postbox: PostboxChannelOutcome | None,
) -> str:
outcome = _delivery_outcome_summary(mail=mail, postbox=postbox)
if outcome.outcome_unknown:
return JobSendStatus.OUTCOME_UNKNOWN.value
if not outcome.accepted_count:
return (
JobSendStatus.FAILED_TEMPORARY.value
if outcome.temporary_rejection
else JobSendStatus.FAILED_PERMANENT.value
)
classifier = _ACCEPTED_DELIVERY_CLASSIFIERS.get(
channel_policy, _classify_fallback_delivery
)
return classifier(outcome)
def _delivery_outcome_summary(
*,
mail: _MailChannelOutcome | None,
postbox: PostboxChannelOutcome | None,
) -> _DeliveryOutcomeSummary:
return _DeliveryOutcomeSummary(
mail_accepted=bool(mail and mail.accepted),
postbox_accepted=int(postbox.accepted_count if postbox else 0),
postbox_rejected=int(postbox.rejected_count if postbox else 0),
outcome_unknown=bool(
(mail and mail.outcome_unknown) or (postbox and postbox.outcome_unknown)
),
temporary_rejection=bool(
(mail and mail.rejected_temporary)
or (postbox and postbox.rejected_temporary)
),
)
def _classify_postbox_delivery(outcome: _DeliveryOutcomeSummary) -> str:
return (
JobSendStatus.PARTIALLY_ACCEPTED.value
if outcome.postbox_rejected
else JobSendStatus.POSTBOX_ACCEPTED.value
)
def _classify_dual_delivery(outcome: _DeliveryOutcomeSummary) -> str:
fully_delivered = (
outcome.mail_accepted
and outcome.postbox_accepted > 0
and not outcome.postbox_rejected
)
return (
JobSendStatus.DELIVERED.value
if fully_delivered
else JobSendStatus.PARTIALLY_ACCEPTED.value
)
def _classify_fallback_delivery(outcome: _DeliveryOutcomeSummary) -> str:
if outcome.postbox_rejected:
return JobSendStatus.PARTIALLY_ACCEPTED.value
return (
JobSendStatus.SMTP_ACCEPTED.value
if outcome.mail_accepted
else JobSendStatus.POSTBOX_ACCEPTED.value
)
_ACCEPTED_DELIVERY_CLASSIFIERS = {
DeliveryChannelPolicy.POSTBOX: _classify_postbox_delivery,
DeliveryChannelPolicy.MAIL_AND_POSTBOX: _classify_dual_delivery,
}
def _multichannel_messages(
mail: _MailChannelOutcome | None,
postbox: PostboxChannelOutcome | None,
) -> list[str]:
values: list[str] = []
if mail and mail.message:
values.append(f"Mail: {mail.message}")
if postbox:
values.extend(f"Postbox: {message}" for message in postbox.messages if message)
return values
def _finalize_multichannel_job(
session: Session,
*,
job_id: str,
channel_policy: DeliveryChannelPolicy,
mail: _MailChannelOutcome | None,
postbox: PostboxChannelOutcome | None,
commit: bool = True,
) -> SendJobResult:
job = session.get(CampaignJob, job_id)
if job is None:
raise SendJobError("Campaign job disappeared while finalizing delivery.")
status = _final_multichannel_status(
channel_policy=channel_policy,
mail=mail,
postbox=postbox,
)
accepted = status in DELIVERY_ACCEPTED_STATUSES
unknown = status == JobSendStatus.OUTCOME_UNKNOWN.value
messages = _multichannel_messages(mail, postbox)
job.queue_status = JobQueueStatus.DRAFT.value
job.send_status = status
job.claim_token = None
job.last_error = "\n".join(messages) or None
job.outcome_unknown_at = _utcnow() if unknown else None
if accepted or (
unknown
and bool((mail and mail.accepted) or (postbox and postbox.accepted_count))
):
job.sent_at = job.sent_at or _utcnow()
files_integration().mark_job_attachment_uses_sent(session, job)
if not (mail and mail.accepted):
job.imap_status = JobImapStatus.NOT_REQUESTED.value
session.add(job)
_update_campaign_after_job(
session,
job.campaign_id,
job.campaign_version_id,
)
if commit:
session.commit()
else:
session.flush()
return SendJobResult(
job_id=job.id,
status=status,
attempt_number=job.attempt_count + job.postbox_attempt_count,
message=job.last_error,
)
def _send_claimed_multichannel_job(
session: Session,
*,
job: CampaignJob,
claim_token: str,
context: _SendJobDeliveryContext,
channel_policy: DeliveryChannelPolicy,
use_rate_limit: bool,
enqueue_imap_task: bool,
) -> SendJobResult:
mail_outcome: _MailChannelOutcome | None = None
postbox_outcome: PostboxChannelOutcome | None = None
if channel_policy == DeliveryChannelPolicy.MAIL_THEN_POSTBOX:
prior_postbox_outcome = _postbox_outcome_from_attempts(session, job)
if prior_postbox_outcome.outcome_unknown:
postbox_outcome = prior_postbox_outcome
elif prior_postbox_outcome.accepted_count:
postbox_outcome = deliver_campaign_job_to_postboxes(
session,
job=job,
message_bytes=context.message_bytes,
classification=context.snapshot.delivery.postbox.classification,
)
else:
mail_outcome = _deliver_mail_channel(
session,
job=job,
claim_token=claim_token,
context=context,
use_rate_limit=use_rate_limit,
enqueue_imap_task=enqueue_imap_task,
)
if mail_outcome.rejected_before_acceptance:
current = session.get(CampaignJob, job.id)
if current is None:
raise SendJobError(
"Campaign job disappeared before Postbox fallback."
)
postbox_outcome = deliver_campaign_job_to_postboxes(
session,
job=current,
message_bytes=context.message_bytes,
classification=context.snapshot.delivery.postbox.classification,
)
else:
current = session.get(CampaignJob, job.id)
if current is not None:
current.postbox_status = JobPostboxStatus.SKIPPED.value
session.add(current)
session.commit()
else:
postbox_outcome = deliver_campaign_job_to_postboxes(
session,
job=job,
message_bytes=context.message_bytes,
classification=context.snapshot.delivery.postbox.classification,
)
should_deliver_mail = (
channel_policy == DeliveryChannelPolicy.MAIL_AND_POSTBOX
or (
channel_policy == DeliveryChannelPolicy.POSTBOX_THEN_MAIL
and postbox_outcome.all_rejected_before_acceptance
)
)
if should_deliver_mail:
current = session.get(CampaignJob, job.id)
if current is None:
raise SendJobError("Campaign job disappeared before Mail delivery.")
mail_outcome = _deliver_mail_channel(
session,
job=current,
claim_token=claim_token,
context=context,
use_rate_limit=use_rate_limit,
enqueue_imap_task=enqueue_imap_task,
)
return _finalize_multichannel_job(
session,
job_id=job.id,
channel_policy=channel_policy,
mail=mail_outcome,
postbox=postbox_outcome,
)
def _record_smtp_send_success(
session: Session,
*,
job: CampaignJob,
attempt: SendAttempt,
snapshot: ExecutionSnapshot,
result: object,
) -> SendJobResult:
refused_warning = _refused_recipient_warning(result)
attempt.finished_at = _utcnow()
attempt.status = (
"smtp_accepted_with_refusals"
if refused_warning
else JobSendStatus.SMTP_ACCEPTED.value
)
attempt.smtp_response = json.dumps(asdict(result), default=str)
job.queue_status = JobQueueStatus.DRAFT.value
job.send_status = JobSendStatus.SMTP_ACCEPTED.value
job.sent_at = _utcnow()
job.claim_token = None
job.outcome_unknown_at = None
job.imap_status = (
JobImapStatus.PENDING.value
if snapshot.delivery.imap_append_sent.enabled
else JobImapStatus.NOT_REQUESTED.value
)
job.last_error = refused_warning
files_integration().mark_job_attachment_uses_sent(session, job)
session.add(attempt)
session.add(job)
_update_campaign_after_job(session, job.campaign_id, job.campaign_version_id)
session.commit()
return SendJobResult(
job_id=job.id,
status=JobSendStatus.SMTP_ACCEPTED.value,
attempt_number=attempt.attempt_number,
message=refused_warning,
)
def _refused_recipient_warning(result: object) -> str | None:
if not result.refused_recipients:
return None
return (
f"SMTP accepted {result.accepted_count}/{len(result.envelope_recipients)} envelope recipient(s); "
f"refused recipients: {json.dumps(result.refused_recipients, default=str, sort_keys=True)}"
)
def _record_smtp_send_error(
session: Session,
*,
job: CampaignJob,
attempt: SendAttempt,
exc: SmtpSendError,
) -> SendJobResult | None:
if getattr(exc, "outcome_unknown", False):
attempt.status = JobSendStatus.OUTCOME_UNKNOWN.value
attempt.finished_at = _utcnow()
attempt.error_type = exc.__class__.__name__
attempt.error_message = str(exc)
session.add(attempt)
job.claim_token = None
session.add(job)
session.flush()
return mark_job_outcome_unknown(session, job, reason=str(exc))
attempt.finished_at = _utcnow()
attempt.error_type = exc.__class__.__name__
attempt.error_message = str(exc)
retryable = bool(getattr(exc, "temporary", False))
job.last_error = str(exc)
job.queue_status = JobQueueStatus.DRAFT.value
job.send_status = (
JobSendStatus.FAILED_TEMPORARY.value
if retryable
else JobSendStatus.FAILED_PERMANENT.value
)
job.claim_token = None
attempt.status = job.send_status
session.add(attempt)
session.add(job)
_update_campaign_after_job(session, job.campaign_id, job.campaign_version_id)
session.commit()
return None
def _record_permanent_send_error(
session: Session,
*,
job: CampaignJob,
attempt: SendAttempt,
exc: Exception,
) -> None:
attempt.finished_at = _utcnow()
attempt.error_type = exc.__class__.__name__
attempt.error_message = str(exc)
attempt.status = JobSendStatus.FAILED_PERMANENT.value
job.last_error = str(exc)
job.queue_status = JobQueueStatus.DRAFT.value
job.send_status = JobSendStatus.FAILED_PERMANENT.value
job.claim_token = None
session.add(attempt)
session.add(job)
_update_campaign_after_job(session, job.campaign_id, job.campaign_version_id)
session.commit()
def _imap_attempt_count(session: Session, job_id: str) -> int:
return (
session.query(ImapAppendAttempt)
.filter(ImapAppendAttempt.job_id == job_id)
.count()
)
def _claim_job_for_imap_append(session: Session, job: CampaignJob) -> str | None:
"""Atomically grant one worker permission to invoke the IMAP provider."""
if not _mail_was_accepted(
session,
job.id,
send_status=getattr(
job,
"send_status",
JobSendStatus.SMTP_ACCEPTED.value,
),
):
return None
claim_token = str(uuid4())
changed = (
session.query(CampaignJob)
.filter(
CampaignJob.id == job.id,
CampaignJob.imap_status.in_(
[JobImapStatus.PENDING.value, JobImapStatus.FAILED.value]
),
CampaignJob.imap_claim_token.is_(None),
)
.update(
{
CampaignJob.imap_status: JobImapStatus.APPENDING.value,
CampaignJob.imap_claimed_at: _utcnow(),
CampaignJob.imap_claim_token: claim_token,
CampaignJob.last_error: None,
},
synchronize_session=False,
)
)
session.commit()
session.expire_all()
return claim_token if changed == 1 else None
def _record_imap_attempt_start(
session: Session,
job: CampaignJob,
claim_token: str,
) -> ImapAppendAttempt:
if (
job.imap_status != JobImapStatus.APPENDING.value
or job.imap_claim_token != claim_token
):
raise SendJobError(
"The IMAP append claim was lost before the provider effect started"
)
existing_count = (
session.query(ImapAppendAttempt)
.filter(ImapAppendAttempt.job_id == job.id)
.count()
)
attempt = ImapAppendAttempt(
job_id=job.id,
attempt_number=existing_count + 1,
status="running",
claim_token=claim_token,
)
job.last_error = None
session.add(attempt)
session.add(job)
session.commit()
return attempt
def _imap_blocked_result(
session: Session,
job: CampaignJob,
*,
dry_run: bool,
) -> AppendSentResult | None:
if job.imap_status not in {
JobImapStatus.NOT_REQUESTED.value,
JobImapStatus.APPENDED.value,
JobImapStatus.SKIPPED.value,
JobImapStatus.OUTCOME_UNKNOWN.value,
JobImapStatus.APPENDING.value,
}:
return None
attempts = _imap_attempt_count(session, job.id)
if job.imap_status == JobImapStatus.NOT_REQUESTED.value:
return AppendSentResult(
job_id=job.id,
status="not_requested",
attempt_number=attempts,
dry_run=dry_run,
)
if job.imap_status == JobImapStatus.APPENDED.value:
return AppendSentResult(
job_id=job.id,
status="already_appended",
attempt_number=attempts,
dry_run=dry_run,
)
if job.imap_status == JobImapStatus.SKIPPED.value:
return AppendSentResult(
job_id=job.id, status="skipped", attempt_number=attempts, dry_run=dry_run
)
if job.imap_status == JobImapStatus.OUTCOME_UNKNOWN.value:
return AppendSentResult(
job_id=job.id,
status=JobImapStatus.OUTCOME_UNKNOWN.value,
attempt_number=attempts,
dry_run=dry_run,
message="The prior IMAP append outcome is unresolved; inspect and reconcile the mailbox before retrying.",
)
if job.imap_status == JobImapStatus.APPENDING.value:
return AppendSentResult(
job_id=job.id,
status="append_in_progress",
attempt_number=attempts,
dry_run=dry_run,
message="Another worker owns the IMAP append claim; automatic duplicate append was stopped.",
)
return None
def _imap_result_after_lost_claim(
session: Session,
*,
job_id: str,
attempt: ImapAppendAttempt,
provider_succeeded: bool,
) -> AppendSentResult:
session.rollback()
current_attempt = session.get(ImapAppendAttempt, attempt.id)
if current_attempt is not None:
current_attempt.status = (
"late_ack_ignored" if provider_succeeded else "claim_lost"
)
current_attempt.error_message = "The provider returned after this attempt no longer owned the durable IMAP claim."
session.add(current_attempt)
if provider_succeeded:
# If no other worker is active, freeze any retryable state. A provider
# success received after claim loss must never be followed by an
# automatic duplicate append.
session.query(CampaignJob).filter(
CampaignJob.id == job_id,
CampaignJob.imap_status.in_(
[
JobImapStatus.PENDING.value,
JobImapStatus.FAILED.value,
JobImapStatus.OUTCOME_UNKNOWN.value,
]
),
CampaignJob.imap_claim_token.is_(None),
).update(
{
CampaignJob.imap_status: JobImapStatus.OUTCOME_UNKNOWN.value,
CampaignJob.last_error: (
"An IMAP provider acknowledgement arrived after the durable claim was lost; operator reconciliation is required."
),
},
synchronize_session=False,
)
session.commit()
session.expire_all()
current = session.get(CampaignJob, job_id)
if current is None:
raise SendJobError(
"Campaign job disappeared while reconciling a lost IMAP claim"
)
blocked = _imap_blocked_result(session, current, dry_run=False)
if blocked is not None:
return blocked
return AppendSentResult(
job_id=current.id,
status="claim_lost",
attempt_number=_imap_attempt_count(session, current.id),
message="The IMAP append claim changed while the provider operation was in progress.",
)
def _record_imap_append_failure(
session: Session,
*,
job: CampaignJob,
attempt: ImapAppendAttempt,
claim_token: str,
folder: str,
message: str,
outcome_unknown: bool,
) -> bool:
next_status = (
JobImapStatus.OUTCOME_UNKNOWN.value
if outcome_unknown
else JobImapStatus.FAILED.value
)
changed = (
session.query(CampaignJob)
.filter(
CampaignJob.id == job.id,
CampaignJob.imap_status == JobImapStatus.APPENDING.value,
CampaignJob.imap_claim_token == claim_token,
)
.update(
{
CampaignJob.imap_status: next_status,
CampaignJob.imap_claimed_at: None,
CampaignJob.imap_claim_token: None,
CampaignJob.last_error: message,
},
synchronize_session=False,
)
)
attempt.status = next_status if outcome_unknown else "failed"
attempt.folder = None if folder == "auto" else folder
attempt.error_message = message
session.add(attempt)
session.commit()
session.expire_all()
return changed == 1
def _record_imap_append_success(
session: Session,
*,
job: CampaignJob,
attempt: ImapAppendAttempt,
claim_token: str,
folder: str,
) -> AppendSentResult:
changed = (
session.query(CampaignJob)
.filter(
CampaignJob.id == job.id,
CampaignJob.imap_status == JobImapStatus.APPENDING.value,
CampaignJob.imap_claim_token == claim_token,
)
.update(
{
CampaignJob.imap_status: JobImapStatus.APPENDED.value,
CampaignJob.imap_claimed_at: None,
CampaignJob.imap_claim_token: None,
CampaignJob.last_error: None,
},
synchronize_session=False,
)
)
if changed != 1:
return _imap_result_after_lost_claim(
session,
job_id=job.id,
attempt=attempt,
provider_succeeded=True,
)
attempt.status = "appended"
attempt.folder = folder
attempt.error_message = None
session.add(attempt)
files_integration().mark_job_attachment_uses_sent(session, job)
session.commit()
session.expire_all()
return AppendSentResult(
job_id=job.id,
status="appended",
attempt_number=attempt.attempt_number,
folder=folder,
)
def _mark_imap_append_outcome_unknown_after_effect(
session: Session,
*,
job_id: str,
attempt_id: str,
claim_token: str,
reason: str,
) -> AppendSentResult:
"""Freeze retry after a provider effect whose local persistence failed."""
session.rollback()
current_attempt = session.get(ImapAppendAttempt, attempt_id)
if current_attempt is not None:
current_attempt.status = JobImapStatus.OUTCOME_UNKNOWN.value
current_attempt.error_message = reason
session.add(current_attempt)
changed = (
session.query(CampaignJob)
.filter(
CampaignJob.id == job_id,
CampaignJob.imap_status == JobImapStatus.APPENDING.value,
CampaignJob.imap_claim_token == claim_token,
)
.update(
{
CampaignJob.imap_status: JobImapStatus.OUTCOME_UNKNOWN.value,
CampaignJob.imap_claimed_at: None,
CampaignJob.imap_claim_token: None,
CampaignJob.last_error: reason,
},
synchronize_session=False,
)
)
if changed != 1:
if current_attempt is not None:
session.rollback()
current_attempt = session.get(ImapAppendAttempt, attempt_id)
if current_attempt is None:
raise SendJobError(
"IMAP append outcome is unknown and its attempt record is unavailable"
)
return _imap_result_after_lost_claim(
session,
job_id=job_id,
attempt=current_attempt,
provider_succeeded=True,
)
session.commit()
session.expire_all()
return AppendSentResult(
job_id=job_id,
status=JobImapStatus.OUTCOME_UNKNOWN.value,
attempt_number=_imap_attempt_count(session, job_id),
message=reason,
)
def _imap_append_precondition(
session: Session,
job: CampaignJob,
*,
dry_run: bool,
) -> AppendSentResult | None:
if not _mail_was_accepted(session, job.id, send_status=job.send_status):
return AppendSentResult(
job_id=job.id,
status="not_sent",
attempt_number=0,
dry_run=dry_run,
message="SMTP has not accepted this job",
)
return _imap_blocked_result(session, job, dry_run=dry_run)
def _prepare_imap_append(
session: Session,
job: CampaignJob,
*,
dry_run: bool,
) -> _ImapAppendContext | AppendSentResult:
version = session.get(CampaignVersion, job.campaign_version_id)
if not version:
raise SendJobError("Campaign version not found")
try:
snapshot = ensure_execution_snapshot(session, version, effect_job=job)
except ExecutionSnapshotError as exc:
raise SendJobError(str(exc)) from exc
if not snapshot.delivery.imap_append_sent.enabled:
job.imap_status = JobImapStatus.NOT_REQUESTED.value
session.add(job)
session.commit()
return AppendSentResult(
job_id=job.id, status="not_requested", attempt_number=0, dry_run=dry_run
)
if not snapshot.imap_transport_revision:
job.imap_status = JobImapStatus.SKIPPED.value
job.last_error = "IMAP append requested, but the selected Mail profile has no IMAP configuration"
session.add(job)
session.commit()
return AppendSentResult(
job_id=job.id,
status="skipped",
attempt_number=0,
dry_run=dry_run,
message=job.last_error,
)
return _ImapAppendContext(
snapshot=snapshot,
message_bytes=_load_eml_bytes_for_job(job),
folder=snapshot.delivery.imap_append_sent.folder or "auto",
)
def _imap_append_dry_run(
session: Session, job: CampaignJob, context: _ImapAppendContext
) -> AppendSentResult:
return AppendSentResult(
job_id=job.id,
status="dry_run",
attempt_number=_imap_attempt_count(session, job.id),
dry_run=True,
folder=context.folder,
message=f"Would append {len(context.message_bytes)} bytes to IMAP folder {context.folder!r}",
)
def _claim_imap_append(
session: Session, job: CampaignJob
) -> _ClaimedImapAppend | AppendSentResult:
claim_token = _claim_job_for_imap_append(session, job)
if claim_token is None:
current = session.get(CampaignJob, job.id)
if current is None:
raise SendJobError(f"Job disappeared while claiming IMAP append: {job.id}")
blocked = _imap_blocked_result(session, current, dry_run=False)
if blocked is not None:
return blocked
return AppendSentResult(
job_id=current.id,
status="not_claimed",
attempt_number=_imap_attempt_count(session, current.id),
message=f"Job is no longer eligible for IMAP append: {current.imap_status}",
)
claimed_job = session.get(CampaignJob, job.id)
if claimed_job is None:
raise SendJobError("Claimed campaign job disappeared before IMAP append")
return _ClaimedImapAppend(
job=claimed_job,
attempt=_record_imap_attempt_start(session, claimed_job, claim_token),
claim_token=claim_token,
)
def _record_imap_provider_error(
session: Session,
claimed: _ClaimedImapAppend,
*,
folder: str,
message: str,
outcome_unknown: bool,
) -> None:
owned = _record_imap_append_failure(
session,
job=claimed.job,
attempt=claimed.attempt,
claim_token=claimed.claim_token,
folder=folder,
message=message,
outcome_unknown=outcome_unknown,
)
if not owned:
_imap_result_after_lost_claim(
session,
job_id=claimed.job.id,
attempt=claimed.attempt,
provider_succeeded=outcome_unknown,
)
def _invoke_imap_append(
session: Session, claimed: _ClaimedImapAppend, context: _ImapAppendContext
):
snapshot = context.snapshot
job = claimed.job
return mail_integration().append_campaign_message_to_sent(
session,
tenant_id=job.tenant_id,
campaign_id=job.campaign_id,
profile_id=snapshot.mail_profile_id,
message_bytes=context.message_bytes,
folder=None if context.folder == "auto" else context.folder,
expected_smtp_transport_revision=snapshot.smtp_transport_revision or "",
expected_imap_transport_revision=snapshot.imap_transport_revision,
smtp_server_id=snapshot.smtp_server_id,
smtp_credential_id=snapshot.smtp_credential_id,
imap_server_id=snapshot.imap_server_id,
imap_credential_id=snapshot.imap_credential_id,
)
def _perform_imap_append(
session: Session,
claimed: _ClaimedImapAppend,
context: _ImapAppendContext,
) -> AppendSentResult:
try:
result = _invoke_imap_append(session, claimed, context)
except (MailProfileError, ImapConfigurationError, ImapAppendError) as exc:
_record_imap_provider_error(
session,
claimed,
folder=context.folder,
message=str(exc),
outcome_unknown=bool(getattr(exc, "outcome_unknown", False)),
)
raise
except Exception:
reason = (
"The Sent-folder append outcome is unknown after an unexpected provider failure; "
"inspect and reconcile the mailbox before retrying."
)
_record_imap_provider_error(
session,
claimed,
folder=context.folder,
message=reason,
outcome_unknown=True,
)
raise ImapAppendError(reason, outcome_unknown=True) from None
try:
return _record_imap_append_success(
session,
job=claimed.job,
attempt=claimed.attempt,
claim_token=claimed.claim_token,
folder=result.folder,
)
except Exception:
return _mark_imap_append_outcome_unknown_after_effect(
session,
job_id=claimed.job.id,
attempt_id=claimed.attempt.id,
claim_token=claimed.claim_token,
reason=(
"The IMAP provider accepted the append, but persisting the outcome failed. "
"Automatic retry is stopped until an operator reconciles the mailbox."
),
)
def append_sent_for_job(
session: Session, *, job_id: str, dry_run: bool = False
) -> AppendSentResult:
"""Append one successfully sent job's exact EML to the configured IMAP Sent folder."""
job = session.get(CampaignJob, job_id)
if not job:
raise SendJobError(f"Job not found: {job_id}")
blocked = _imap_append_precondition(session, job, dry_run=dry_run)
if blocked is not None:
return blocked
prepared = _prepare_imap_append(session, job, dry_run=dry_run)
if isinstance(prepared, AppendSentResult):
return prepared
if dry_run:
return _imap_append_dry_run(session, job, prepared)
claimed = _claim_imap_append(session, job)
if isinstance(claimed, AppendSentResult):
return claimed
return _perform_imap_append(session, claimed, prepared)
def enqueue_pending_imap_appends(
session: Session,
*,
tenant_id: str,
campaign_id: str,
enqueue_celery: bool = True,
run_inline: bool = False,
dry_run: bool = False,
) -> dict[str, Any]:
campaign = _get_campaign_for_tenant(
session, campaign_id=campaign_id, tenant_id=tenant_id
)
jobs = (
session.query(CampaignJob)
.filter(
CampaignJob.tenant_id == tenant_id,
CampaignJob.campaign_id == campaign.id,
CampaignJob.imap_status.in_(
[JobImapStatus.PENDING.value, JobImapStatus.FAILED.value]
),
)
.order_by(CampaignJob.entry_index.asc())
.all()
)
jobs = [
job
for job in jobs
if _mail_was_accepted(
session,
job.id,
send_status=job.send_status,
)
]
should_enqueue = (
_should_enqueue_celery(enqueue_celery) and not dry_run and not run_inline
)
results: list[dict[str, Any]] = []
appended_count = 0
failed_count = 0
skipped_count = 0
if run_inline or dry_run:
for job in jobs:
try:
result = append_sent_for_job(session, job_id=job.id, dry_run=dry_run)
payload = result.as_dict()
results.append(payload)
if result.status == JobImapStatus.APPENDED.value:
appended_count += 1
elif result.status in {
"skipped",
"not_requested",
"not_sent",
"already_appended",
"dry_run",
}:
skipped_count += 1
except (
Exception
) as exc: # keep processing later jobs and expose per-job details
failed_count += 1
results.append(
{"job_id": job.id, "status": "failed", "message": str(exc)}
)
elif should_enqueue:
for job in jobs:
_celery_enqueue_append_sent_job(job.id)
return {
"campaign_id": campaign.id,
"pending_count": len(jobs),
"enqueued_count": len(jobs) if should_enqueue else 0,
"processed_count": len(results) if run_inline and not dry_run else 0,
"appended_count": appended_count,
"failed_count": failed_count,
"skipped_count": skipped_count,
"dry_run": dry_run,
"run_inline": run_inline,
"results": results,
}
def next_retry_delay(snapshot: ExecutionSnapshot, attempt_count: int) -> int:
delays = snapshot.delivery.retry.backoff_seconds or [60]
index = max(0, min(attempt_count - 1, len(delays) - 1))
return int(delays[index])