196 lines
6.6 KiB
Python
196 lines
6.6 KiB
Python
from __future__ import annotations
|
|
|
|
import unittest
|
|
from types import SimpleNamespace
|
|
from unittest.mock import patch
|
|
|
|
from govoplan_campaign.backend.db.models import (
|
|
JobBuildStatus,
|
|
JobQueueStatus,
|
|
JobSendStatus,
|
|
JobValidationStatus,
|
|
)
|
|
from govoplan_campaign.backend.sending.jobs import (
|
|
SendJobResult,
|
|
_queue_validation_statuses,
|
|
_select_campaign_jobs_for_queue,
|
|
_send_claimed_campaign_job,
|
|
)
|
|
|
|
|
|
class FakeSession:
|
|
def __init__(self) -> None:
|
|
self.added: list[object] = []
|
|
|
|
def add(self, value: object) -> None:
|
|
self.added.append(value)
|
|
|
|
|
|
def _job(entry_id: str, **overrides):
|
|
values = {
|
|
"entry_id": entry_id,
|
|
"entry_index": int(entry_id),
|
|
"send_status": JobSendStatus.NOT_QUEUED.value,
|
|
"queue_status": JobQueueStatus.DRAFT.value,
|
|
"validation_status": JobValidationStatus.READY.value,
|
|
"build_status": JobBuildStatus.BUILT.value,
|
|
"eml_local_path": f"{entry_id}.eml",
|
|
"eml_storage_key": None,
|
|
"last_error": "old error",
|
|
"queued_at": None,
|
|
"claimed_at": "claimed",
|
|
"claim_token": "token",
|
|
"smtp_started_at": "started",
|
|
"outcome_unknown_at": "unknown",
|
|
}
|
|
values.update(overrides)
|
|
return SimpleNamespace(**values)
|
|
|
|
|
|
class CampaignQueueSelectionTests(unittest.TestCase):
|
|
def test_selects_queueable_jobs_without_reclassifying_retry_states(self):
|
|
skipped_send = _job("1", send_status=JobSendStatus.FAILED_TEMPORARY.value)
|
|
skipped_queue = _job("2", queue_status=JobQueueStatus.PAUSED.value)
|
|
blocked_validation = _job(
|
|
"3",
|
|
validation_status=JobValidationStatus.WARNING.value,
|
|
)
|
|
blocked_missing_eml = _job(
|
|
"4",
|
|
eml_local_path=None,
|
|
eml_storage_key=None,
|
|
)
|
|
ready = _job("5")
|
|
reviewed = _job(
|
|
"6",
|
|
validation_status=JobValidationStatus.NEEDS_REVIEW.value,
|
|
)
|
|
session = FakeSession()
|
|
|
|
queued, skipped_count, blocked_count = _select_campaign_jobs_for_queue(
|
|
session,
|
|
jobs=[
|
|
skipped_send,
|
|
skipped_queue,
|
|
blocked_validation,
|
|
blocked_missing_eml,
|
|
ready,
|
|
reviewed,
|
|
],
|
|
allowed_validation=_queue_validation_statuses(include_warnings=False),
|
|
reviewed_needs_review_keys={"6"},
|
|
dry_run=False,
|
|
)
|
|
|
|
self.assertEqual(queued, [ready, reviewed])
|
|
self.assertEqual(skipped_count, 2)
|
|
self.assertEqual(blocked_count, 2)
|
|
self.assertIn("generated EML", blocked_missing_eml.last_error)
|
|
self.assertEqual(session.added, [ready, reviewed])
|
|
for job in queued:
|
|
self.assertEqual(job.queue_status, JobQueueStatus.QUEUED.value)
|
|
self.assertEqual(job.send_status, JobSendStatus.QUEUED.value)
|
|
self.assertIsNotNone(job.queued_at)
|
|
self.assertIsNone(job.claimed_at)
|
|
self.assertIsNone(job.claim_token)
|
|
self.assertIsNone(job.smtp_started_at)
|
|
self.assertIsNone(job.outcome_unknown_at)
|
|
self.assertIsNone(job.last_error)
|
|
|
|
def test_dry_run_does_not_mutate_queueable_job(self):
|
|
warning = _job(
|
|
"1",
|
|
validation_status=JobValidationStatus.WARNING.value,
|
|
)
|
|
session = FakeSession()
|
|
|
|
queued, skipped_count, blocked_count = _select_campaign_jobs_for_queue(
|
|
session,
|
|
jobs=[warning],
|
|
allowed_validation=_queue_validation_statuses(include_warnings=True),
|
|
reviewed_needs_review_keys=set(),
|
|
dry_run=True,
|
|
)
|
|
|
|
self.assertEqual(queued, [warning])
|
|
self.assertEqual((skipped_count, blocked_count), (0, 0))
|
|
self.assertEqual(warning.queue_status, JobQueueStatus.DRAFT.value)
|
|
self.assertEqual(warning.send_status, JobSendStatus.NOT_QUEUED.value)
|
|
self.assertEqual(session.added, [])
|
|
|
|
def test_post_smtp_persistence_failure_is_outcome_unknown_not_retryable_failure(self):
|
|
job = SimpleNamespace(
|
|
id="job-1",
|
|
tenant_id="tenant-1",
|
|
campaign_id="campaign-1",
|
|
campaign_version_id="version-1",
|
|
imap_status="not_requested",
|
|
resolved_recipients={"from": {"email": "sender@example.test"}},
|
|
)
|
|
snapshot = SimpleNamespace(
|
|
mail_profile_id="profile-1",
|
|
smtp_transport_revision="frozen",
|
|
delivery=SimpleNamespace(
|
|
rate_limit=SimpleNamespace(messages_per_minute=60),
|
|
),
|
|
)
|
|
context = SimpleNamespace(
|
|
snapshot=snapshot,
|
|
message_bytes=b"message",
|
|
envelope_from="sender@example.test",
|
|
envelope_recipients=["recipient@example.test"],
|
|
)
|
|
current = SimpleNamespace(id="job-1")
|
|
|
|
class Session:
|
|
def __init__(self) -> None:
|
|
self.rolled_back = False
|
|
|
|
def rollback(self) -> None:
|
|
self.rolled_back = True
|
|
|
|
def get(self, _model, _id):
|
|
return current
|
|
|
|
class Mail:
|
|
def wait_for_rate_limit(self, **_kwargs):
|
|
return None
|
|
|
|
def send_campaign_email_bytes(self, *_args, **_kwargs):
|
|
return SimpleNamespace(accepted_count=1)
|
|
|
|
session = Session()
|
|
expected = SendJobResult(
|
|
job_id="job-1",
|
|
status=JobSendStatus.OUTCOME_UNKNOWN.value,
|
|
attempt_number=1,
|
|
)
|
|
with (
|
|
patch("govoplan_campaign.backend.sending.jobs.mail_integration", return_value=Mail()),
|
|
patch("govoplan_campaign.backend.sending.jobs._record_attempt_start", return_value=object()),
|
|
patch(
|
|
"govoplan_campaign.backend.sending.jobs._record_smtp_send_success",
|
|
side_effect=OSError("storage unavailable"),
|
|
),
|
|
patch(
|
|
"govoplan_campaign.backend.sending.jobs.mark_job_outcome_unknown",
|
|
return_value=expected,
|
|
) as mark_unknown,
|
|
):
|
|
result = _send_claimed_campaign_job(
|
|
session, # type: ignore[arg-type]
|
|
job=job, # type: ignore[arg-type]
|
|
claim_token="claim-1",
|
|
context=context, # type: ignore[arg-type]
|
|
use_rate_limit=False,
|
|
enqueue_imap_task=False,
|
|
)
|
|
|
|
self.assertIs(result, expected)
|
|
self.assertTrue(session.rolled_back)
|
|
self.assertIn("Automatic retry is stopped", mark_unknown.call_args.kwargs["reason"])
|
|
|
|
|
|
if __name__ == "__main__":
|
|
unittest.main()
|