intermittent commit
This commit is contained in:
@@ -219,8 +219,10 @@ def create_execution_snapshot(
|
||||
queueable_statuses = {JobValidationStatus.READY.value, JobValidationStatus.WARNING.value}
|
||||
redacted_smtp = _redacted_transport_config(smtp)
|
||||
redacted_imap = _redacted_transport_config(imap)
|
||||
assert isinstance(redacted_smtp, SmtpConfig)
|
||||
assert redacted_imap is None or isinstance(redacted_imap, ImapConfig)
|
||||
if not isinstance(redacted_smtp, SmtpConfig):
|
||||
raise ExecutionSnapshotError("Redacted SMTP configuration is invalid.")
|
||||
if redacted_imap is not None and not isinstance(redacted_imap, ImapConfig):
|
||||
raise ExecutionSnapshotError("Redacted IMAP configuration is invalid.")
|
||||
payload = ExecutionSnapshot(
|
||||
campaign_version_id=version.id,
|
||||
campaign_json_sha256=_sha256(raw_json),
|
||||
|
||||
@@ -12,6 +12,7 @@ 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,
|
||||
@@ -28,6 +29,7 @@ from govoplan_campaign.backend.db.models import (
|
||||
SendAttempt,
|
||||
)
|
||||
from govoplan_campaign.backend.sending.execution import ExecutionSnapshot, ExecutionSnapshotError, ensure_execution_snapshot, runtime_imap_config, runtime_smtp_config
|
||||
from govoplan_campaign.backend.runtime import get_registry
|
||||
from govoplan_campaign.backend.integrations import (
|
||||
ImapAppendError,
|
||||
ImapConfigurationError,
|
||||
@@ -133,6 +135,25 @@ class AppendSentResult:
|
||||
}
|
||||
|
||||
|
||||
@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,
|
||||
@@ -140,6 +161,15 @@ QUEUEABLE_VALIDATION_STATUSES = {
|
||||
SMTP_ACCEPTED_STATUSES = {JobSendStatus.SMTP_ACCEPTED.value, JobSendStatus.SENT.value}
|
||||
AUTOMATICALLY_SENDABLE_STATUSES = {JobSendStatus.QUEUED.value}
|
||||
EXPLICIT_RETRY_STATUSES = {JobSendStatus.FAILED_TEMPORARY.value, JobSendStatus.FAILED_PERMANENT.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:
|
||||
@@ -210,16 +240,72 @@ def _should_enqueue_celery(enqueue_celery: bool) -> bool:
|
||||
return bool(enqueue_celery and _celery_enabled())
|
||||
|
||||
|
||||
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("multimailer.send_email", args=[job_id], queue="send_email")
|
||||
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("multimailer.append_sent", args=[job_id], queue="append_sent")
|
||||
celery.send_task("govoplan.campaigns.append_sent", args=[job_id], queue="append_sent")
|
||||
|
||||
|
||||
def queue_campaign_jobs(
|
||||
@@ -296,11 +382,20 @@ def queue_campaign_jobs(
|
||||
|
||||
if not dry_run:
|
||||
if queued:
|
||||
previous_status = campaign.status
|
||||
campaign.status = CampaignStatus.QUEUED.value
|
||||
version.workflow_state = CampaignVersionWorkflowState.QUEUED.value
|
||||
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()
|
||||
|
||||
@@ -469,8 +564,16 @@ def resume_campaign_jobs(session: Session, *, tenant_id: str, campaign_id: str,
|
||||
job.send_status = JobSendStatus.QUEUED.value
|
||||
session.add(job)
|
||||
if jobs:
|
||||
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
|
||||
@@ -909,27 +1012,54 @@ def _update_campaign_after_job(session: Session, campaign_id: str, version_id: s
|
||||
campaign = session.get(Campaign, campaign_id)
|
||||
if not campaign:
|
||||
return
|
||||
base_filters = [CampaignJob.campaign_id == campaign_id]
|
||||
if version_id:
|
||||
base_filters.append(CampaignJob.campaign_version_id == version_id)
|
||||
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}
|
||||
}
|
||||
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)
|
||||
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),
|
||||
)
|
||||
cancelled = counts.get(JobSendStatus.CANCELLED.value, 0)
|
||||
# Blocked/excluded build rows are reportable but are not delivery attempts.
|
||||
# Only a built, queueable message counts as an unattempted part of an
|
||||
# execution when deriving a partial outcome.
|
||||
not_started = (
|
||||
|
||||
|
||||
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,
|
||||
@@ -940,37 +1070,35 @@ def _update_campaign_after_job(session: Session, campaign_id: str, version_id: s
|
||||
.count()
|
||||
)
|
||||
|
||||
version = session.get(CampaignVersion, version_id) if version_id else None
|
||||
if active:
|
||||
campaign.status = CampaignStatus.SENDING.value
|
||||
if version:
|
||||
version.workflow_state = CampaignVersionWorkflowState.SENDING.value
|
||||
elif accepted and (failed or unknown or cancelled or not_started):
|
||||
campaign.status = CampaignStatus.PARTIALLY_COMPLETED.value
|
||||
if version:
|
||||
version.workflow_state = CampaignVersionWorkflowState.PARTIALLY_COMPLETED.value
|
||||
elif unknown:
|
||||
campaign.status = CampaignStatus.OUTCOME_UNKNOWN.value
|
||||
if version:
|
||||
version.workflow_state = CampaignVersionWorkflowState.OUTCOME_UNKNOWN.value
|
||||
elif failed:
|
||||
campaign.status = CampaignStatus.FAILED.value
|
||||
if version:
|
||||
version.workflow_state = CampaignVersionWorkflowState.FAILED.value
|
||||
elif accepted:
|
||||
campaign.status = CampaignStatus.SENT.value
|
||||
if version:
|
||||
version.workflow_state = CampaignVersionWorkflowState.COMPLETED.value
|
||||
elif cancelled:
|
||||
campaign.status = CampaignStatus.CANCELLED.value
|
||||
if version:
|
||||
version.workflow_state = CampaignVersionWorkflowState.CANCELLED.value
|
||||
|
||||
if version and (accepted or unknown):
|
||||
if version.locked_at is None:
|
||||
version.locked_at = _utcnow()
|
||||
session.add(version)
|
||||
session.add(campaign)
|
||||
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(
|
||||
@@ -984,15 +1112,49 @@ def send_campaign_job(
|
||||
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)
|
||||
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)
|
||||
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)
|
||||
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,
|
||||
job_id=job.id,
|
||||
status=JobSendStatus.OUTCOME_UNKNOWN.value,
|
||||
attempt_number=job.attempt_count,
|
||||
dry_run=dry_run,
|
||||
@@ -1014,7 +1176,10 @@ def send_campaign_job(
|
||||
)
|
||||
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")
|
||||
@@ -1028,134 +1193,196 @@ def send_campaign_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,
|
||||
)
|
||||
|
||||
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(envelope_recipients)} recipient(s) from {envelope_from}",
|
||||
)
|
||||
|
||||
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:
|
||||
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}",
|
||||
)
|
||||
return _not_claimed_send_job_result(session, job)
|
||||
|
||||
job = session.get(CampaignJob, job.id)
|
||||
assert job is not None
|
||||
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=snapshot.delivery.rate_limit.messages_per_minute,
|
||||
messages_per_minute=context.snapshot.delivery.rate_limit.messages_per_minute,
|
||||
enabled=use_rate_limit,
|
||||
)
|
||||
attempt = _record_attempt_start(session, job, claim_token)
|
||||
try:
|
||||
smtp_config = runtime_smtp_config(session, version, snapshot)
|
||||
smtp_config = runtime_smtp_config(session, context.version, context.snapshot)
|
||||
mail_integration().assert_mail_policy_allows_send(
|
||||
session,
|
||||
tenant_id=job.tenant_id,
|
||||
campaign_id=job.campaign_id,
|
||||
smtp=smtp_config,
|
||||
envelope_sender=envelope_from,
|
||||
from_header=_from_header_from_job(job, snapshot),
|
||||
recipients=envelope_recipients,
|
||||
envelope_sender=context.envelope_from,
|
||||
from_header=_from_header_from_job(job, context.snapshot),
|
||||
recipients=context.envelope_recipients,
|
||||
)
|
||||
result = mail_integration().send_email_bytes(
|
||||
message_bytes,
|
||||
context.message_bytes,
|
||||
smtp_config=smtp_config,
|
||||
envelope_from=envelope_from,
|
||||
envelope_recipients=envelope_recipients,
|
||||
envelope_from=context.envelope_from,
|
||||
envelope_recipients=context.envelope_recipients,
|
||||
)
|
||||
if result.accepted_count <= 0:
|
||||
raise SmtpSendError("SMTP did not accept any envelope recipients", temporary=False)
|
||||
refused_warning = None
|
||||
if result.refused_recipients:
|
||||
refused_warning = (
|
||||
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)}"
|
||||
)
|
||||
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
|
||||
if snapshot.delivery.imap_append_sent.enabled:
|
||||
job.imap_status = JobImapStatus.PENDING.value
|
||||
else:
|
||||
job.imap_status = 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()
|
||||
if enqueue_imap_task and _celery_enabled() and job.imap_status == JobImapStatus.PENDING.value:
|
||||
_celery_enqueue_append_sent_job(job.id)
|
||||
return SendJobResult(
|
||||
job_id=job.id,
|
||||
status=JobSendStatus.SMTP_ACCEPTED.value,
|
||||
attempt_number=attempt.attempt_number,
|
||||
message=refused_warning,
|
||||
return _record_smtp_send_success(
|
||||
session,
|
||||
job=job,
|
||||
attempt=attempt,
|
||||
snapshot=context.snapshot,
|
||||
result=result,
|
||||
enqueue_imap_task=enqueue_imap_task,
|
||||
)
|
||||
|
||||
except SmtpSendError as exc:
|
||||
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()
|
||||
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
|
||||
|
||||
|
||||
def _record_smtp_send_success(
|
||||
session: Session,
|
||||
*,
|
||||
job: CampaignJob,
|
||||
attempt: SendAttempt,
|
||||
snapshot: ExecutionSnapshot,
|
||||
result: object,
|
||||
enqueue_imap_task: bool,
|
||||
) -> SendJobResult:
|
||||
if result.accepted_count <= 0:
|
||||
raise SmtpSendError("SMTP did not accept any envelope recipients", temporary=False)
|
||||
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()
|
||||
if enqueue_imap_task and _celery_enabled() and job.imap_status == JobImapStatus.PENDING.value:
|
||||
_celery_enqueue_append_sent_job(job.id)
|
||||
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)
|
||||
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)
|
||||
job.claim_token = None
|
||||
session.add(job)
|
||||
_update_campaign_after_job(session, job.campaign_id, job.campaign_version_id)
|
||||
session.commit()
|
||||
raise
|
||||
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 _record_imap_attempt_start(session: Session, job: CampaignJob) -> ImapAppendAttempt:
|
||||
|
||||
Reference in New Issue
Block a user