feat(campaign): persist delivery execution mode

This commit is contained in:
2026-07-22 08:37:58 +02:00
parent aa4ec66b7b
commit 62a68792a4
9 changed files with 373 additions and 4 deletions

View File

@@ -201,6 +201,14 @@ QUEUEABLE_VALIDATION_STATUSES = {
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 | {
@@ -273,6 +281,13 @@ 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:
@@ -311,6 +326,10 @@ 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.",
@@ -601,11 +620,13 @@ def _persist_campaign_queue(
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)
@@ -638,9 +659,13 @@ def queue_campaign_jobs(
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)
@@ -667,6 +692,7 @@ def queue_campaign_jobs(
campaign=campaign,
version=version,
queued=queued,
delivery_mode=selected_delivery_mode,
)
enqueued_count = _enqueue_campaign_jobs(
queued,
@@ -680,7 +706,7 @@ def queue_campaign_jobs(
skipped_count=skipped_count,
blocked_count=blocked_count,
enqueued_count=enqueued_count,
delivery_mode="worker_queue" if enqueue_celery else "database_queue",
delivery_mode=selected_delivery_mode,
worker_queue_available=_celery_enabled(),
dry_run=dry_run,
)
@@ -738,6 +764,7 @@ def send_campaign_now(
include_warnings=include_warnings,
enqueue_celery=False,
dry_run=dry_run,
delivery_mode=DELIVERY_MODE_SYNCHRONOUS,
)
if dry_run:
return SendCampaignNowResult(
@@ -939,6 +966,12 @@ def resume_campaign_jobs(session: Session, *, tenant_id: str, campaign_id: str,
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)
@@ -1065,6 +1098,10 @@ def queue_failed_jobs_for_retry(
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()
@@ -1134,6 +1171,10 @@ def queue_unattempted_jobs(
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()
@@ -1202,6 +1243,7 @@ def send_single_campaign_job(
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()