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

2528 lines
91 KiB
Python

from __future__ import annotations
import hashlib
import json
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.orm import Session
from govoplan_core.core.notifications import NotificationDispatchRequest, notification_dispatch_provider
from govoplan_core.settings import settings as core_settings
from govoplan_campaign.backend.db.models import (
Campaign,
CampaignJob,
CampaignStatus,
CampaignVersion,
CampaignVersionWorkflowState,
JobBuildStatus,
JobImapStatus,
JobQueueStatus,
JobSendStatus,
JobValidationStatus,
ImapAppendAttempt,
SendAttempt,
)
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.runtime import get_registry
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
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
envelope_recipients: list[str]
QUEUEABLE_VALIDATION_STATUSES = {
JobValidationStatus.READY.value,
JobValidationStatus.WARNING.value,
}
SMTP_ACCEPTED_STATUSES = {JobSendStatus.SMTP_ACCEPTED.value, JobSendStatus.SENT.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 = SMTP_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 _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,
) -> 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)
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,
) -> 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)
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,
)
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,
)
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,
)
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 {JobSendStatus.SMTP_ACCEPTED.value, "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)
profile_summary = profile_delivery_summary(session, version)
expected_revision = next(iter(contexts.values())).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 SMTP_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]:
"""Explicitly queue failed jobs; never includes uncertain or accepted jobs."""
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}
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
if not force_max_attempts and job.attempt_count >= 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)
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.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,
include_warnings: bool = True,
dry_run: bool = False,
use_rate_limit: bool = True,
enqueue_imap_task: bool = False,
) -> dict[str, Any]:
"""Explicitly queue and send one built recipient message through the audit send path."""
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")
version = _get_current_version(session, campaign, version_id=job.campaign_version_id)
_ensure_version_validated_and_locked(version)
ensure_execution_snapshot(session, version)
queue_action = _prepare_single_job_for_send(
session,
version=version,
job=job,
include_warnings=include_warnings,
dry_run=dry_run,
)
if dry_run:
context = _send_job_delivery_context(session, job)
return {
"campaign_id": campaign.id,
"version_id": version.id,
"job_id": job.id,
"action": "single_send",
"queue_action": queue_action,
"dry_run": True,
"result": SendJobResult(
job_id=job.id,
status="dry_run",
attempt_number=job.attempt_count,
dry_run=True,
message=f"Would send to {len(context.envelope_recipients)} recipient(s) from {context.envelope_from}",
).as_dict(),
}
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,
)
result = send_campaign_job(
session,
job_id=job.id,
dry_run=False,
use_rate_limit=use_rate_limit,
enqueue_imap_task=enqueue_imap_task,
)
return {
"campaign_id": campaign.id,
"version_id": version.id,
"job_id": job.id,
"action": "single_send",
"queue_action": queue_action,
"dry_run": False,
"result": result.as_dict(),
}
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 SMTP_ACCEPTED_STATUSES:
raise QueueingError("This message has already been accepted by SMTP 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:
raise QueueingError("This message already has SMTP 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,
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,
)
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 _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 job.send_status not in SMTP_ACCEPTED_STATUSES:
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_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("Only local EML paths are supported for sending in this implementation step")
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 SMTP_ACCEPTED_STATUSES),
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.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:
return SendJobResult(
job_id=job.id,
status="dry_run",
attempt_number=job.attempt_count,
dry_run=True,
message=f"Would send to {len(context.envelope_recipients)} recipient(s) from {context.envelope_from}",
)
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 SMTP_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="SMTP 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 SMTP attempt was still marked in progress. Automatic resend 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)
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:
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 "",
)
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 _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."""
claim_token = str(uuid4())
changed = (
session.query(CampaignJob)
.filter(
CampaignJob.id == job.id,
CampaignJob.send_status.in_(list(SMTP_ACCEPTED_STATUSES)),
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 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}")
if job.send_status not in SMTP_ACCEPTED_STATUSES:
return AppendSentResult(job_id=job_id, status="not_sent", attempt_number=0, dry_run=dry_run, message="SMTP has not accepted this job")
blocked = _imap_blocked_result(session, job, dry_run=dry_run)
if blocked is not None:
return blocked
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)
message_bytes = _load_eml_bytes_for_job(job)
folder = snapshot.delivery.imap_append_sent.folder or "auto"
if dry_run:
attempts = session.query(ImapAppendAttempt).filter(ImapAppendAttempt.job_id == job.id).count()
return AppendSentResult(
job_id=job.id,
status="dry_run",
attempt_number=attempts,
dry_run=True,
folder=folder,
message=f"Would append {len(message_bytes)} bytes to IMAP folder {folder!r}",
)
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}",
)
job = session.get(CampaignJob, job.id)
if job is None:
raise SendJobError("Claimed campaign job disappeared before IMAP append")
attempt = _record_imap_attempt_start(session, job, claim_token)
try:
result = 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=message_bytes,
folder=None if folder == "auto" else folder,
expected_smtp_transport_revision=snapshot.smtp_transport_revision or "",
expected_imap_transport_revision=snapshot.imap_transport_revision,
)
except (MailProfileError, ImapConfigurationError, ImapAppendError) as exc:
outcome_unknown = bool(getattr(exc, "outcome_unknown", False))
owned = _record_imap_append_failure(
session,
job=job,
attempt=attempt,
claim_token=claim_token,
folder=folder,
message=str(exc),
outcome_unknown=outcome_unknown,
)
if not owned:
_imap_result_after_lost_claim(
session,
job_id=job.id,
attempt=attempt,
provider_succeeded=outcome_unknown,
)
raise
except Exception:
reason = (
"The Sent-folder append outcome is unknown after an unexpected provider failure; "
"inspect and reconcile the mailbox before retrying."
)
owned = _record_imap_append_failure(
session,
job=job,
attempt=attempt,
claim_token=claim_token,
folder=folder,
message=reason,
outcome_unknown=True,
)
if not owned:
_imap_result_after_lost_claim(
session,
job_id=job.id,
attempt=attempt,
provider_succeeded=True,
)
raise ImapAppendError(reason, outcome_unknown=True) from None
try:
return _record_imap_append_success(
session,
job=job,
attempt=attempt,
claim_token=claim_token,
folder=result.folder,
)
except Exception:
return _mark_imap_append_outcome_unknown_after_effect(
session,
job_id=job.id,
attempt_id=attempt.id,
claim_token=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 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.send_status.in_(list(SMTP_ACCEPTED_STATUSES)),
CampaignJob.imap_status.in_([JobImapStatus.PENDING.value, JobImapStatus.FAILED.value]),
)
.order_by(CampaignJob.entry_index.asc())
.all()
)
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])