test(campaign): complete delivery fault drills
This commit is contained in:
@@ -20,6 +20,7 @@ import json
|
||||
import os
|
||||
import shutil
|
||||
import socketserver
|
||||
import subprocess
|
||||
import sys
|
||||
import tempfile
|
||||
import threading
|
||||
@@ -37,7 +38,7 @@ from uuid import uuid4
|
||||
SCRIPT_ROOT = Path(__file__).resolve().parent
|
||||
REPOSITORY_ROOT = SCRIPT_ROOT.parents[1]
|
||||
DEFAULT_FIXTURE = REPOSITORY_ROOT / "examples" / "greenmail-delivery" / "campaign.json"
|
||||
EVIDENCE_SCHEMA = "govoplan.campaign.greenmail-acceptance.v1"
|
||||
EVIDENCE_SCHEMA = "govoplan.campaign.greenmail-acceptance.v2"
|
||||
EXPECTED_AUDIT_ACTIONS = frozenset(
|
||||
{
|
||||
"campaign.created",
|
||||
@@ -47,6 +48,7 @@ EXPECTED_AUDIT_ACTIONS = frozenset(
|
||||
"campaign.sent_now",
|
||||
"campaign.send_now_rejected",
|
||||
"campaign.append_sent_enqueued",
|
||||
"campaign.queued",
|
||||
}
|
||||
)
|
||||
FORBIDDEN_CAMPAIGN_KEYS = frozenset(
|
||||
@@ -114,6 +116,35 @@ REPORT_IMAP_STATUSES = frozenset(
|
||||
"skipped",
|
||||
}
|
||||
)
|
||||
DURABLE_ATTEMPT_STATUSES = frozenset(
|
||||
{
|
||||
"failed_permanent",
|
||||
"failed_temporary",
|
||||
"outcome_unknown",
|
||||
"smtp_accepted",
|
||||
"smtp_accepted_with_refusals",
|
||||
"smtp_in_progress",
|
||||
}
|
||||
)
|
||||
SMTP_FAULT_MODES = frozenset(
|
||||
{
|
||||
"temporary_data_response",
|
||||
"partial_recipient_refusal",
|
||||
"post_data_disconnect",
|
||||
"post_data_hold",
|
||||
}
|
||||
)
|
||||
WORKER_TASK_CODE = """
|
||||
import sys
|
||||
from govoplan_core.celery_app import send_email
|
||||
|
||||
result = send_email.run(sys.argv[1])
|
||||
if not isinstance(result, dict) or result.get("status") not in {
|
||||
"outcome_unknown",
|
||||
"smtp_accepted",
|
||||
}:
|
||||
raise SystemExit(2)
|
||||
"""
|
||||
|
||||
|
||||
class AcceptanceError(RuntimeError):
|
||||
@@ -195,6 +226,23 @@ class CampaignBoundary:
|
||||
}
|
||||
|
||||
|
||||
@dataclass(frozen=True, slots=True)
|
||||
class PreparedCampaignScenario:
|
||||
campaign_id: str
|
||||
version_id: str
|
||||
subject: str
|
||||
validation: dict[str, Any]
|
||||
build: dict[str, Any]
|
||||
boundary: CampaignBoundary
|
||||
|
||||
def public_evidence(self) -> dict[str, Any]:
|
||||
return {
|
||||
"validation": dict(self.validation),
|
||||
"build": dict(self.build),
|
||||
"campaign_mail_boundary": self.boundary.as_dict(),
|
||||
}
|
||||
|
||||
|
||||
def _env(name: str, default: str) -> str:
|
||||
return os.environ.get(name, default).strip() or default
|
||||
|
||||
@@ -337,6 +385,7 @@ def materialize_campaign_fixture(
|
||||
settings: TestbedSettings,
|
||||
scenario: str,
|
||||
run_token: str,
|
||||
additional_envelope_recipient: bool = False,
|
||||
) -> tuple[dict[str, Any], str]:
|
||||
raw = json.loads(fixture_path.read_text(encoding="utf-8"))
|
||||
if not isinstance(raw, dict):
|
||||
@@ -358,6 +407,18 @@ def materialize_campaign_fixture(
|
||||
raise AcceptanceError("Campaign fixture address shape is invalid")
|
||||
from_addresses[0]["email"] = settings.sender
|
||||
to_addresses[0]["email"] = settings.recipient
|
||||
if additional_envelope_recipient:
|
||||
try:
|
||||
local_part, domain = settings.recipient.rsplit("@", 1)
|
||||
except ValueError as exc:
|
||||
raise AcceptanceError("Campaign acceptance recipient must be an email address") from exc
|
||||
to_addresses.append(
|
||||
{
|
||||
"email": f"{local_part}+partial-{run_token}@{domain}",
|
||||
"name": "Controlled partial-refusal recipient",
|
||||
"type": "to",
|
||||
}
|
||||
)
|
||||
fields["acceptance_run"] = run_token
|
||||
template = raw.get("template")
|
||||
if not isinstance(template, dict):
|
||||
@@ -550,7 +611,39 @@ def _allowlisted_count_map(
|
||||
return projected
|
||||
|
||||
|
||||
def execute_campaign_scenario(
|
||||
def _durable_state_evidence(payload: Mapping[str, Any]) -> dict[str, Any]:
|
||||
job_count = payload.get("job_count")
|
||||
unfinished_attempt_count = payload.get("unfinished_attempt_count")
|
||||
if (
|
||||
not isinstance(job_count, int)
|
||||
or isinstance(job_count, bool)
|
||||
or job_count < 0
|
||||
or not isinstance(unfinished_attempt_count, int)
|
||||
or isinstance(unfinished_attempt_count, bool)
|
||||
or unfinished_attempt_count < 0
|
||||
):
|
||||
raise AcceptanceError("Campaign durable delivery state count is invalid")
|
||||
send_statuses = payload.get("send_status_counts")
|
||||
attempt_statuses = payload.get("attempt_status_counts")
|
||||
if not isinstance(send_statuses, dict) or not isinstance(attempt_statuses, dict):
|
||||
raise AcceptanceError("Campaign durable delivery state is invalid")
|
||||
return {
|
||||
"job_count": job_count,
|
||||
"send_status_counts": _allowlisted_count_map(
|
||||
send_statuses,
|
||||
allowed=REPORT_SEND_STATUSES,
|
||||
label="Campaign durable send status",
|
||||
),
|
||||
"attempt_status_counts": _allowlisted_count_map(
|
||||
attempt_statuses,
|
||||
allowed=DURABLE_ATTEMPT_STATUSES,
|
||||
label="Campaign durable attempt status",
|
||||
),
|
||||
"unfinished_attempt_count": unfinished_attempt_count,
|
||||
}
|
||||
|
||||
|
||||
def prepare_campaign_scenario(
|
||||
client: ApiClient,
|
||||
headers: Mapping[str, str],
|
||||
*,
|
||||
@@ -559,10 +652,8 @@ def execute_campaign_scenario(
|
||||
settings: TestbedSettings,
|
||||
scenario: str,
|
||||
snapshot_probe: Callable[[str], tuple[Mapping[str, Any], Mapping[str, Any]]],
|
||||
audit_probe: Callable[[str, str], Mapping[str, int]],
|
||||
append_sent: bool,
|
||||
repeat_send: bool,
|
||||
) -> tuple[dict[str, Any], str]:
|
||||
additional_envelope_recipient: bool = False,
|
||||
) -> PreparedCampaignScenario:
|
||||
run_token = uuid4().hex[:12]
|
||||
raw, subject = materialize_campaign_fixture(
|
||||
fixture_path,
|
||||
@@ -570,6 +661,7 @@ def execute_campaign_scenario(
|
||||
settings=settings,
|
||||
scenario=scenario,
|
||||
run_token=run_token,
|
||||
additional_envelope_recipient=additional_envelope_recipient,
|
||||
)
|
||||
created = _expect(
|
||||
client.post("/api/v1/campaigns", headers=dict(headers), json={"config": raw}),
|
||||
@@ -608,6 +700,51 @@ def execute_campaign_scenario(
|
||||
|
||||
stored_raw, execution_snapshot = snapshot_probe(version_id)
|
||||
boundary = assert_campaign_boundary(stored_raw, execution_snapshot, profile_id=profile_id)
|
||||
return PreparedCampaignScenario(
|
||||
campaign_id=campaign_id,
|
||||
version_id=version_id,
|
||||
subject=subject,
|
||||
validation={
|
||||
"ok": True,
|
||||
"error_count": int(validation.get("error_count") or 0),
|
||||
"warning_count": int(validation.get("warning_count") or 0),
|
||||
},
|
||||
build={
|
||||
"built_count": int(build.get("built_count") or 0),
|
||||
"build_failed_count": int(build.get("build_failed_count") or 0),
|
||||
"queueable_count": int(build.get("queueable_count") or 0),
|
||||
},
|
||||
boundary=boundary,
|
||||
)
|
||||
|
||||
|
||||
def execute_campaign_scenario(
|
||||
client: ApiClient,
|
||||
headers: Mapping[str, str],
|
||||
*,
|
||||
fixture_path: Path,
|
||||
profile_id: str,
|
||||
settings: TestbedSettings,
|
||||
scenario: str,
|
||||
snapshot_probe: Callable[[str], tuple[Mapping[str, Any], Mapping[str, Any]]],
|
||||
audit_probe: Callable[[str, str], Mapping[str, int]],
|
||||
append_sent: bool,
|
||||
repeat_send: bool,
|
||||
delivery_probe: Callable[[str, str], Mapping[str, Any]] | None = None,
|
||||
additional_envelope_recipient: bool = False,
|
||||
) -> tuple[dict[str, Any], str]:
|
||||
prepared = prepare_campaign_scenario(
|
||||
client,
|
||||
headers,
|
||||
fixture_path=fixture_path,
|
||||
profile_id=profile_id,
|
||||
settings=settings,
|
||||
scenario=scenario,
|
||||
snapshot_probe=snapshot_probe,
|
||||
additional_envelope_recipient=additional_envelope_recipient,
|
||||
)
|
||||
campaign_id = prepared.campaign_id
|
||||
version_id = prepared.version_id
|
||||
send_payload = _expect(
|
||||
client.post(
|
||||
f"/api/v1/campaigns/{campaign_id}/send-now",
|
||||
@@ -699,17 +836,7 @@ def execute_campaign_scenario(
|
||||
raise AcceptanceError(f"{scenario} Campaign audit evidence is incomplete")
|
||||
|
||||
evidence = {
|
||||
"validation": {
|
||||
"ok": True,
|
||||
"error_count": int(validation.get("error_count") or 0),
|
||||
"warning_count": int(validation.get("warning_count") or 0),
|
||||
},
|
||||
"build": {
|
||||
"built_count": int(build.get("built_count") or 0),
|
||||
"build_failed_count": int(build.get("build_failed_count") or 0),
|
||||
"queueable_count": int(build.get("queueable_count") or 0),
|
||||
},
|
||||
"campaign_mail_boundary": boundary.as_dict(),
|
||||
**prepared.public_evidence(),
|
||||
"send": _send_evidence(send_payload),
|
||||
"report": _report_evidence(report),
|
||||
"audit_actions": audit_actions,
|
||||
@@ -718,7 +845,9 @@ def execute_campaign_scenario(
|
||||
evidence["repeated_send"] = repeat_evidence
|
||||
if append_evidence is not None:
|
||||
evidence["append_sent"] = append_evidence
|
||||
return evidence, subject
|
||||
if delivery_probe is not None:
|
||||
evidence["durable_state"] = _durable_state_evidence(delivery_probe(campaign_id, version_id))
|
||||
return evidence, prepared.subject
|
||||
|
||||
|
||||
def _mailbox_state(
|
||||
@@ -800,6 +929,111 @@ class _LoopbackServer(socketserver.TCPServer):
|
||||
allow_reuse_address = True
|
||||
|
||||
|
||||
class _SmtpFaultServer(_LoopbackServer):
|
||||
def __init__(self, mode: str) -> None:
|
||||
if mode not in SMTP_FAULT_MODES:
|
||||
raise AcceptanceError("Unsupported SMTP fault mode")
|
||||
self.mode = mode
|
||||
self.data_received = threading.Event()
|
||||
self.release_data_response = threading.Event()
|
||||
self._counter_lock = threading.Lock()
|
||||
self._counters = {
|
||||
"connection_count": 0,
|
||||
"accepted_rcpt_commands": 0,
|
||||
"refused_rcpt_commands": 0,
|
||||
"data_transactions": 0,
|
||||
}
|
||||
super().__init__(("127.0.0.1", 0), _FaultSmtpHandler)
|
||||
|
||||
def increment(self, key: str) -> int:
|
||||
with self._counter_lock:
|
||||
self._counters[key] += 1
|
||||
return self._counters[key]
|
||||
|
||||
def evidence(self) -> dict[str, int]:
|
||||
with self._counter_lock:
|
||||
return dict(self._counters)
|
||||
|
||||
|
||||
class _FaultSmtpHandler(socketserver.StreamRequestHandler):
|
||||
def handle(self) -> None:
|
||||
server = self.server
|
||||
if not isinstance(server, _SmtpFaultServer):
|
||||
return
|
||||
server.increment("connection_count")
|
||||
self._write(b"220 local acceptance SMTP\r\n")
|
||||
while line := self.rfile.readline(65_536):
|
||||
command = line.split(None, 1)[0].upper()
|
||||
if command in {b"EHLO", b"HELO"}:
|
||||
self._write(b"250-localhost\r\n250 AUTH PLAIN\r\n")
|
||||
elif command == b"AUTH":
|
||||
self._write(b"235 2.7.0 Authentication successful\r\n")
|
||||
elif command == b"MAIL":
|
||||
self._write(b"250 2.1.0 Sender accepted\r\n")
|
||||
elif command == b"RCPT":
|
||||
rcpt_number = (
|
||||
server.evidence()["accepted_rcpt_commands"]
|
||||
+ server.evidence()["refused_rcpt_commands"]
|
||||
+ 1
|
||||
)
|
||||
if server.mode == "partial_recipient_refusal" and rcpt_number > 1:
|
||||
server.increment("refused_rcpt_commands")
|
||||
self._write(b"550 5.1.1 Controlled recipient refusal\r\n")
|
||||
else:
|
||||
server.increment("accepted_rcpt_commands")
|
||||
self._write(b"250 2.1.5 Recipient accepted\r\n")
|
||||
elif command == b"DATA":
|
||||
self._write(b"354 End data with <CR><LF>.<CR><LF>\r\n")
|
||||
if not self._read_data_transaction():
|
||||
return
|
||||
server.increment("data_transactions")
|
||||
server.data_received.set()
|
||||
if server.mode == "temporary_data_response":
|
||||
self._write(b"451 4.3.0 Controlled temporary rejection\r\n")
|
||||
elif server.mode == "post_data_disconnect":
|
||||
return
|
||||
elif server.mode == "post_data_hold":
|
||||
server.release_data_response.wait()
|
||||
return
|
||||
else:
|
||||
self._write(b"250 2.0.0 Message accepted\r\n")
|
||||
elif command == b"RSET":
|
||||
self._write(b"250 2.0.0 Reset\r\n")
|
||||
elif command == b"NOOP":
|
||||
self._write(b"250 2.0.0 OK\r\n")
|
||||
elif command == b"QUIT":
|
||||
self._write(b"221 2.0.0 Bye\r\n")
|
||||
return
|
||||
else:
|
||||
self._write(b"500 5.5.1 Unsupported acceptance command\r\n")
|
||||
|
||||
def _read_data_transaction(self) -> bool:
|
||||
while line := self.rfile.readline(65_536):
|
||||
if line in {b".\n", b".\r\n"}:
|
||||
return True
|
||||
return False
|
||||
|
||||
def _write(self, value: bytes) -> None:
|
||||
self.wfile.write(value)
|
||||
self.wfile.flush()
|
||||
|
||||
|
||||
@dataclass(frozen=True, slots=True)
|
||||
class SmtpFaultEndpoint:
|
||||
host: str
|
||||
port: int
|
||||
_server: _SmtpFaultServer
|
||||
|
||||
def wait_for_data(self, timeout_seconds: float) -> bool:
|
||||
return self._server.data_received.wait(timeout_seconds)
|
||||
|
||||
def evidence(self) -> dict[str, int]:
|
||||
return self._server.evidence()
|
||||
|
||||
def release_held_connection(self) -> None:
|
||||
self._server.release_data_response.set()
|
||||
|
||||
|
||||
class _RejectingSmtpHandler(socketserver.StreamRequestHandler):
|
||||
def handle(self) -> None:
|
||||
self.wfile.write(b"220 local acceptance SMTP\r\n")
|
||||
@@ -866,6 +1100,25 @@ def rejecting_imap_endpoint() -> Iterator[tuple[str, int]]:
|
||||
yield endpoint
|
||||
|
||||
|
||||
@contextmanager
|
||||
def smtp_fault_endpoint(mode: str) -> Iterator[SmtpFaultEndpoint]:
|
||||
server = _SmtpFaultServer(mode)
|
||||
thread = threading.Thread(
|
||||
target=server.serve_forever,
|
||||
name=f"govoplan-smtp-{mode.replace('_', '-')}",
|
||||
daemon=True,
|
||||
)
|
||||
thread.start()
|
||||
try:
|
||||
host, port = server.server_address
|
||||
yield SmtpFaultEndpoint(str(host), int(port), server)
|
||||
finally:
|
||||
server.release_data_response.set()
|
||||
server.shutdown()
|
||||
server.server_close()
|
||||
thread.join(timeout=2)
|
||||
|
||||
|
||||
@contextmanager
|
||||
def _scripted_endpoint(
|
||||
handler: type[socketserver.BaseRequestHandler],
|
||||
@@ -911,6 +1164,187 @@ def ensure_sent_folder(settings: TestbedSettings) -> None:
|
||||
raise AcceptanceError("GreenMail IMAP did not become ready for the Campaign acceptance run") from last_error
|
||||
|
||||
|
||||
def _start_campaign_worker_task(job_id: str) -> subprocess.Popen[bytes]:
|
||||
"""Start the exact Celery send task body in an isolated OS process.
|
||||
|
||||
The acceptance environment deliberately has no Redis broker. Calling the
|
||||
registered task's ``run`` method exercises the same task/capability path a
|
||||
Celery child process uses, without claiming broker redelivery coverage.
|
||||
"""
|
||||
|
||||
return subprocess.Popen(
|
||||
[sys.executable, "-c", WORKER_TASK_CODE, job_id],
|
||||
stdin=subprocess.DEVNULL,
|
||||
stdout=subprocess.DEVNULL,
|
||||
stderr=subprocess.DEVNULL,
|
||||
close_fds=True,
|
||||
)
|
||||
|
||||
|
||||
def _terminate_worker_process(process: subprocess.Popen[bytes]) -> None:
|
||||
if process.poll() is not None:
|
||||
raise AcceptanceError("Campaign worker task exited before controlled interruption")
|
||||
process.terminate()
|
||||
try:
|
||||
process.wait(timeout=5)
|
||||
except subprocess.TimeoutExpired:
|
||||
process.kill()
|
||||
try:
|
||||
process.wait(timeout=5)
|
||||
except subprocess.TimeoutExpired as exc:
|
||||
raise AcceptanceError("Campaign worker task could not be stopped") from exc
|
||||
|
||||
|
||||
def _wait_for_worker_process(process: subprocess.Popen[bytes], *, timeout_seconds: int) -> None:
|
||||
try:
|
||||
return_code = process.wait(timeout=timeout_seconds)
|
||||
except subprocess.TimeoutExpired as exc:
|
||||
process.kill()
|
||||
process.wait(timeout=5)
|
||||
raise AcceptanceError("Restarted Campaign worker task did not complete") from exc
|
||||
if return_code != 0:
|
||||
raise AcceptanceError("Restarted Campaign worker task failed")
|
||||
|
||||
|
||||
def execute_worker_interruption_scenario(
|
||||
client: ApiClient,
|
||||
headers: Mapping[str, str],
|
||||
*,
|
||||
fixture_path: Path,
|
||||
profile_id: str,
|
||||
settings: TestbedSettings,
|
||||
endpoint: SmtpFaultEndpoint,
|
||||
snapshot_probe: Callable[[str], tuple[Mapping[str, Any], Mapping[str, Any]]],
|
||||
audit_probe: Callable[[str, str], Mapping[str, int]],
|
||||
delivery_probe: Callable[[str, str], Mapping[str, Any]],
|
||||
worker_job_probe: Callable[[str], str],
|
||||
) -> dict[str, Any]:
|
||||
prepared = prepare_campaign_scenario(
|
||||
client,
|
||||
headers,
|
||||
fixture_path=fixture_path,
|
||||
profile_id=profile_id,
|
||||
settings=settings,
|
||||
scenario="worker_interruption",
|
||||
snapshot_probe=snapshot_probe,
|
||||
)
|
||||
queue_payload = _expect(
|
||||
client.post(
|
||||
f"/api/v1/campaigns/{prepared.campaign_id}/queue",
|
||||
headers=dict(headers),
|
||||
json={
|
||||
"version_id": prepared.version_id,
|
||||
"include_warnings": True,
|
||||
"enqueue_celery": False,
|
||||
"dry_run": False,
|
||||
},
|
||||
),
|
||||
200,
|
||||
"Worker-interruption Campaign queue",
|
||||
)
|
||||
expected_queue = {
|
||||
"queued_count": 1,
|
||||
"skipped_count": 0,
|
||||
"blocked_count": 0,
|
||||
"enqueued_count": 0,
|
||||
"delivery_mode": "database_queue",
|
||||
"worker_queue_available": False,
|
||||
"dry_run": False,
|
||||
}
|
||||
queue_evidence = {key: queue_payload.get(key) for key in expected_queue}
|
||||
if queue_evidence != expected_queue:
|
||||
raise AcceptanceError("Worker-interruption Campaign was not durably queued once")
|
||||
|
||||
job_id = worker_job_probe(prepared.version_id)
|
||||
first_worker = _start_campaign_worker_task(job_id)
|
||||
try:
|
||||
if not endpoint.wait_for_data(settings.provider_timeout_seconds):
|
||||
raise AcceptanceError("Campaign worker did not reach the controlled SMTP DATA boundary")
|
||||
_terminate_worker_process(first_worker)
|
||||
finally:
|
||||
if first_worker.poll() is None:
|
||||
first_worker.kill()
|
||||
first_worker.wait(timeout=5)
|
||||
endpoint.release_held_connection()
|
||||
|
||||
interrupted_state = _durable_state_evidence(
|
||||
delivery_probe(prepared.campaign_id, prepared.version_id)
|
||||
)
|
||||
if interrupted_state != {
|
||||
"job_count": 1,
|
||||
"send_status_counts": {"sending": 1},
|
||||
"attempt_status_counts": {"smtp_in_progress": 1},
|
||||
"unfinished_attempt_count": 1,
|
||||
}:
|
||||
raise AcceptanceError("Interrupted worker state was not durably retained at SMTP in-progress")
|
||||
|
||||
restarted_worker = _start_campaign_worker_task(job_id)
|
||||
_wait_for_worker_process(
|
||||
restarted_worker,
|
||||
timeout_seconds=settings.provider_timeout_seconds,
|
||||
)
|
||||
restarted_state = _durable_state_evidence(
|
||||
delivery_probe(prepared.campaign_id, prepared.version_id)
|
||||
)
|
||||
expected_restarted_state = {
|
||||
"job_count": 1,
|
||||
"send_status_counts": {"outcome_unknown": 1},
|
||||
"attempt_status_counts": {"outcome_unknown": 1},
|
||||
"unfinished_attempt_count": 0,
|
||||
}
|
||||
if restarted_state != expected_restarted_state:
|
||||
raise AcceptanceError("Restarted worker did not freeze the unfinished SMTP attempt")
|
||||
|
||||
protocol_evidence = endpoint.evidence()
|
||||
if protocol_evidence != {
|
||||
"connection_count": 1,
|
||||
"accepted_rcpt_commands": 1,
|
||||
"refused_rcpt_commands": 0,
|
||||
"data_transactions": 1,
|
||||
}:
|
||||
raise AcceptanceError("Restarted worker contacted SMTP or produced an unexpected transaction")
|
||||
|
||||
report = _expect(
|
||||
client.get(
|
||||
f"/api/v1/campaigns/{prepared.campaign_id}/report",
|
||||
headers=dict(headers),
|
||||
params={"version_id": prepared.version_id},
|
||||
),
|
||||
200,
|
||||
"Worker-interruption Campaign report",
|
||||
)
|
||||
report_evidence = _report_evidence(report)
|
||||
if report_evidence["send_status_counts"] != {"outcome_unknown": 1}:
|
||||
raise AcceptanceError("Worker-interruption Campaign report did not retain outcome unknown")
|
||||
|
||||
audit_actions = dict(
|
||||
sorted(audit_probe(prepared.campaign_id, prepared.version_id).items())
|
||||
)
|
||||
if not {
|
||||
"campaign.created",
|
||||
"campaign.validated",
|
||||
"campaign.messages_built",
|
||||
"campaign.queued",
|
||||
}.issubset(audit_actions):
|
||||
raise AcceptanceError("Worker-interruption Campaign audit evidence is incomplete")
|
||||
|
||||
return {
|
||||
**prepared.public_evidence(),
|
||||
"queue": queue_evidence,
|
||||
"interrupted_durable_state": interrupted_state,
|
||||
"restarted_durable_state": restarted_state,
|
||||
"protocol": protocol_evidence,
|
||||
"report": report_evidence,
|
||||
"audit_actions": audit_actions,
|
||||
"process_boundary": {
|
||||
"dedicated_task_process_terminated_after_data": True,
|
||||
"fresh_task_process_completed": True,
|
||||
"duplicate_smtp_transaction_prevented": True,
|
||||
"celery_broker_redelivery_exercised": False,
|
||||
},
|
||||
}
|
||||
|
||||
|
||||
def run_acceptance(
|
||||
client: ApiClient,
|
||||
headers: Mapping[str, str],
|
||||
@@ -919,6 +1353,8 @@ def run_acceptance(
|
||||
fixture_path: Path,
|
||||
snapshot_probe: Callable[[str], tuple[Mapping[str, Any], Mapping[str, Any]]],
|
||||
audit_probe: Callable[[str, str], Mapping[str, int]],
|
||||
delivery_probe: Callable[[str, str], Mapping[str, Any]],
|
||||
worker_job_probe: Callable[[str], str],
|
||||
include_failure_drills: bool,
|
||||
module_versions: Mapping[str, str],
|
||||
) -> dict[str, Any]:
|
||||
@@ -954,6 +1390,7 @@ def run_acceptance(
|
||||
audit_probe=audit_probe,
|
||||
append_sent=True,
|
||||
repeat_send=True,
|
||||
delivery_probe=delivery_probe,
|
||||
)
|
||||
success["provider_verification"] = _wait_for_provider_evidence(
|
||||
client,
|
||||
@@ -989,6 +1426,7 @@ def run_acceptance(
|
||||
audit_probe=audit_probe,
|
||||
append_sent=False,
|
||||
repeat_send=False,
|
||||
delivery_probe=delivery_probe,
|
||||
)
|
||||
if temporary["report"]["send_status_counts"] != {"failed_temporary": 1}:
|
||||
raise AcceptanceError("SMTP connection failure was not classified as temporary")
|
||||
@@ -1014,6 +1452,7 @@ def run_acceptance(
|
||||
audit_probe=audit_probe,
|
||||
append_sent=False,
|
||||
repeat_send=False,
|
||||
delivery_probe=delivery_probe,
|
||||
)
|
||||
if permanent["report"]["send_status_counts"] != {"failed_permanent": 1}:
|
||||
raise AcceptanceError("SMTP authentication failure was not classified as permanent")
|
||||
@@ -1039,11 +1478,150 @@ def run_acceptance(
|
||||
audit_probe=audit_probe,
|
||||
append_sent=True,
|
||||
repeat_send=True,
|
||||
delivery_probe=delivery_probe,
|
||||
)
|
||||
if imap_failure["report"]["imap_status_counts"] != {"failed": 1}:
|
||||
raise AcceptanceError("IMAP authentication failure was not retained as a failed append")
|
||||
drills["imap_authentication_failure"] = imap_failure
|
||||
|
||||
with smtp_fault_endpoint("temporary_data_response") as endpoint:
|
||||
temporary_response_profile = create_mail_profile(
|
||||
client,
|
||||
headers,
|
||||
settings,
|
||||
name="Explicit temporary SMTP response drill",
|
||||
smtp_host=endpoint.host,
|
||||
smtp_port=endpoint.port,
|
||||
)
|
||||
temporary_response, _ = execute_campaign_scenario(
|
||||
client,
|
||||
headers,
|
||||
fixture_path=fixture_path,
|
||||
profile_id=temporary_response_profile,
|
||||
settings=settings,
|
||||
scenario="smtp_temporary_response",
|
||||
snapshot_probe=snapshot_probe,
|
||||
audit_probe=audit_probe,
|
||||
append_sent=False,
|
||||
repeat_send=False,
|
||||
delivery_probe=delivery_probe,
|
||||
)
|
||||
temporary_response["protocol"] = endpoint.evidence()
|
||||
if (
|
||||
temporary_response["report"]["send_status_counts"] != {"failed_temporary": 1}
|
||||
or temporary_response["durable_state"]["attempt_status_counts"]
|
||||
!= {"failed_temporary": 1}
|
||||
or temporary_response["protocol"]
|
||||
!= {
|
||||
"connection_count": 1,
|
||||
"accepted_rcpt_commands": 1,
|
||||
"refused_rcpt_commands": 0,
|
||||
"data_transactions": 1,
|
||||
}
|
||||
):
|
||||
raise AcceptanceError("Explicit SMTP 451 response was not retained as temporary")
|
||||
drills["smtp_temporary_response"] = temporary_response
|
||||
|
||||
with smtp_fault_endpoint("partial_recipient_refusal") as endpoint:
|
||||
partial_profile = create_mail_profile(
|
||||
client,
|
||||
headers,
|
||||
settings,
|
||||
name="Partial SMTP recipient-refusal drill",
|
||||
smtp_host=endpoint.host,
|
||||
smtp_port=endpoint.port,
|
||||
)
|
||||
partial, _ = execute_campaign_scenario(
|
||||
client,
|
||||
headers,
|
||||
fixture_path=fixture_path,
|
||||
profile_id=partial_profile,
|
||||
settings=settings,
|
||||
scenario="partial_envelope_refusal",
|
||||
snapshot_probe=snapshot_probe,
|
||||
audit_probe=audit_probe,
|
||||
append_sent=False,
|
||||
repeat_send=False,
|
||||
delivery_probe=delivery_probe,
|
||||
additional_envelope_recipient=True,
|
||||
)
|
||||
partial["protocol"] = endpoint.evidence()
|
||||
if (
|
||||
partial["report"]["send_status_counts"] != {"smtp_accepted": 1}
|
||||
or partial["durable_state"]["attempt_status_counts"]
|
||||
!= {"smtp_accepted_with_refusals": 1}
|
||||
or partial["protocol"]
|
||||
!= {
|
||||
"connection_count": 1,
|
||||
"accepted_rcpt_commands": 1,
|
||||
"refused_rcpt_commands": 1,
|
||||
"data_transactions": 1,
|
||||
}
|
||||
):
|
||||
raise AcceptanceError("Partial SMTP envelope refusal was not durably distinguished")
|
||||
drills["partial_envelope_refusal"] = partial
|
||||
|
||||
with smtp_fault_endpoint("post_data_disconnect") as endpoint:
|
||||
ambiguous_profile = create_mail_profile(
|
||||
client,
|
||||
headers,
|
||||
settings,
|
||||
name="Post-DATA SMTP ambiguity drill",
|
||||
smtp_host=endpoint.host,
|
||||
smtp_port=endpoint.port,
|
||||
)
|
||||
ambiguous, _ = execute_campaign_scenario(
|
||||
client,
|
||||
headers,
|
||||
fixture_path=fixture_path,
|
||||
profile_id=ambiguous_profile,
|
||||
settings=settings,
|
||||
scenario="post_data_connection_loss",
|
||||
snapshot_probe=snapshot_probe,
|
||||
audit_probe=audit_probe,
|
||||
append_sent=False,
|
||||
repeat_send=False,
|
||||
delivery_probe=delivery_probe,
|
||||
)
|
||||
ambiguous["protocol"] = endpoint.evidence()
|
||||
if (
|
||||
ambiguous["report"]["send_status_counts"] != {"outcome_unknown": 1}
|
||||
or ambiguous["durable_state"]["attempt_status_counts"]
|
||||
!= {"outcome_unknown": 1}
|
||||
or ambiguous["protocol"]
|
||||
!= {
|
||||
"connection_count": 1,
|
||||
"accepted_rcpt_commands": 1,
|
||||
"refused_rcpt_commands": 0,
|
||||
"data_transactions": 1,
|
||||
}
|
||||
):
|
||||
raise AcceptanceError("Post-DATA connection loss was not frozen as outcome unknown")
|
||||
drills["post_data_connection_loss"] = ambiguous
|
||||
|
||||
with smtp_fault_endpoint("post_data_hold") as endpoint:
|
||||
worker_profile = create_mail_profile(
|
||||
client,
|
||||
headers,
|
||||
settings,
|
||||
name="Campaign worker interruption drill",
|
||||
smtp_host=endpoint.host,
|
||||
smtp_port=endpoint.port,
|
||||
)
|
||||
worker_interruption = execute_worker_interruption_scenario(
|
||||
client,
|
||||
headers,
|
||||
fixture_path=fixture_path,
|
||||
profile_id=worker_profile,
|
||||
settings=settings,
|
||||
endpoint=endpoint,
|
||||
snapshot_probe=snapshot_probe,
|
||||
audit_probe=audit_probe,
|
||||
delivery_probe=delivery_probe,
|
||||
worker_job_probe=worker_job_probe,
|
||||
)
|
||||
drills["worker_interruption"] = worker_interruption
|
||||
|
||||
evidence = {
|
||||
"schema_version": EVIDENCE_SCHEMA,
|
||||
"generated_at": datetime.now(timezone.utc).isoformat(),
|
||||
@@ -1062,13 +1640,14 @@ def run_acceptance(
|
||||
"imap_append": True,
|
||||
"repeat_ordinary_send_no_duplicate": True,
|
||||
"smtp_connection_failure": include_failure_drills,
|
||||
"smtp_temporary_response": False,
|
||||
"smtp_temporary_response": include_failure_drills,
|
||||
"smtp_permanent_authentication_failure": include_failure_drills,
|
||||
"partial_envelope_refusal": False,
|
||||
"partial_envelope_refusal": include_failure_drills,
|
||||
"imap_authentication_failure": include_failure_drills,
|
||||
"post_data_connection_loss_outcome_unknown": False,
|
||||
"post_data_connection_loss_outcome_unknown": include_failure_drills,
|
||||
"source_artifact_provenance": False,
|
||||
"worker_restart_interruption": False,
|
||||
"worker_restart_interruption": include_failure_drills,
|
||||
"celery_broker_redelivery": False,
|
||||
},
|
||||
}
|
||||
_assert_evidence_safe(evidence, settings=settings)
|
||||
@@ -1235,6 +1814,48 @@ def _bootstrap_and_run(
|
||||
]
|
||||
return dict(Counter(actions))
|
||||
|
||||
def delivery_probe(campaign_id: str, version_id: str) -> Mapping[str, Any]:
|
||||
from govoplan_campaign.backend.db.models import CampaignJob, SendAttempt
|
||||
|
||||
with database.SessionLocal() as session:
|
||||
jobs = (
|
||||
session.query(CampaignJob)
|
||||
.filter(
|
||||
CampaignJob.campaign_id == campaign_id,
|
||||
CampaignJob.campaign_version_id == version_id,
|
||||
)
|
||||
.all()
|
||||
)
|
||||
job_ids = [job.id for job in jobs]
|
||||
attempts = (
|
||||
session.query(SendAttempt)
|
||||
.filter(SendAttempt.job_id.in_(job_ids))
|
||||
.all()
|
||||
if job_ids
|
||||
else []
|
||||
)
|
||||
return {
|
||||
"job_count": len(jobs),
|
||||
"send_status_counts": dict(Counter(job.send_status for job in jobs)),
|
||||
"attempt_status_counts": dict(Counter(attempt.status for attempt in attempts)),
|
||||
"unfinished_attempt_count": sum(
|
||||
1 for attempt in attempts if attempt.finished_at is None
|
||||
),
|
||||
}
|
||||
|
||||
def worker_job_probe(version_id: str) -> str:
|
||||
from govoplan_campaign.backend.db.models import CampaignJob
|
||||
|
||||
with database.SessionLocal() as session:
|
||||
jobs = (
|
||||
session.query(CampaignJob)
|
||||
.filter(CampaignJob.campaign_version_id == version_id)
|
||||
.all()
|
||||
)
|
||||
if len(jobs) != 1:
|
||||
raise AcceptanceError("Worker-interruption Campaign did not contain one job")
|
||||
return jobs[0].id
|
||||
|
||||
with TestClient(app) as client:
|
||||
login = _expect(
|
||||
client.post(
|
||||
@@ -1265,6 +1886,8 @@ def _bootstrap_and_run(
|
||||
fixture_path=fixture_path,
|
||||
snapshot_probe=snapshot_probe,
|
||||
audit_probe=audit_probe,
|
||||
delivery_probe=delivery_probe,
|
||||
worker_job_probe=worker_job_probe,
|
||||
include_failure_drills=include_failure_drills,
|
||||
module_versions=module_versions,
|
||||
)
|
||||
|
||||
Reference in New Issue
Block a user