315 lines
10 KiB
Python
315 lines
10 KiB
Python
from __future__ import annotations
|
|
|
|
import importlib.util
|
|
import json
|
|
from pathlib import Path
|
|
from types import SimpleNamespace
|
|
import sys
|
|
import tempfile
|
|
from unittest import mock
|
|
|
|
import pytest
|
|
|
|
|
|
REPOSITORY_ROOT = Path(__file__).resolve().parents[1]
|
|
RUNNER_PATH = (
|
|
REPOSITORY_ROOT
|
|
/ "dev"
|
|
/ "mail-testbed"
|
|
/ "run_celery_redelivery_acceptance.py"
|
|
)
|
|
COMPOSE_PATH = REPOSITORY_ROOT / "dev" / "mail-testbed" / "docker-compose.yml"
|
|
FIXTURE_PATH = REPOSITORY_ROOT / "examples" / "greenmail-delivery" / "campaign.json"
|
|
TASK_ID = "12345678-1234-4234-8234-123456789abc"
|
|
|
|
|
|
def _load_runner():
|
|
spec = importlib.util.spec_from_file_location(
|
|
"govoplan_campaign_celery_redelivery_acceptance",
|
|
RUNNER_PATH,
|
|
)
|
|
assert spec is not None and spec.loader is not None
|
|
module = importlib.util.module_from_spec(spec)
|
|
sys.modules[spec.name] = module
|
|
spec.loader.exec_module(module)
|
|
return module
|
|
|
|
|
|
runner = _load_runner()
|
|
|
|
|
|
class _Response:
|
|
def __init__(self, status_code: int, payload: dict) -> None:
|
|
self.status_code = status_code
|
|
self._payload = payload
|
|
|
|
def json(self) -> dict:
|
|
return self._payload
|
|
|
|
|
|
class _Client:
|
|
def post(self, path: str, **_kwargs) -> _Response:
|
|
assert path.endswith("/queue")
|
|
return _Response(
|
|
200,
|
|
{
|
|
"queued_count": 1,
|
|
"skipped_count": 0,
|
|
"blocked_count": 0,
|
|
"enqueued_count": 1,
|
|
"delivery_mode": "worker_queue",
|
|
"worker_queue_available": True,
|
|
"dry_run": False,
|
|
},
|
|
)
|
|
|
|
def get(self, path: str, **_kwargs) -> _Response:
|
|
assert path.endswith("/report")
|
|
return _Response(
|
|
200,
|
|
{
|
|
"cards": {
|
|
"jobs_total": 1,
|
|
"outcome_unknown": 1,
|
|
"needs_attention": 1,
|
|
},
|
|
"status_counts": {
|
|
"send": {"outcome_unknown": 1},
|
|
"imap": {"pending": 1},
|
|
},
|
|
},
|
|
)
|
|
|
|
|
|
class _Endpoint:
|
|
host = "127.0.0.1"
|
|
port = 4025
|
|
|
|
def __init__(self) -> None:
|
|
self.release_count = 0
|
|
|
|
def wait_for_data(self, _timeout_seconds: int) -> bool:
|
|
return True
|
|
|
|
def release_held_connection(self) -> None:
|
|
self.release_count += 1
|
|
|
|
def evidence(self) -> dict[str, int]:
|
|
return {
|
|
"connection_count": 1,
|
|
"accepted_rcpt_commands": 1,
|
|
"refused_rcpt_commands": 0,
|
|
"data_transactions": 1,
|
|
}
|
|
|
|
|
|
def _settings():
|
|
return runner.TestbedSettings(
|
|
smtp_host="127.0.0.1",
|
|
smtp_port=3025,
|
|
imap_host="127.0.0.1",
|
|
imap_port=3143,
|
|
username="campaign-test@govoplan.test",
|
|
password="local-test-password",
|
|
sender="campaign-test@govoplan.test",
|
|
recipient="campaign-test@govoplan.test",
|
|
sent_folder="Sent",
|
|
provider_timeout_seconds=5,
|
|
)
|
|
|
|
|
|
def test_compose_redis_is_isolated_durable_and_health_checked() -> None:
|
|
compose = COMPOSE_PATH.read_text(encoding="utf-8")
|
|
|
|
assert "redis:7-alpine" in compose
|
|
assert '"--appendonly", "yes"' in compose
|
|
assert "127.0.0.1:${GOVOPLAN_CAMPAIGN_TEST_REDIS_PORT:-36379}:6379" in compose
|
|
assert 'test: ["CMD", "redis-cli", "ping"]' in compose
|
|
assert "campaign-redis-data:/data" in compose
|
|
|
|
|
|
def test_runbook_keeps_local_redelivery_distinct_from_production_supervision() -> None:
|
|
testbed = (REPOSITORY_ROOT / "dev" / "mail-testbed" / "README.md").read_text(
|
|
encoding="utf-8"
|
|
)
|
|
runbook = (REPOSITORY_ROOT / "docs" / "CAMPAIGN_DELIVERY_RUNBOOK.md").read_text(
|
|
encoding="utf-8"
|
|
)
|
|
|
|
assert "run_celery_redelivery_acceptance.py" in testbed
|
|
assert "same Celery task identity must be redelivered" in testbed
|
|
assert "production daemon supervision" in testbed
|
|
assert "empty broker queue/unacked set" in runbook
|
|
assert "production process manager" in runbook
|
|
|
|
|
|
def test_compose_lifecycle_targets_only_isolated_redis_service() -> None:
|
|
up = runner._compose_command(
|
|
compose_file=COMPOSE_PATH,
|
|
project_name="govoplan-campaign-redelivery-test",
|
|
operation="up",
|
|
)
|
|
down = runner._compose_command(
|
|
compose_file=COMPOSE_PATH,
|
|
project_name="govoplan-campaign-redelivery-test",
|
|
operation="down",
|
|
)
|
|
|
|
assert up[-3:] == ["up", "--detach", "redis"]
|
|
assert down[-3:] == ["down", "--volumes", "--remove-orphans"]
|
|
assert "greenmail" not in up
|
|
assert "--project-name" in up
|
|
|
|
|
|
def test_worker_bootstrap_uses_real_late_ack_solo_celery_worker() -> None:
|
|
source = runner.WORKER_BOOTSTRAP
|
|
|
|
assert "celery.worker_main" in source
|
|
assert '"--pool=solo"' in source
|
|
assert '"--queues=send_email"' in source
|
|
assert '"visibility_timeout"' in source
|
|
assert '"polling_interval"' in source
|
|
assert "send_email.run" not in source
|
|
|
|
|
|
def test_runtime_root_uses_platform_temp_selection() -> None:
|
|
with mock.patch(
|
|
"govoplan_campaign_celery_redelivery_acceptance.tempfile.mkdtemp",
|
|
return_value="/selected-temp/govoplan-campaign-celery-redelivery-test",
|
|
) as mkdtemp:
|
|
runtime_root = runner._create_runtime_root()
|
|
|
|
assert runtime_root == Path(
|
|
"/selected-temp/govoplan-campaign-celery-redelivery-test"
|
|
)
|
|
mkdtemp.assert_called_once_with(prefix="govoplan-campaign-celery-redelivery-")
|
|
|
|
|
|
def test_worker_log_projection_matches_redelivered_task_without_retaining_id() -> None:
|
|
with tempfile.TemporaryDirectory() as temporary_directory:
|
|
log_path = Path(temporary_directory) / "worker.log"
|
|
log_path.write_text(
|
|
"\n".join(
|
|
[
|
|
f"Task govoplan.campaigns.send_email[{TASK_ID}] received",
|
|
f"Task govoplan.campaigns.send_email[{TASK_ID}] succeeded in 0.1s",
|
|
]
|
|
),
|
|
encoding="utf-8",
|
|
)
|
|
with log_path.open("ab") as handle:
|
|
worker = runner.WorkerProcess(
|
|
process=SimpleNamespace(),
|
|
log_path=log_path,
|
|
log_handle=handle,
|
|
)
|
|
|
|
assert worker.received_task_ids() == (TASK_ID,)
|
|
assert worker.succeeded_task_ids() == (TASK_ID,)
|
|
|
|
|
|
def test_queue_projection_fails_closed_if_no_task_was_published() -> None:
|
|
with pytest.raises(runner.AcceptanceError, match="one Celery task"):
|
|
runner._queue_evidence(
|
|
{
|
|
"queued_count": 1,
|
|
"skipped_count": 0,
|
|
"blocked_count": 0,
|
|
"enqueued_count": 0,
|
|
"delivery_mode": "database_queue",
|
|
"worker_queue_available": False,
|
|
"dry_run": False,
|
|
}
|
|
)
|
|
|
|
|
|
def test_redelivery_orchestration_requires_same_task_and_no_second_smtp_effect(
|
|
monkeypatch,
|
|
) -> None:
|
|
first_worker = mock.Mock()
|
|
first_worker.received_task_ids.return_value = (TASK_ID,)
|
|
replacement_worker = mock.Mock()
|
|
replacement_worker.received_task_ids.return_value = (TASK_ID,)
|
|
workers = iter([first_worker, replacement_worker])
|
|
endpoint = _Endpoint()
|
|
durable_states = iter(
|
|
[
|
|
{
|
|
"job_count": 1,
|
|
"send_status_counts": {"sending": 1},
|
|
"attempt_status_counts": {"smtp_in_progress": 1},
|
|
"unfinished_attempt_count": 1,
|
|
},
|
|
{
|
|
"job_count": 1,
|
|
"send_status_counts": {"outcome_unknown": 1},
|
|
"attempt_status_counts": {"outcome_unknown": 1},
|
|
"unfinished_attempt_count": 0,
|
|
},
|
|
]
|
|
)
|
|
prepared = SimpleNamespace(
|
|
campaign_id="campaign-internal",
|
|
version_id="version-internal",
|
|
public_evidence=lambda: {
|
|
"validation": {"ok": True},
|
|
"build": {"built_count": 1},
|
|
"campaign_mail_boundary": {
|
|
"profile_reference_only": True,
|
|
"smtp_revision_frozen": True,
|
|
"imap_revision_frozen": True,
|
|
"resolved_transport_material_present": False,
|
|
},
|
|
},
|
|
)
|
|
|
|
monkeypatch.setattr(runner, "create_mail_profile", lambda *args, **kwargs: "profile-internal")
|
|
monkeypatch.setattr(runner, "prepare_campaign_scenario", lambda *args, **kwargs: prepared)
|
|
monkeypatch.setattr(runner, "_start_worker", lambda *args, **kwargs: next(workers))
|
|
monkeypatch.setattr(runner, "_wait_for_worker_ready", lambda *args, **kwargs: None)
|
|
received = iter([TASK_ID, TASK_ID])
|
|
monkeypatch.setattr(runner, "_wait_for_received_task", lambda *args, **kwargs: next(received))
|
|
monkeypatch.setattr(runner, "_wait_for_task_success", lambda *args, **kwargs: None)
|
|
monkeypatch.setattr(runner, "_kill_worker", lambda *args, **kwargs: -9)
|
|
monkeypatch.setattr(runner, "_stop_worker", lambda *args, **kwargs: None)
|
|
monkeypatch.setattr(
|
|
runner,
|
|
"_wait_for_broker_drained",
|
|
lambda *args, **kwargs: runner.RedisBrokerState(0, 0, 0),
|
|
)
|
|
|
|
evidence = runner.execute_redelivery_scenario(
|
|
_Client(),
|
|
{"Authorization": "not-retained"},
|
|
fixture_path=FIXTURE_PATH,
|
|
settings=_settings(),
|
|
endpoint=endpoint,
|
|
redis_url="redis://127.0.0.1:36379/0",
|
|
runtime_root=Path("/not-used"),
|
|
snapshot_probe=lambda _version_id: ({}, {}),
|
|
audit_probe=lambda _campaign_id, _version_id: {
|
|
"campaign.created": 1,
|
|
"campaign.validated": 1,
|
|
"campaign.messages_built": 1,
|
|
"campaign.queued": 1,
|
|
},
|
|
delivery_probe=lambda _campaign_id, _version_id: next(durable_states),
|
|
)
|
|
|
|
assert evidence["broker"] == {
|
|
"transport": "redis",
|
|
"same_task_identity_redelivered": True,
|
|
"first_worker_received_count": 1,
|
|
"replacement_worker_received_count": 1,
|
|
"queue_depth": 0,
|
|
"unacked_hash_count": 0,
|
|
"unacked_index_count": 0,
|
|
}
|
|
assert evidence["protocol"]["connection_count"] == 1
|
|
assert evidence["protocol"]["data_transactions"] == 1
|
|
assert evidence["recovered_durable_state"]["send_status_counts"] == {
|
|
"outcome_unknown": 1
|
|
}
|
|
assert TASK_ID not in json.dumps(evidence, sort_keys=True)
|
|
assert endpoint.release_count >= 1
|