506 lines
32 KiB
Python
506 lines
32 KiB
Python
"""Workerless recovery operates the real job/attempt ledger, never resend actions."""
|
|
from contextlib import nullcontext
|
|
from dataclasses import dataclass, field
|
|
from datetime import datetime, timedelta, timezone
|
|
from types import SimpleNamespace
|
|
from unittest.mock import Mock, patch
|
|
|
|
import pytest
|
|
from fastapi import FastAPI, HTTPException
|
|
from fastapi.testclient import TestClient
|
|
from sqlalchemy import Column, String, Table, create_engine, event
|
|
from sqlalchemy.orm import Session, sessionmaker
|
|
|
|
from govoplan_campaign.backend.db.models import (
|
|
Campaign, CampaignVersion, CampaignJob, SendAttempt, ImapAppendAttempt,
|
|
PostboxDeliveryAttempt, PrintOutputAttempt, CampaignMessageAction,
|
|
)
|
|
from govoplan_campaign.backend.delivery_policy import SynchronousSendPolicy
|
|
from govoplan_campaign.backend.integrations import MailProfileError, SmtpSendError, ImapAppendError
|
|
from govoplan_campaign.backend.routes import delivery as routes
|
|
from govoplan_campaign.backend.schemas import CampaignRetryJobsRequest, CampaignSendUnattemptedRequest, CampaignRecoverClaimRequest
|
|
from govoplan_campaign.backend.services.delivery_progress import campaign_delivery_progress
|
|
from govoplan_campaign.backend.services.delivery_recovery import job_recovery_metadata, recover_stale_delivery_claim, RecoveryStateConflict
|
|
from govoplan_campaign.backend.sending import jobs
|
|
from govoplan_core.core.change_sequence import ChangeSequenceEntry
|
|
from govoplan_core.core.recovery import RecoveryOperation, RecoveryCheckpoint, verify_recovery_evidence_chain
|
|
from govoplan_core.core.runtime_coordination import DistributedLease, RuntimeNode, process_runtime_identity
|
|
from govoplan_core.db.base import Base
|
|
from govoplan_core.db.session import get_session
|
|
from govoplan_core.auth import get_api_principal, ApiPrincipal
|
|
from govoplan_core.core.access import PrincipalRef
|
|
|
|
|
|
@dataclass
|
|
class _SmtpResult:
|
|
accepted_count: int = 1
|
|
refused_recipients: dict = field(default_factory=dict)
|
|
envelope_recipients: list = field(default_factory=lambda: ["recipient@example.test"])
|
|
|
|
|
|
@pytest.fixture
|
|
def recovery(tmp_path, monkeypatch):
|
|
engine = create_engine(f"sqlite+pysqlite:///{tmp_path / 'recovery.db'}")
|
|
for name in ("access_users", "access_groups"):
|
|
if name not in Base.metadata.tables:
|
|
Table(name, Base.metadata, Column("id", String(36), primary_key=True))
|
|
Base.metadata.create_all(engine, tables=[Base.metadata.tables[name] for name in ("access_users", "access_groups")] + [
|
|
model.__table__ for model in (ChangeSequenceEntry, Campaign, CampaignVersion, CampaignJob,
|
|
SendAttempt, ImapAppendAttempt, PostboxDeliveryAttempt, PrintOutputAttempt, CampaignMessageAction,
|
|
DistributedLease, RuntimeNode, RecoveryOperation, RecoveryCheckpoint)
|
|
])
|
|
factory = sessionmaker(engine)
|
|
session = factory()
|
|
campaign = Campaign(id="campaign", tenant_id="tenant", external_id="C", name="Recovery", current_version_id="version")
|
|
version = CampaignVersion(id="version", campaign_id="campaign", version_number=1, raw_json={},
|
|
locked_at=datetime.now(timezone.utc), validation_summary={"ok": True},
|
|
build_summary={"build_token": "build"}, execution_snapshot_hash="f" * 64,
|
|
editor_state={"review_send": {"build_token": "build", "inspection_complete": True, "reviewed_message_keys": ["reviewed"]}})
|
|
session.add_all([campaign, version])
|
|
session.commit()
|
|
snapshot = SimpleNamespace(
|
|
mail_profile_id="profile", smtp_server_id="smtp", smtp_credential_id="credential",
|
|
smtp_transport_revision="smtp-revision", imap_transport_revision="imap-revision", uses_mail=True,
|
|
delivery=SimpleNamespace(retry=SimpleNamespace(max_attempts=3),
|
|
rate_limit=SimpleNamespace(messages_per_minute=60), imap_append_sent=SimpleNamespace(enabled=True)),
|
|
)
|
|
provider = Mock()
|
|
provider.wait_for_rate_limit.return_value = None
|
|
provider.send_campaign_email_bytes.return_value = _SmtpResult()
|
|
provider.campaign_imap_batch.side_effect = lambda **_: nullcontext(SimpleNamespace(connection_count=1, reconnect_count=0))
|
|
monkeypatch.setattr(jobs, "get_database", lambda: SimpleNamespace(SessionLocal=factory))
|
|
monkeypatch.setattr(jobs, "ensure_execution_snapshot", lambda *_args, **_kw: snapshot)
|
|
monkeypatch.setattr(jobs, "_ensure_campaign_approval_gate", Mock())
|
|
monkeypatch.setattr(jobs, "_emit_campaign_status_notification", Mock())
|
|
monkeypatch.setattr(jobs, "_mark_accepted_job_artifacts", Mock())
|
|
monkeypatch.setattr(jobs, "mail_integration", lambda: provider)
|
|
monkeypatch.setattr(jobs, "_celery_enabled", lambda: False)
|
|
monkeypatch.setattr(jobs, "effective_synchronous_send_policy", lambda *_args, **_kw: SynchronousSendPolicy(2, "system", 500, system_max_recipient_jobs=2))
|
|
monkeypatch.setattr(jobs, "_synchronous_smtp_batch_manager", lambda *_args, **_kw: nullcontext(SimpleNamespace(connection_count=1, reconnect_count=0)))
|
|
monkeypatch.setattr(jobs, "_send_job_delivery_context", lambda _session, job: SimpleNamespace(
|
|
version=version, snapshot=snapshot, message_bytes=b"immutable message",
|
|
envelope_from="sender@example.test", envelope_recipients=["recipient@example.test"],
|
|
))
|
|
monkeypatch.setattr(jobs, "profile_delivery_summary", lambda *_: {"smtp_transport_revision": "smtp-revision"})
|
|
yield SimpleNamespace(session=session, factory=factory, engine=engine, campaign=campaign, version=version, provider=provider, snapshot=snapshot)
|
|
session.close()
|
|
engine.dispose()
|
|
|
|
|
|
def _add(recovery, name, *, status="not_queued", attempt=0, **kwargs):
|
|
job = CampaignJob(id=name, tenant_id="tenant", campaign_id="campaign", campaign_version_id="version",
|
|
entry_index=recovery.session.query(CampaignJob).count() + 1, entry_id=name,
|
|
build_status="built", validation_status="ready", send_status=status,
|
|
queue_status="draft", attempt_count=attempt, eml_sha256="e" * 64,
|
|
eml_local_path="unused-exact-message.eml",
|
|
resolved_recipients={"from": {"email": "sender@example.test"}, "to": [{"email": "recipient@example.test"}]})
|
|
for key, value in kwargs.items():
|
|
setattr(job, key, value)
|
|
recovery.session.add(job)
|
|
recovery.session.commit()
|
|
return job
|
|
|
|
|
|
def _retry(recovery, **kw):
|
|
return jobs.queue_failed_jobs_for_retry(recovery.session, tenant_id="tenant", campaign_id="campaign", version_id="version", enqueue_celery=False, run_inline=True, **kw)
|
|
|
|
|
|
def _continue(recovery, **kw):
|
|
return jobs.queue_unattempted_jobs(recovery.session, tenant_id="tenant", campaign_id="campaign", version_id="version", enqueue_celery=False, run_inline=True, **kw)
|
|
|
|
|
|
def test_inline_retry_records_canonical_acceptance_and_never_creates_resend_action(recovery):
|
|
failed = _add(recovery, "failed", status="failed_temporary", attempt=1)
|
|
recovery.session.add(SendAttempt(job_id=failed.id, attempt_number=1, status="failed_temporary", finished_at=datetime.now(timezone.utc)))
|
|
recovery.session.commit()
|
|
result = _retry(recovery)
|
|
assert result["sent_count"] == result["attempted_count"] == 1
|
|
with recovery.factory() as check:
|
|
current = check.get(CampaignJob, failed.id)
|
|
assert current.send_status == "smtp_accepted" and current.attempt_count == 2 and current.imap_status == "pending"
|
|
assert [a.status for a in check.query(SendAttempt).order_by(SendAttempt.attempt_number)] == ["failed_temporary", "smtp_accepted"]
|
|
assert check.query(CampaignMessageAction).count() == 0
|
|
ledger = check.query(RecoveryOperation).one()
|
|
assert ledger.status == "succeeded" and verify_recovery_evidence_chain(check, ledger.id)
|
|
second = _retry(recovery)
|
|
assert second["selected_count"] == 0 and recovery.provider.send_campaign_email_bytes.call_count == 1
|
|
|
|
|
|
def test_continue_is_bounded_and_skips_accepted_excluded_unknown_and_active(recovery):
|
|
for name in ("one", "two", "three"):
|
|
_add(recovery, name)
|
|
_add(recovery, "accepted", status="smtp_accepted", attempt=1)
|
|
_add(recovery, "excluded", status="skipped", validation_status="excluded")
|
|
_add(recovery, "unknown", status="outcome_unknown", attempt=1)
|
|
_add(recovery, "active", status="sending", attempt=1, claim_token="live")
|
|
_add(recovery, "claimed", status="claimed", claim_token="live-before-smtp")
|
|
first = _continue(recovery)
|
|
assert first["selected_count"] == first["sent_count"] == 2 and first["remaining_count"] == 1
|
|
assert recovery.session.get(CampaignJob, "three").send_status == "not_queued"
|
|
second = _continue(recovery)
|
|
assert second["sent_count"] == 1 and second["remaining_count"] == 0
|
|
assert recovery.provider.send_campaign_email_bytes.call_count == 3
|
|
assert recovery.session.get(CampaignJob, "active").send_status == "sending"
|
|
assert recovery.session.get(CampaignJob, "claimed").send_status == "claimed"
|
|
assert recovery.session.get(CampaignJob, "unknown").send_status == "outcome_unknown"
|
|
|
|
|
|
def test_continue_drains_queued_unattempted_remainder_but_never_repeats_prior_print_effect(recovery):
|
|
_add(recovery, "queued", status="queued", queue_status="queued")
|
|
_add(recovery, "printed", status="queued", queue_status="queued", print_attempt_count=1)
|
|
result = _continue(recovery)
|
|
assert result["sent_count"] == 1
|
|
assert recovery.session.get(CampaignJob, "queued").send_status == "smtp_accepted"
|
|
assert recovery.session.get(CampaignJob, "printed").send_status == "queued"
|
|
assert recovery.provider.send_campaign_email_bytes.call_count == 1
|
|
|
|
|
|
def test_unattempted_review_exception_requires_completed_same_build_review(recovery):
|
|
_add(recovery, "reviewed", validation_status="needs_review")
|
|
_add(recovery, "not-reviewed", validation_status="needs_review")
|
|
result = _continue(recovery)
|
|
assert result["selected_count"] == 1 and result["sent_count"] == 1
|
|
assert recovery.session.get(CampaignJob, "not-reviewed").send_status == "not_queued"
|
|
|
|
|
|
@pytest.mark.parametrize("reason", ["max_attempts", "permanent", "blocked", "approval", "limit_zero"])
|
|
def test_recovery_preserves_delivery_gates(recovery, monkeypatch, reason):
|
|
job = _add(recovery, "candidate", status="failed_temporary", attempt=1)
|
|
if reason == "max_attempts":
|
|
job.attempt_count = 3
|
|
elif reason == "permanent":
|
|
job.send_status = "failed_permanent"
|
|
elif reason == "blocked":
|
|
job.validation_status = "blocked"
|
|
elif reason == "approval":
|
|
monkeypatch.setattr(jobs, "_ensure_campaign_approval_gate", Mock(side_effect=jobs.QueueingError("Approval missing")))
|
|
else:
|
|
monkeypatch.setattr(jobs, "effective_synchronous_send_policy", lambda *_args, **_kw: SynchronousSendPolicy(0, "system", 500))
|
|
recovery.session.commit()
|
|
if reason in {"approval", "limit_zero"}:
|
|
with pytest.raises(jobs.QueueingError):
|
|
_retry(recovery)
|
|
else:
|
|
assert _retry(recovery)["selected_count"] == 0
|
|
assert recovery.provider.send_campaign_email_bytes.call_count == 0
|
|
|
|
|
|
def test_preflight_failure_rolls_back_queue_changes_without_provider_effect(recovery, monkeypatch):
|
|
job = _add(recovery, "failed", status="failed_temporary", attempt=1)
|
|
class Refuse:
|
|
def __enter__(self):
|
|
raise MailProfileError("Revoked profile selection")
|
|
def __exit__(self, *_):
|
|
return False
|
|
monkeypatch.setattr(jobs, "_synchronous_smtp_batch_manager", lambda *_args, **_kw: Refuse())
|
|
with pytest.raises(jobs.SynchronousSendRejected):
|
|
_retry(recovery)
|
|
recovery.session.expire_all()
|
|
assert job.send_status == "failed_temporary" and job.attempt_count == 1
|
|
recovery.provider.send_campaign_email_bytes.assert_not_called()
|
|
|
|
|
|
def test_repeated_worker_task_does_not_mutate_a_live_smtp_claim(recovery):
|
|
job = _add(recovery, "live", status="sending", attempt=1, claim_token="live", queue_status="sending")
|
|
result = jobs.send_campaign_job(recovery.session, job_id=job.id)
|
|
assert result.status == "already_sending"
|
|
recovery.session.expire_all()
|
|
assert job.send_status == "sending" and job.claim_token == "live"
|
|
recovery.provider.send_campaign_email_bytes.assert_not_called()
|
|
|
|
|
|
def test_progress_counts_include_active_and_partition_each_channel_without_loading_jobs(recovery):
|
|
_add(recovery, "accepted", status="smtp_accepted", imap_status="appended")
|
|
_add(recovery, "imap-active", status="smtp_accepted", imap_status="appending")
|
|
_add(recovery, "sending", status="sending", queue_status="sending", imap_status="pending")
|
|
_add(recovery, "claimed", status="claimed", imap_status="pending")
|
|
_add(recovery, "pending", imap_status="pending")
|
|
_add(recovery, "failed", status="failed_temporary", imap_status="failed")
|
|
_add(recovery, "unknown", status="outcome_unknown", imap_status="outcome_unknown")
|
|
_add(recovery, "paused", status="queued", queue_status="paused", imap_status="pending")
|
|
_add(recovery, "cancelled", status="cancelled", imap_status="skipped")
|
|
_add(recovery, "excluded", status="skipped", validation_status="excluded", imap_status="skipped")
|
|
recovery.session.expunge_all()
|
|
loaded = []
|
|
event.listen(recovery.session, "loaded_as_persistent", lambda _session, row: loaded.append(row))
|
|
progress = campaign_delivery_progress(recovery.session, tenant_id="tenant", campaign_id="campaign")
|
|
assert progress["total_jobs"] == 10
|
|
assert progress["smtp"] == dict(total=9, processed=5, accepted=2, active=2, pending=1, failed=1, outcome_unknown=1, excluded=1, paused=1, cancelled=1)
|
|
assert progress["imap"] == dict(total=8, processed=3, appended=1, active=1, pending=4, failed=1, outcome_unknown=1, excluded=2)
|
|
assert not any(isinstance(row, CampaignJob) for row in loaded)
|
|
assert progress["status_counts"]["send"]["claimed"] == 1
|
|
assert "recipient@example.test" not in repr(progress)
|
|
|
|
|
|
def test_progress_scopes_version_and_preserves_partial_multichannel_smtp_acceptance(recovery):
|
|
job = _add(recovery, "partial", status="partially_accepted", delivery_channel_policy="mail_and_postbox")
|
|
recovery.session.add(SendAttempt(job_id=job.id, attempt_number=1, status="smtp_accepted"))
|
|
recovery.session.add(CampaignVersion(id="old", campaign_id="campaign", version_number=2, raw_json={}))
|
|
recovery.session.commit()
|
|
result = campaign_delivery_progress(recovery.session, tenant_id="tenant", campaign_id="campaign", version_id="version")
|
|
assert result["smtp"]["accepted"] == 1 and result["smtp"]["failed"] == 0
|
|
assert campaign_delivery_progress(recovery.session, tenant_id="tenant", campaign_id="campaign", version_id="old")["total_jobs"] == 0
|
|
with pytest.raises(jobs.QueueingError):
|
|
campaign_delivery_progress(recovery.session, tenant_id="other", campaign_id="campaign")
|
|
|
|
|
|
def test_append_processes_only_selected_version_and_uses_one_lazy_batch(recovery, monkeypatch):
|
|
first = _add(recovery, "current", status="smtp_accepted", imap_status="pending")
|
|
recovery.session.add(CampaignVersion(id="old", campaign_id="campaign", version_number=2, raw_json={}))
|
|
recovery.session.commit()
|
|
old = _add(recovery, "old-job", status="smtp_accepted", imap_status="pending", campaign_version_id="old")
|
|
appended = []
|
|
def append(session, *, job_id, dry_run):
|
|
appended.append(job_id)
|
|
return jobs.AppendSentResult(job_id=job_id, status="appended", attempt_number=1)
|
|
monkeypatch.setattr(jobs, "append_sent_for_job", append)
|
|
result = jobs.enqueue_pending_imap_appends(recovery.session, tenant_id="tenant", campaign_id="campaign", enqueue_celery=False, run_inline=True)
|
|
assert appended == [first.id] and result["version_id"] == "version" and result["imap_connection_count"] == 1
|
|
jobs.enqueue_pending_imap_appends(recovery.session, tenant_id="tenant", campaign_id="campaign", version_id="old", enqueue_celery=False, run_inline=True)
|
|
assert appended == [first.id, old.id]
|
|
|
|
|
|
@pytest.mark.parametrize("outcome", ["failed", "outcome_unknown", "raised_unknown"])
|
|
def test_inline_imap_summary_keeps_failed_and_unknown_results_distinct(recovery, monkeypatch, outcome):
|
|
job = _add(recovery, "imap", status="smtp_accepted", imap_status="pending")
|
|
def append(*_args, **_kwargs):
|
|
if outcome == "raised_unknown":
|
|
raise ImapAppendError("Provider response lost", outcome_unknown=True)
|
|
return jobs.AppendSentResult(job_id=job.id, status=outcome, attempt_number=1)
|
|
monkeypatch.setattr(jobs, "append_sent_for_job", append)
|
|
result = jobs.enqueue_pending_imap_appends(recovery.session, tenant_id="tenant", campaign_id="campaign", run_inline=True)
|
|
assert result["appended_count"] == 0 and result["processed_count"] == 1
|
|
assert result["outcome_unknown_count"] == int(outcome != "failed")
|
|
assert result["failed_count"] == int(outcome == "failed")
|
|
assert result["results"][0]["status"] == ("outcome_unknown" if outcome == "raised_unknown" else outcome)
|
|
|
|
|
|
def _stale_claim(recovery, *, channel="smtp", expired=True, owner="stopped", policy="mail"):
|
|
job = _add(recovery, "stale", status="sending" if channel == "smtp" else "smtp_accepted", queue_status="sending" if channel == "smtp" else "draft",
|
|
attempt=1, claim_token="stale-token" if channel == "smtp" else None,
|
|
imap_status="appending" if channel == "imap" else "pending", imap_claim_token="stale-token" if channel == "imap" else None,
|
|
delivery_channel_policy=policy)
|
|
identity = process_runtime_identity()
|
|
context = SimpleNamespace(version=recovery.version, snapshot=recovery.snapshot, folder="Sent")
|
|
begin = jobs._begin_job_delivery_recovery if channel == "smtp" else jobs._begin_imap_append_recovery
|
|
started = begin(job=job, context=context, claim_token="stale-token")
|
|
recovery.session.expire_all()
|
|
lease = recovery.session.query(DistributedLease).one()
|
|
lease.holder_node_id = "old-process"
|
|
lease.holder_incarnation = "old-incarnation"
|
|
lease.expires_at = datetime.now(timezone.utc) + timedelta(minutes=-1 if expired else 10)
|
|
node = RuntimeNode(installation_id=identity.installation_id, node_id="old-process",
|
|
incarnation="replacement" if owner == "replaced" else "old-incarnation", role="worker", software_version="test", composition_hash="c" * 64,
|
|
state="stopped" if owner == "stopped" else "active")
|
|
recovery.session.add(node)
|
|
attempt = SendAttempt(job_id=job.id, attempt_number=1, status="smtp_in_progress", claim_token="stale-token") if channel == "smtp" else ImapAppendAttempt(job_id=job.id, attempt_number=1, status="appending", claim_token="stale-token")
|
|
recovery.session.add(attempt)
|
|
recovery.session.commit()
|
|
return job, lease, node, started.operation_id
|
|
|
|
|
|
@pytest.mark.parametrize("channel", ["smtp", "imap"])
|
|
@pytest.mark.parametrize("owner", ["stopped", "replaced"])
|
|
def test_stale_claim_recovery_is_fenced_unknown_then_explicit_evidence_reconciliation(recovery, channel, owner):
|
|
job, lease, node, operation_id = _stale_claim(recovery, channel=channel, owner=owner)
|
|
metadata = job_recovery_metadata(recovery.session, [job])[job.id][channel]
|
|
assert metadata["eligible"] is True
|
|
result = recover_stale_delivery_claim(recovery.session, tenant_id="tenant", campaign_id="campaign", job_id=job.id,
|
|
channel=channel, expected_revision=metadata["revision"], note="Verified original worker process stopped; inspect provider evidence next.")
|
|
recovery.session.commit()
|
|
assert result["reconciliation_required"] is True
|
|
assert (job.send_status if channel == "smtp" else job.imap_status) == "outcome_unknown"
|
|
assert recovery.session.get(RecoveryOperation, operation_id).status == "outcome_unknown"
|
|
assert verify_recovery_evidence_chain(recovery.session, operation_id)
|
|
assert _retry(recovery)["selected_count"] == 0
|
|
decision = "not_sent" if channel == "smtp" else "imap_not_appended"
|
|
jobs.reconcile_job_outcome(recovery.session, tenant_id="tenant", campaign_id="campaign", job_id=job.id, decision=decision, note="Provider logs and target mailbox confirm the effect did not occur.")
|
|
assert (job.send_status if channel == "smtp" else job.imap_status) == ("failed_temporary" if channel == "smtp" else "failed")
|
|
assert recovery.session.get(RecoveryOperation, operation_id).status == "recovered"
|
|
assert verify_recovery_evidence_chain(recovery.session, operation_id)
|
|
original_attempt = recovery.session.query(SendAttempt if channel == "smtp" else ImapAppendAttempt).filter_by(job_id=job.id).one()
|
|
assert original_attempt.status == ("reconciled_not_sent" if channel == "smtp" else "reconciled_imap_not_appended")
|
|
assert "Provider logs" in original_attempt.error_message
|
|
recovery.provider.send_campaign_email_bytes.assert_not_called()
|
|
|
|
|
|
@pytest.mark.parametrize("case", ["live_lease", "active_owner", "changed_revision", "wrong_tenant", "missing_ledger"])
|
|
def test_claim_recovery_rejects_live_ambiguous_changed_or_unauthorized_state(recovery, case):
|
|
job, lease, node, operation_id = _stale_claim(recovery, expired=case != "live_lease", owner="active" if case == "active_owner" else "stopped")
|
|
metadata = job_recovery_metadata(recovery.session, [job])[job.id]["smtp"]
|
|
if case == "missing_ledger":
|
|
recovery.session.get(RecoveryOperation, operation_id).idempotency_key = "another-operation"
|
|
recovery.session.commit()
|
|
with pytest.raises(jobs.QueueingError):
|
|
recover_stale_delivery_claim(recovery.session, tenant_id="other" if case == "wrong_tenant" else "tenant", campaign_id="campaign", job_id=job.id, channel="smtp", expected_revision="f" * 64 if case == "changed_revision" else metadata["revision"], note="Evidence")
|
|
recovery.session.rollback()
|
|
assert job.send_status == "sending" and job.claim_token == "stale-token"
|
|
recovery.provider.send_campaign_email_bytes.assert_not_called()
|
|
|
|
|
|
class _Principal:
|
|
tenant_id = "tenant"
|
|
user = SimpleNamespace(id="operator")
|
|
|
|
def __init__(self, *scopes):
|
|
self.scopes = set(scopes)
|
|
|
|
def has(self, scope):
|
|
return scope in self.scopes
|
|
|
|
|
|
@pytest.mark.parametrize("kind", ["retry", "unattempted"])
|
|
@pytest.mark.parametrize("missing", ["send", "recipient"])
|
|
def test_inline_recovery_requires_send_and_recipient_permission(recovery, kind, missing):
|
|
principal = _Principal(*({"campaigns:campaign:send", "campaigns:recipient:read"} - {"campaigns:campaign:send" if missing == "send" else "campaigns:recipient:read"}))
|
|
endpoint = routes.retry_campaign_jobs if kind == "retry" else routes.send_unattempted_campaign_jobs
|
|
schema = CampaignRetryJobsRequest if kind == "retry" else CampaignSendUnattemptedRequest
|
|
with patch.object(routes, "_get_campaign_for_principal", return_value=recovery.campaign):
|
|
with pytest.raises(HTTPException) as error:
|
|
endpoint("campaign", schema(version_id="version", run_inline=True), session=recovery.session, principal=principal)
|
|
assert error.value.status_code == 403
|
|
recovery.provider.send_campaign_email_bytes.assert_not_called()
|
|
|
|
|
|
def test_claim_recovery_and_audit_are_atomic(recovery):
|
|
job, lease, node, operation_id = _stale_claim(recovery)
|
|
metadata = job_recovery_metadata(recovery.session, [job])[job.id]["smtp"]
|
|
original_fence = lease.fencing_token
|
|
principal = _Principal("campaigns:recipient:read", "campaigns:campaign:reconcile")
|
|
with patch.object(routes, "_get_campaign_for_principal", return_value=recovery.campaign), patch.object(routes, "audit_from_principal", autospec=True, side_effect=RuntimeError("Audit unavailable")):
|
|
with pytest.raises(RuntimeError, match="Audit unavailable"):
|
|
routes.recover_campaign_job_claim("campaign", job.id, CampaignRecoverClaimRequest(channel="smtp", expected_revision=metadata["revision"], note="Stopped process verified"), session=recovery.session, principal=principal)
|
|
with recovery.factory() as check:
|
|
assert check.get(CampaignJob, job.id).send_status == "sending"
|
|
assert check.get(CampaignJob, job.id).claim_token == "stale-token"
|
|
assert check.get(RecoveryOperation, operation_id).status == "running"
|
|
assert check.get(DistributedLease, lease.id).fencing_token == original_fence
|
|
assert verify_recovery_evidence_chain(check, operation_id)
|
|
|
|
|
|
def test_existing_worker_queue_path_remains_supported(recovery, monkeypatch):
|
|
job = _add(recovery, "retry", status="failed_temporary", attempt=1)
|
|
enqueue = Mock()
|
|
monkeypatch.setattr(jobs, "_celery_enabled", lambda: True)
|
|
monkeypatch.setattr(jobs, "_celery_enqueue_send_job", enqueue)
|
|
result = jobs.queue_failed_jobs_for_retry(recovery.session, tenant_id="tenant", campaign_id="campaign", enqueue_celery=True)
|
|
assert result["enqueued_count"] == 1 and result["run_inline"] is False
|
|
enqueue.assert_called_once_with(job.id)
|
|
assert job.send_status == "queued" and job.attempt_count == 1
|
|
recovery.provider.send_campaign_email_bytes.assert_not_called()
|
|
|
|
|
|
def test_retry_dry_run_and_foreign_ids_never_change_selected_state(recovery):
|
|
job = _add(recovery, "retry", status="failed_temporary", attempt=1)
|
|
assert _retry(recovery, dry_run=True)["selected_count"] == 1
|
|
assert job.send_status == "failed_temporary" and job.attempt_count == 1
|
|
assert _retry(recovery, job_ids=["foreign-job"])["selected_count"] == 0
|
|
recovery.provider.send_campaign_email_bytes.assert_not_called()
|
|
|
|
|
|
def test_recovery_public_response_never_returns_provider_diagnostics():
|
|
projected = routes._public_recovery_result({"run_inline": True, "campaign_id": "campaign", "version_id": "version",
|
|
"selected_count": 1, "remaining_count": 0, "sent_count": 0, "failed_count": 1,
|
|
"results": [{"job_id": "job", "status": "failed", "message": "Provider rejected hidden@example.test"}]})
|
|
assert projected["selected_count"] == 1
|
|
assert projected["results"] == [{"job_id": "job", "status": "failed"}]
|
|
assert "hidden@example.test" not in repr(projected)
|
|
|
|
|
|
@pytest.mark.parametrize("result_status", ["failed_temporary", "failed_permanent", "outcome_unknown"])
|
|
def test_inline_summary_counts_returned_failures_and_unknown_separately(recovery, monkeypatch, result_status):
|
|
job = _add(recovery, "candidate", status="failed_temporary", attempt=1)
|
|
monkeypatch.setattr(jobs, "_deliver_job_with_recovery", lambda *_args, **_kw: jobs.SendJobResult(job_id=job.id, status=result_status, attempt_number=2))
|
|
result = _retry(recovery)
|
|
assert result["failed_count"] == int(result_status != "outcome_unknown")
|
|
assert result["outcome_unknown_count"] == int(result_status == "outcome_unknown")
|
|
assert result["skipped_count"] == 0
|
|
|
|
|
|
@pytest.mark.parametrize("path,scope,method,payload", [
|
|
("delivery-progress?version_id=version", "campaigns:campaign:read", "get", None),
|
|
("jobs/retry", "campaigns:campaign:retry", "post", {"version_id": "version", "dry_run": True}),
|
|
("jobs/send-unattempted", "campaigns:campaign:queue", "post", {"version_id": "version", "dry_run": True}),
|
|
("jobs/job/recover-claim", "campaigns:campaign:reconcile", "post", {"channel": "smtp", "expected_revision": "f" * 64, "note": "Evidence"}),
|
|
])
|
|
def test_http_scope_dependencies_reject_missing_route_permission(recovery, path, scope, method, payload):
|
|
app = FastAPI()
|
|
app.include_router(routes.router, prefix="/api/v1")
|
|
scopes = {"campaigns:campaign:send", "campaigns:recipient:read"}
|
|
actor = SimpleNamespace(id="operator")
|
|
app.dependency_overrides[get_api_principal] = lambda: ApiPrincipal(principal=PrincipalRef(account_id="operator", membership_id="operator", tenant_id="tenant", scopes=frozenset(scopes)), user=actor, account=actor)
|
|
app.dependency_overrides[get_session] = lambda: recovery.session
|
|
with TestClient(app) as client:
|
|
response = client.request(method, f"/api/v1/campaigns/campaign/{path}", **({"json": payload} if payload else {}))
|
|
assert response.status_code == 403
|
|
recovery.provider.send_campaign_email_bytes.assert_not_called()
|
|
|
|
|
|
@pytest.mark.parametrize("channel", ["smtp", "imap"])
|
|
def test_accepted_reconciliation_updates_only_matching_original_campaign_operation(recovery, channel):
|
|
job, lease, node, operation_id = _stale_claim(recovery, channel=channel)
|
|
meta = job_recovery_metadata(recovery.session, [job])[job.id][channel]
|
|
recover_stale_delivery_claim(recovery.session, tenant_id="tenant", campaign_id="campaign", job_id=job.id, channel=channel, expected_revision=meta["revision"], note="Stopped process verified")
|
|
recovery.session.commit()
|
|
jobs.reconcile_job_outcome(recovery.session, tenant_id="tenant", campaign_id="campaign", job_id=job.id,
|
|
decision="smtp_accepted" if channel == "smtp" else "imap_appended", note="Provider log confirms the exact Message-ID was accepted")
|
|
with recovery.factory() as check:
|
|
assert check.get(RecoveryOperation, operation_id).status == "succeeded"
|
|
assert verify_recovery_evidence_chain(check, operation_id)
|
|
assert check.query(RecoveryOperation).count() == 1
|
|
assert (check.get(CampaignJob, job.id).send_status if channel == "smtp" else check.get(CampaignJob, job.id).imap_status) == ("smtp_accepted" if channel == "smtp" else "appended")
|
|
recovery.provider.send_campaign_email_bytes.assert_not_called()
|
|
|
|
|
|
@pytest.mark.parametrize("channel", ["smtp", "imap"])
|
|
def test_reconciliation_audit_failure_rolls_back_campaign_attempt_and_core_operation(recovery, channel):
|
|
job, lease, node, operation_id = _stale_claim(recovery, channel=channel)
|
|
meta = job_recovery_metadata(recovery.session, [job])[job.id][channel]
|
|
recover_stale_delivery_claim(recovery.session, tenant_id="tenant", campaign_id="campaign", job_id=job.id, channel=channel, expected_revision=meta["revision"], note="Stopped process verified")
|
|
recovery.session.commit()
|
|
from govoplan_campaign.backend.schemas import CampaignResolveOutcomeRequest
|
|
with patch.object(routes, "_get_campaign_for_principal", return_value=recovery.campaign), patch.object(routes, "audit_from_principal", autospec=True, side_effect=RuntimeError("Audit unavailable")):
|
|
with pytest.raises(RuntimeError, match="Audit unavailable"):
|
|
routes.resolve_campaign_job_outcome("campaign", job.id, CampaignResolveOutcomeRequest(decision="not_sent" if channel == "smtp" else "imap_not_appended", note="Verified effect absent"), session=recovery.session, principal=_Principal("campaigns:recipient:read"))
|
|
with recovery.factory() as check:
|
|
assert check.get(RecoveryOperation, operation_id).status == "outcome_unknown"
|
|
assert (check.get(CampaignJob, job.id).send_status if channel == "smtp" else check.get(CampaignJob, job.id).imap_status) == "outcome_unknown"
|
|
assert check.query(SendAttempt if channel == "smtp" else ImapAppendAttempt).filter_by(job_id=job.id).one().status == "outcome_unknown"
|
|
|
|
|
|
@pytest.mark.parametrize("channel_policy", ["mail_then_postbox", "postbox_then_mail", "mail_then_print", "mail_and_postbox"])
|
|
def test_mixed_delivered_status_without_smtp_attempt_evidence_never_claims_mail_acceptance(recovery, channel_policy):
|
|
_add(recovery, "mixed", status="delivered", delivery_channel_policy=channel_policy)
|
|
result = campaign_delivery_progress(recovery.session, tenant_id="tenant", campaign_id="campaign")
|
|
assert result["smtp"]["accepted"] == 0
|
|
if channel_policy == "mail_and_postbox":
|
|
assert result["smtp"]["outcome_unknown"] == 1
|
|
else:
|
|
assert result["smtp"]["excluded"] == 1
|
|
|
|
|
|
def test_concurrent_legacy_reconciliation_cannot_overwrite_already_accepted_state(recovery):
|
|
job = _add(recovery, "legacy", status="outcome_unknown", attempt=1)
|
|
# Keep an old ORM projection, then commit another operator's decision.
|
|
with recovery.factory() as concurrent:
|
|
current = concurrent.get(CampaignJob, job.id)
|
|
current.send_status = "smtp_accepted"
|
|
concurrent.commit()
|
|
assert job.send_status == "outcome_unknown"
|
|
with pytest.raises(jobs.QueueingError, match="changed"):
|
|
jobs.reconcile_job_outcome(recovery.session, tenant_id="tenant", campaign_id="campaign", job_id=job.id, decision="not_sent", note="Stale operator evidence")
|
|
recovery.session.rollback()
|
|
assert job.send_status == "smtp_accepted"
|
|
|
|
|
|
def test_one_smtp_decision_never_resolves_compound_postbox_operation(recovery):
|
|
job, lease, node, operation_id = _stale_claim(recovery, policy="mail_and_postbox")
|
|
meta = job_recovery_metadata(recovery.session, [job])[job.id]["smtp"]
|
|
recover_stale_delivery_claim(recovery.session, tenant_id="tenant", campaign_id="campaign", job_id=job.id, channel="smtp", expected_revision=meta["revision"], note="Original process stopped")
|
|
recovery.session.commit()
|
|
jobs.reconcile_job_outcome(recovery.session, tenant_id="tenant", campaign_id="campaign", job_id=job.id,
|
|
decision="smtp_accepted", note="SMTP acceptance verified; Postbox remains separately unresolved")
|
|
assert job.send_status == "smtp_accepted"
|
|
assert recovery.session.get(RecoveryOperation, operation_id).status == "outcome_unknown"
|
|
assert verify_recovery_evidence_chain(recovery.session, operation_id)
|