Compare commits
11 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 4774a8025c | |||
| 689dc1fd6b | |||
| 90677348ff | |||
| 06786e86ef | |||
| 6fda123fc3 | |||
| 2fa91bb943 | |||
| e3cc476508 | |||
| 01ef541917 | |||
| d567257311 | |||
| 3075ef7f5b | |||
| 1ab4e91ffd |
@@ -9,3 +9,5 @@ GOVOPLAN_MAIL_TEST_IMAP_PORT=3143
|
|||||||
GOVOPLAN_MAIL_TEST_SENT_FOLDER=Sent
|
GOVOPLAN_MAIL_TEST_SENT_FOLDER=Sent
|
||||||
GOVOPLAN_MAIL_TEST_ZIP_PASSWORD=zip-test-password
|
GOVOPLAN_MAIL_TEST_ZIP_PASSWORD=zip-test-password
|
||||||
GOVOPLAN_MAIL_TEST_READY_TIMEOUT_SECONDS=45
|
GOVOPLAN_MAIL_TEST_READY_TIMEOUT_SECONDS=45
|
||||||
|
GOVOPLAN_CAMPAIGN_TEST_REDIS_PORT=36379
|
||||||
|
GOVOPLAN_CAMPAIGN_TEST_REDIS_VISIBILITY_TIMEOUT_SECONDS=3
|
||||||
|
|||||||
@@ -44,6 +44,111 @@ If the smoke is started immediately after `docker compose up -d`, GreenMail may
|
|||||||
bind the SMTP/IMAP ports before the services are fully ready. The smoke retries
|
bind the SMTP/IMAP ports before the services are fully ready. The smoke retries
|
||||||
login and folder setup for `GOVOPLAN_MAIL_TEST_READY_TIMEOUT_SECONDS`.
|
login and folder setup for `GOVOPLAN_MAIL_TEST_READY_TIMEOUT_SECONDS`.
|
||||||
|
|
||||||
|
## Campaign Acceptance
|
||||||
|
|
||||||
|
The transport smoke proves the Mail adapters. The Campaign acceptance runner
|
||||||
|
proves the public composition: it creates an isolated temporary Core database,
|
||||||
|
creates a Mail-owned encrypted profile through the API, materializes the
|
||||||
|
credential-free [`greenmail-delivery`](../../examples/greenmail-delivery/campaign.json)
|
||||||
|
fixture with only that profile reference, validates/builds it, sends the exact
|
||||||
|
generated EML through Campaign once, appends it once, and cross-checks Campaign
|
||||||
|
report/audit state with one unique-subject message in the GreenMail INBOX and
|
||||||
|
Sent folders. Provider mailbox verification is not a byte-for-byte comparison
|
||||||
|
after provider-side header or storage transformations.
|
||||||
|
|
||||||
|
```bash
|
||||||
|
cd /mnt/DATA/git/govoplan-campaign/dev/mail-testbed
|
||||||
|
set -a
|
||||||
|
. ./.env
|
||||||
|
set +a
|
||||||
|
/mnt/DATA/git/govoplan/.venv/bin/python run_campaign_acceptance.py \
|
||||||
|
--evidence /tmp/govoplan-campaign-greenmail-evidence.json
|
||||||
|
```
|
||||||
|
|
||||||
|
The default run also uses controlled loopback protocol endpoints to prove that
|
||||||
|
an SMTP connection loss before transmission is temporary, an explicit SMTP
|
||||||
|
authentication rejection is permanent, a final `451` response after DATA is
|
||||||
|
temporary, one accepted and one refused RCPT command is retained as partial
|
||||||
|
envelope acceptance, a connection loss after complete DATA is frozen as
|
||||||
|
`outcome_unknown`, and an IMAP authentication rejection after SMTP acceptance
|
||||||
|
leaves the send accepted while the append fails. A
|
||||||
|
second ordinary send must be rejected before another provider effect.
|
||||||
|
|
||||||
|
The worker drill queues one job, starts the registered
|
||||||
|
`govoplan.campaigns.send_email` task body in a dedicated OS process, waits until
|
||||||
|
the controlled endpoint has received complete DATA, terminates that process,
|
||||||
|
and starts the same task body in a fresh process. It proves the durable
|
||||||
|
`sending`/unfinished-attempt boundary is recovered as `outcome_unknown`
|
||||||
|
without a second SMTP connection or DATA transaction. This is a real process
|
||||||
|
and task-boundary interruption, but it does not start a Celery daemon, Redis
|
||||||
|
broker, or broker redelivery; `celery_broker_redelivery` therefore remains
|
||||||
|
`false` in the evidence.
|
||||||
|
|
||||||
|
The bounded JSON contains no endpoint, account, address, credential, profile,
|
||||||
|
campaign, version, or job identifiers. It records module versions, the fixture
|
||||||
|
hash, normalized classifications/counts, the Mail-profile boundary, required
|
||||||
|
audit actions, provider mailbox increments, and coverage flags. Runtime version
|
||||||
|
declarations identify the exercised composition; they do not claim that the
|
||||||
|
sources are clean, tagged, signed, or release-provenanced. The evidence names
|
||||||
|
them `declared_module_versions` and keeps `source_artifact_provenance` false;
|
||||||
|
exact commit/artifact provenance belongs to the package and release gate.
|
||||||
|
|
||||||
|
This runner is restricted to literal loopback IP addresses and the synchronous
|
||||||
|
Campaign delivery mode. Hostnames such as `localhost` and every non-loopback
|
||||||
|
address fail before profile creation, avoiding a DNS change between validation
|
||||||
|
and connection. It is local target-like evidence, not approval of an
|
||||||
|
institution's SMTP/IMAP service. The controlled post-DATA, temporary-response,
|
||||||
|
partial-refusal, and task-process interruption drills are local effect-level
|
||||||
|
proof, not proof of a target provider's behavior or Redis/Celery broker
|
||||||
|
redelivery. Use `--success-only` only when testing the success journey without
|
||||||
|
the local failure endpoints.
|
||||||
|
|
||||||
|
## Redis/Celery Redelivery Acceptance
|
||||||
|
|
||||||
|
Run the maintained broker/worker-loss acceptance separately from the GreenMail
|
||||||
|
journey:
|
||||||
|
|
||||||
|
```bash
|
||||||
|
cd /mnt/DATA/git/govoplan-campaign/dev/mail-testbed
|
||||||
|
set -a
|
||||||
|
. ./.env
|
||||||
|
set +a
|
||||||
|
/mnt/DATA/git/govoplan/.venv/bin/python run_celery_redelivery_acceptance.py \
|
||||||
|
--evidence /tmp/govoplan-campaign-celery-redelivery-evidence.json
|
||||||
|
```
|
||||||
|
|
||||||
|
The runner creates a unique Compose project, starts only its loopback-bound,
|
||||||
|
AOF-enabled Redis service, creates an isolated temporary GovOPlaN database,
|
||||||
|
and starts a real Celery worker subscribed to `send_email`. A controlled SMTP
|
||||||
|
server holds the transaction after complete DATA and before the final response.
|
||||||
|
The runner kills that solo worker with the task still unacknowledged and starts
|
||||||
|
a replacement worker. After the configured Redis visibility timeout, the same Celery task identity must be redelivered.
|
||||||
|
The replacement must turn the durable
|
||||||
|
unfinished attempt into `outcome_unknown`, acknowledge the task, drain the
|
||||||
|
broker queue/unacked records, and leave the SMTP endpoint at exactly one
|
||||||
|
connection and one DATA transaction.
|
||||||
|
|
||||||
|
Only bounded counts, classifications, and booleans are retained. Worker logs,
|
||||||
|
task IDs, database identifiers, endpoints, credentials, and raw diagnostics are
|
||||||
|
kept in the temporary runtime and deleted. The evidence proves the local Redis
|
||||||
|
transport, real Celery process boundary, runner-supervised replacement, and
|
||||||
|
Campaign's duplicate-effect guard. It deliberately keeps production daemon supervision,
|
||||||
|
target-provider behavior, and source-artifact provenance false.
|
||||||
|
It does not claim that systemd, Kubernetes, another container orchestrator, or
|
||||||
|
an institution's Redis/SMTP deployment behaves identically.
|
||||||
|
|
||||||
|
The default run requires Docker CLI/Compose/daemon access and permission to
|
||||||
|
pull `redis:7-alpine`; the Celery workers execute from the current Python
|
||||||
|
environment. The isolated Compose project and volume are removed on exit.
|
||||||
|
`GOVOPLAN_CAMPAIGN_TEST_REDIS_VISIBILITY_TIMEOUT_SECONDS` defaults to three
|
||||||
|
seconds only to make this destructive local drill finish promptly; it is not a
|
||||||
|
production recommendation.
|
||||||
|
|
||||||
|
Starting the maintained test bed requires a working Docker CLI, Compose plugin,
|
||||||
|
daemon/socket access, and permission to pull `greenmail/standalone:2.1.9`. The
|
||||||
|
Campaign runner needs only the already-running loopback endpoints; it neither
|
||||||
|
starts Docker nor claims that it did.
|
||||||
|
|
||||||
## Use With A Campaign
|
## Use With A Campaign
|
||||||
|
|
||||||
Use the same settings in a campaign mail profile:
|
Use the same settings in a campaign mail profile:
|
||||||
|
|||||||
@@ -15,3 +15,19 @@ services:
|
|||||||
- "${GOVOPLAN_MAIL_TEST_SMTP_PORT:-3025}:3025"
|
- "${GOVOPLAN_MAIL_TEST_SMTP_PORT:-3025}:3025"
|
||||||
- "${GOVOPLAN_MAIL_TEST_IMAP_PORT:-3143}:3143"
|
- "${GOVOPLAN_MAIL_TEST_IMAP_PORT:-3143}:3143"
|
||||||
- "127.0.0.1:38080:8080"
|
- "127.0.0.1:38080:8080"
|
||||||
|
|
||||||
|
redis:
|
||||||
|
image: redis:7-alpine
|
||||||
|
command: ["redis-server", "--appendonly", "yes"]
|
||||||
|
ports:
|
||||||
|
- "127.0.0.1:${GOVOPLAN_CAMPAIGN_TEST_REDIS_PORT:-36379}:6379"
|
||||||
|
healthcheck:
|
||||||
|
test: ["CMD", "redis-cli", "ping"]
|
||||||
|
interval: 1s
|
||||||
|
timeout: 1s
|
||||||
|
retries: 30
|
||||||
|
volumes:
|
||||||
|
- campaign-redis-data:/data
|
||||||
|
|
||||||
|
volumes:
|
||||||
|
campaign-redis-data:
|
||||||
|
|||||||
1935
dev/mail-testbed/run_campaign_acceptance.py
Normal file
1935
dev/mail-testbed/run_campaign_acceptance.py
Normal file
File diff suppressed because it is too large
Load Diff
856
dev/mail-testbed/run_celery_redelivery_acceptance.py
Normal file
856
dev/mail-testbed/run_celery_redelivery_acceptance.py
Normal file
@@ -0,0 +1,856 @@
|
|||||||
|
#!/usr/bin/env python3
|
||||||
|
"""Prove Redis/Celery redelivery does not repeat an ambiguous SMTP effect.
|
||||||
|
|
||||||
|
The default run starts an isolated Redis Compose service, two successive real
|
||||||
|
Celery worker processes, and a controlled loopback SMTP endpoint. It kills the
|
||||||
|
first worker after complete DATA but before a final SMTP response. The same
|
||||||
|
unacknowledged broker task must be delivered to the replacement worker, which
|
||||||
|
must freeze the unfinished durable attempt as ``outcome_unknown`` without a
|
||||||
|
second SMTP connection or DATA transaction.
|
||||||
|
"""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import argparse
|
||||||
|
from collections import Counter
|
||||||
|
from contextlib import contextmanager
|
||||||
|
from dataclasses import dataclass, replace
|
||||||
|
from datetime import datetime, timezone
|
||||||
|
import hashlib
|
||||||
|
import json
|
||||||
|
import os
|
||||||
|
from pathlib import Path
|
||||||
|
import re
|
||||||
|
import shutil
|
||||||
|
import socket
|
||||||
|
import subprocess
|
||||||
|
import sys
|
||||||
|
import tempfile
|
||||||
|
import time
|
||||||
|
from typing import Any, Callable, Iterator, Mapping
|
||||||
|
from uuid import uuid4
|
||||||
|
|
||||||
|
from redis import Redis
|
||||||
|
from redis.exceptions import RedisError
|
||||||
|
|
||||||
|
|
||||||
|
SCRIPT_ROOT = Path(__file__).resolve().parent
|
||||||
|
REPOSITORY_ROOT = SCRIPT_ROOT.parents[1]
|
||||||
|
if str(SCRIPT_ROOT) not in sys.path:
|
||||||
|
sys.path.insert(0, str(SCRIPT_ROOT))
|
||||||
|
|
||||||
|
from run_campaign_acceptance import ( # noqa: E402
|
||||||
|
AcceptanceError,
|
||||||
|
DEFAULT_FIXTURE,
|
||||||
|
EXPECTED_AUDIT_ACTIONS,
|
||||||
|
TestbedSettings,
|
||||||
|
_assert_evidence_safe,
|
||||||
|
_core_package_version,
|
||||||
|
_durable_state_evidence,
|
||||||
|
_expect,
|
||||||
|
_report_evidence,
|
||||||
|
create_mail_profile,
|
||||||
|
prepare_campaign_scenario,
|
||||||
|
required_composition_versions,
|
||||||
|
smtp_fault_endpoint,
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
EVIDENCE_SCHEMA = "govoplan.campaign.celery-redelivery-acceptance.v1"
|
||||||
|
MAX_EVIDENCE_BYTES = 128 * 1024
|
||||||
|
DEFAULT_COMPOSE_FILE = SCRIPT_ROOT / "docker-compose.yml"
|
||||||
|
TASK_RECEIVED_PATTERN = re.compile(
|
||||||
|
r"Task govoplan[.]campaigns[.]send_email\[([0-9a-f-]{36})\] received",
|
||||||
|
re.IGNORECASE,
|
||||||
|
)
|
||||||
|
TASK_SUCCEEDED_PATTERN = re.compile(
|
||||||
|
r"Task govoplan[.]campaigns[.]send_email\[([0-9a-f-]{36})\] succeeded",
|
||||||
|
re.IGNORECASE,
|
||||||
|
)
|
||||||
|
WORKER_BOOTSTRAP = r"""
|
||||||
|
import os
|
||||||
|
import sys
|
||||||
|
|
||||||
|
from govoplan_core.celery_app import celery
|
||||||
|
|
||||||
|
visibility_timeout = int(os.environ["GOVOPLAN_CAMPAIGN_TEST_REDIS_VISIBILITY_TIMEOUT_SECONDS"])
|
||||||
|
celery.conf.broker_transport_options = {
|
||||||
|
**dict(celery.conf.broker_transport_options or {}),
|
||||||
|
"polling_interval": 0.25,
|
||||||
|
"visibility_timeout": visibility_timeout,
|
||||||
|
}
|
||||||
|
celery.worker_main(
|
||||||
|
[
|
||||||
|
"worker",
|
||||||
|
"--loglevel=INFO",
|
||||||
|
"--pool=solo",
|
||||||
|
"--concurrency=1",
|
||||||
|
"--queues=send_email",
|
||||||
|
f"--hostname={sys.argv[1]}@%h",
|
||||||
|
"--without-gossip",
|
||||||
|
"--without-mingle",
|
||||||
|
"--without-heartbeat",
|
||||||
|
]
|
||||||
|
)
|
||||||
|
"""
|
||||||
|
|
||||||
|
|
||||||
|
@dataclass(slots=True)
|
||||||
|
class WorkerProcess:
|
||||||
|
process: subprocess.Popen[bytes]
|
||||||
|
log_path: Path
|
||||||
|
log_handle: Any
|
||||||
|
|
||||||
|
def text(self) -> str:
|
||||||
|
self.log_handle.flush()
|
||||||
|
try:
|
||||||
|
return self.log_path.read_text(encoding="utf-8", errors="replace")
|
||||||
|
except OSError as exc:
|
||||||
|
raise AcceptanceError("Celery worker evidence log could not be read") from exc
|
||||||
|
|
||||||
|
def received_task_ids(self) -> tuple[str, ...]:
|
||||||
|
return tuple(TASK_RECEIVED_PATTERN.findall(self.text()))
|
||||||
|
|
||||||
|
def succeeded_task_ids(self) -> tuple[str, ...]:
|
||||||
|
return tuple(TASK_SUCCEEDED_PATTERN.findall(self.text()))
|
||||||
|
|
||||||
|
|
||||||
|
@dataclass(frozen=True, slots=True)
|
||||||
|
class RedisBrokerState:
|
||||||
|
queue_depth: int
|
||||||
|
unacked_hash_count: int
|
||||||
|
unacked_index_count: int
|
||||||
|
|
||||||
|
def as_dict(self) -> dict[str, int]:
|
||||||
|
return {
|
||||||
|
"queue_depth": self.queue_depth,
|
||||||
|
"unacked_hash_count": self.unacked_hash_count,
|
||||||
|
"unacked_index_count": self.unacked_index_count,
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
|
def _positive_int(value: str, *, label: str) -> int:
|
||||||
|
try:
|
||||||
|
parsed = int(value)
|
||||||
|
except ValueError as exc:
|
||||||
|
raise argparse.ArgumentTypeError(f"{label} must be a positive integer") from exc
|
||||||
|
if parsed <= 0:
|
||||||
|
raise argparse.ArgumentTypeError(f"{label} must be a positive integer")
|
||||||
|
return parsed
|
||||||
|
|
||||||
|
|
||||||
|
def _unused_loopback_port() -> int:
|
||||||
|
with socket.socket(socket.AF_INET, socket.SOCK_STREAM) as probe:
|
||||||
|
probe.bind(("127.0.0.1", 0))
|
||||||
|
return int(probe.getsockname()[1])
|
||||||
|
|
||||||
|
|
||||||
|
def _compose_command(
|
||||||
|
*, compose_file: Path, project_name: str, operation: str
|
||||||
|
) -> list[str]:
|
||||||
|
prefix = [
|
||||||
|
"docker",
|
||||||
|
"compose",
|
||||||
|
"--file",
|
||||||
|
str(compose_file),
|
||||||
|
"--project-name",
|
||||||
|
project_name,
|
||||||
|
]
|
||||||
|
if operation == "up":
|
||||||
|
return [*prefix, "up", "--detach", "redis"]
|
||||||
|
if operation == "down":
|
||||||
|
return [*prefix, "down", "--volumes", "--remove-orphans"]
|
||||||
|
raise AcceptanceError("Unsupported Redis Compose operation")
|
||||||
|
|
||||||
|
|
||||||
|
def _run_compose(
|
||||||
|
command: list[str],
|
||||||
|
*,
|
||||||
|
environment: Mapping[str, str],
|
||||||
|
timeout_seconds: int,
|
||||||
|
) -> None:
|
||||||
|
try:
|
||||||
|
completed = subprocess.run(
|
||||||
|
command,
|
||||||
|
env=dict(environment),
|
||||||
|
stdin=subprocess.DEVNULL,
|
||||||
|
stdout=subprocess.PIPE,
|
||||||
|
stderr=subprocess.PIPE,
|
||||||
|
timeout=timeout_seconds,
|
||||||
|
check=False,
|
||||||
|
)
|
||||||
|
except (OSError, subprocess.TimeoutExpired) as exc:
|
||||||
|
raise AcceptanceError("Redis Compose lifecycle command failed") from exc
|
||||||
|
if completed.returncode != 0:
|
||||||
|
raise AcceptanceError("Redis Compose lifecycle command failed")
|
||||||
|
|
||||||
|
|
||||||
|
def _wait_for_redis(redis_url: str, *, timeout_seconds: int) -> None:
|
||||||
|
deadline = time.monotonic() + timeout_seconds
|
||||||
|
client = Redis.from_url(
|
||||||
|
redis_url,
|
||||||
|
socket_connect_timeout=1,
|
||||||
|
socket_timeout=1,
|
||||||
|
decode_responses=False,
|
||||||
|
)
|
||||||
|
try:
|
||||||
|
while time.monotonic() < deadline:
|
||||||
|
try:
|
||||||
|
if client.ping() is True:
|
||||||
|
return
|
||||||
|
except RedisError:
|
||||||
|
pass
|
||||||
|
time.sleep(0.25)
|
||||||
|
finally:
|
||||||
|
client.close()
|
||||||
|
raise AcceptanceError("Isolated Redis broker did not become ready")
|
||||||
|
|
||||||
|
|
||||||
|
@contextmanager
|
||||||
|
def isolated_redis_broker(
|
||||||
|
*,
|
||||||
|
compose_file: Path,
|
||||||
|
timeout_seconds: int,
|
||||||
|
requested_port: int | None = None,
|
||||||
|
) -> Iterator[str]:
|
||||||
|
if shutil.which("docker") is None:
|
||||||
|
raise AcceptanceError("Docker CLI is required to start the isolated Redis broker")
|
||||||
|
if not compose_file.is_file():
|
||||||
|
raise AcceptanceError("Redis Compose definition is unavailable")
|
||||||
|
port = requested_port or _unused_loopback_port()
|
||||||
|
if port <= 0 or port > 65_535:
|
||||||
|
raise AcceptanceError("Redis test port is invalid")
|
||||||
|
project_name = f"govoplan-campaign-redelivery-{uuid4().hex[:12]}"
|
||||||
|
environment = {
|
||||||
|
**os.environ,
|
||||||
|
"GOVOPLAN_CAMPAIGN_TEST_REDIS_PORT": str(port),
|
||||||
|
}
|
||||||
|
lifecycle_attempted = False
|
||||||
|
try:
|
||||||
|
lifecycle_attempted = True
|
||||||
|
_run_compose(
|
||||||
|
_compose_command(
|
||||||
|
compose_file=compose_file,
|
||||||
|
project_name=project_name,
|
||||||
|
operation="up",
|
||||||
|
),
|
||||||
|
environment=environment,
|
||||||
|
timeout_seconds=timeout_seconds,
|
||||||
|
)
|
||||||
|
redis_url = f"redis://127.0.0.1:{port}/0"
|
||||||
|
_wait_for_redis(redis_url, timeout_seconds=timeout_seconds)
|
||||||
|
yield redis_url
|
||||||
|
finally:
|
||||||
|
if lifecycle_attempted:
|
||||||
|
_run_compose(
|
||||||
|
_compose_command(
|
||||||
|
compose_file=compose_file,
|
||||||
|
project_name=project_name,
|
||||||
|
operation="down",
|
||||||
|
),
|
||||||
|
environment=environment,
|
||||||
|
timeout_seconds=timeout_seconds,
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def _start_worker(runtime_root: Path, *, label: str) -> WorkerProcess:
|
||||||
|
log_path = runtime_root / f"{label}.log"
|
||||||
|
log_handle = log_path.open("wb")
|
||||||
|
environment = {**os.environ, "PYTHONUNBUFFERED": "1"}
|
||||||
|
try:
|
||||||
|
process = subprocess.Popen(
|
||||||
|
[sys.executable, "-c", WORKER_BOOTSTRAP, label],
|
||||||
|
env=environment,
|
||||||
|
stdin=subprocess.DEVNULL,
|
||||||
|
stdout=log_handle,
|
||||||
|
stderr=subprocess.STDOUT,
|
||||||
|
close_fds=True,
|
||||||
|
)
|
||||||
|
except Exception:
|
||||||
|
log_handle.close()
|
||||||
|
raise
|
||||||
|
return WorkerProcess(process=process, log_path=log_path, log_handle=log_handle)
|
||||||
|
|
||||||
|
|
||||||
|
def _wait_for_worker_ready(worker: WorkerProcess, *, timeout_seconds: int) -> None:
|
||||||
|
deadline = time.monotonic() + timeout_seconds
|
||||||
|
while time.monotonic() < deadline:
|
||||||
|
if worker.process.poll() is not None:
|
||||||
|
raise AcceptanceError("Celery worker exited before becoming ready")
|
||||||
|
if " ready." in worker.text():
|
||||||
|
return
|
||||||
|
time.sleep(0.2)
|
||||||
|
raise AcceptanceError("Celery worker did not become ready")
|
||||||
|
|
||||||
|
|
||||||
|
def _wait_for_received_task(
|
||||||
|
worker: WorkerProcess,
|
||||||
|
*,
|
||||||
|
timeout_seconds: int,
|
||||||
|
expected_task_id: str | None = None,
|
||||||
|
) -> str:
|
||||||
|
deadline = time.monotonic() + timeout_seconds
|
||||||
|
while time.monotonic() < deadline:
|
||||||
|
received = worker.received_task_ids()
|
||||||
|
if received:
|
||||||
|
if len(set(received)) != 1:
|
||||||
|
raise AcceptanceError("Celery worker received more than one task identity")
|
||||||
|
task_id = received[0]
|
||||||
|
if expected_task_id is not None and task_id != expected_task_id:
|
||||||
|
raise AcceptanceError("Replacement worker received a different broker task")
|
||||||
|
return task_id
|
||||||
|
if worker.process.poll() is not None:
|
||||||
|
raise AcceptanceError("Celery worker exited before receiving the task")
|
||||||
|
time.sleep(0.2)
|
||||||
|
raise AcceptanceError("Celery worker did not receive the broker task")
|
||||||
|
|
||||||
|
|
||||||
|
def _wait_for_task_success(
|
||||||
|
worker: WorkerProcess,
|
||||||
|
*,
|
||||||
|
task_id: str,
|
||||||
|
timeout_seconds: int,
|
||||||
|
) -> None:
|
||||||
|
deadline = time.monotonic() + timeout_seconds
|
||||||
|
while time.monotonic() < deadline:
|
||||||
|
succeeded = worker.succeeded_task_ids()
|
||||||
|
if task_id in succeeded:
|
||||||
|
return
|
||||||
|
if worker.process.poll() is not None:
|
||||||
|
raise AcceptanceError("Replacement Celery worker exited before task success")
|
||||||
|
time.sleep(0.2)
|
||||||
|
raise AcceptanceError("Redelivered Celery task did not complete")
|
||||||
|
|
||||||
|
|
||||||
|
def _kill_worker(worker: WorkerProcess, *, timeout_seconds: int) -> int:
|
||||||
|
if worker.process.poll() is not None:
|
||||||
|
raise AcceptanceError("Celery worker exited before controlled termination")
|
||||||
|
worker.process.kill()
|
||||||
|
try:
|
||||||
|
return_code = worker.process.wait(timeout=timeout_seconds)
|
||||||
|
except subprocess.TimeoutExpired as exc:
|
||||||
|
raise AcceptanceError("Celery worker could not be killed") from exc
|
||||||
|
if return_code == 0:
|
||||||
|
raise AcceptanceError("Celery worker termination was not forced")
|
||||||
|
return return_code
|
||||||
|
|
||||||
|
|
||||||
|
def _stop_worker(worker: WorkerProcess, *, timeout_seconds: int) -> None:
|
||||||
|
if worker.process.poll() is None:
|
||||||
|
worker.process.terminate()
|
||||||
|
try:
|
||||||
|
worker.process.wait(timeout=timeout_seconds)
|
||||||
|
except subprocess.TimeoutExpired:
|
||||||
|
worker.process.kill()
|
||||||
|
worker.process.wait(timeout=timeout_seconds)
|
||||||
|
worker.log_handle.close()
|
||||||
|
|
||||||
|
|
||||||
|
def _broker_state(redis_url: str) -> RedisBrokerState:
|
||||||
|
client = Redis.from_url(
|
||||||
|
redis_url,
|
||||||
|
socket_connect_timeout=2,
|
||||||
|
socket_timeout=2,
|
||||||
|
decode_responses=False,
|
||||||
|
)
|
||||||
|
try:
|
||||||
|
return RedisBrokerState(
|
||||||
|
queue_depth=int(client.llen("send_email")),
|
||||||
|
unacked_hash_count=int(client.hlen("unacked")),
|
||||||
|
unacked_index_count=int(client.zcard("unacked_index")),
|
||||||
|
)
|
||||||
|
except RedisError as exc:
|
||||||
|
raise AcceptanceError("Redis broker state could not be inspected") from exc
|
||||||
|
finally:
|
||||||
|
client.close()
|
||||||
|
|
||||||
|
|
||||||
|
def _wait_for_broker_drained(
|
||||||
|
redis_url: str,
|
||||||
|
*,
|
||||||
|
timeout_seconds: int,
|
||||||
|
) -> RedisBrokerState:
|
||||||
|
deadline = time.monotonic() + timeout_seconds
|
||||||
|
last = RedisBrokerState(0, 0, 0)
|
||||||
|
while time.monotonic() < deadline:
|
||||||
|
last = _broker_state(redis_url)
|
||||||
|
if last == RedisBrokerState(0, 0, 0):
|
||||||
|
return last
|
||||||
|
time.sleep(0.2)
|
||||||
|
raise AcceptanceError("Redis broker retained delivery state after recovery")
|
||||||
|
|
||||||
|
|
||||||
|
def _queue_evidence(payload: Mapping[str, Any]) -> dict[str, Any]:
|
||||||
|
expected = {
|
||||||
|
"queued_count": 1,
|
||||||
|
"skipped_count": 0,
|
||||||
|
"blocked_count": 0,
|
||||||
|
"enqueued_count": 1,
|
||||||
|
"delivery_mode": "worker_queue",
|
||||||
|
"worker_queue_available": True,
|
||||||
|
"dry_run": False,
|
||||||
|
}
|
||||||
|
evidence = {key: payload.get(key) for key in expected}
|
||||||
|
if evidence != expected:
|
||||||
|
raise AcceptanceError("Campaign was not durably queued to one Celery task")
|
||||||
|
return evidence
|
||||||
|
|
||||||
|
|
||||||
|
def execute_redelivery_scenario(
|
||||||
|
client: Any,
|
||||||
|
headers: Mapping[str, str],
|
||||||
|
*,
|
||||||
|
fixture_path: Path,
|
||||||
|
settings: TestbedSettings,
|
||||||
|
endpoint: Any,
|
||||||
|
redis_url: str,
|
||||||
|
runtime_root: 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]],
|
||||||
|
) -> dict[str, Any]:
|
||||||
|
profile_id = create_mail_profile(
|
||||||
|
client,
|
||||||
|
headers,
|
||||||
|
settings,
|
||||||
|
name="Campaign Redis Celery redelivery drill",
|
||||||
|
smtp_host=endpoint.host,
|
||||||
|
smtp_port=endpoint.port,
|
||||||
|
)
|
||||||
|
prepared = prepare_campaign_scenario(
|
||||||
|
client,
|
||||||
|
headers,
|
||||||
|
fixture_path=fixture_path,
|
||||||
|
profile_id=profile_id,
|
||||||
|
settings=settings,
|
||||||
|
scenario="celery_broker_redelivery",
|
||||||
|
snapshot_probe=snapshot_probe,
|
||||||
|
)
|
||||||
|
|
||||||
|
first_worker = _start_worker(runtime_root, label="first-worker")
|
||||||
|
replacement_worker: WorkerProcess | None = None
|
||||||
|
try:
|
||||||
|
_wait_for_worker_ready(
|
||||||
|
first_worker,
|
||||||
|
timeout_seconds=settings.provider_timeout_seconds,
|
||||||
|
)
|
||||||
|
queued = _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": True,
|
||||||
|
"dry_run": False,
|
||||||
|
},
|
||||||
|
),
|
||||||
|
200,
|
||||||
|
"Celery-redelivery Campaign queue",
|
||||||
|
)
|
||||||
|
queue_evidence = _queue_evidence(queued)
|
||||||
|
first_task_id = _wait_for_received_task(
|
||||||
|
first_worker,
|
||||||
|
timeout_seconds=settings.provider_timeout_seconds,
|
||||||
|
)
|
||||||
|
if not endpoint.wait_for_data(settings.provider_timeout_seconds):
|
||||||
|
raise AcceptanceError("Celery worker did not reach complete SMTP DATA")
|
||||||
|
first_exit_code = _kill_worker(
|
||||||
|
first_worker,
|
||||||
|
timeout_seconds=settings.provider_timeout_seconds,
|
||||||
|
)
|
||||||
|
endpoint.release_held_connection()
|
||||||
|
|
||||||
|
interrupted_state = _durable_state_evidence(
|
||||||
|
delivery_probe(prepared.campaign_id, prepared.version_id)
|
||||||
|
)
|
||||||
|
expected_interrupted = {
|
||||||
|
"job_count": 1,
|
||||||
|
"send_status_counts": {"sending": 1},
|
||||||
|
"attempt_status_counts": {"smtp_in_progress": 1},
|
||||||
|
"unfinished_attempt_count": 1,
|
||||||
|
}
|
||||||
|
if interrupted_state != expected_interrupted:
|
||||||
|
raise AcceptanceError("Killed worker state was not durably SMTP-in-progress")
|
||||||
|
|
||||||
|
replacement_worker = _start_worker(runtime_root, label="replacement-worker")
|
||||||
|
_wait_for_worker_ready(
|
||||||
|
replacement_worker,
|
||||||
|
timeout_seconds=settings.provider_timeout_seconds,
|
||||||
|
)
|
||||||
|
redelivered_task_id = _wait_for_received_task(
|
||||||
|
replacement_worker,
|
||||||
|
timeout_seconds=settings.provider_timeout_seconds,
|
||||||
|
expected_task_id=first_task_id,
|
||||||
|
)
|
||||||
|
_wait_for_task_success(
|
||||||
|
replacement_worker,
|
||||||
|
task_id=redelivered_task_id,
|
||||||
|
timeout_seconds=settings.provider_timeout_seconds,
|
||||||
|
)
|
||||||
|
recovered_state = _durable_state_evidence(
|
||||||
|
delivery_probe(prepared.campaign_id, prepared.version_id)
|
||||||
|
)
|
||||||
|
expected_recovered = {
|
||||||
|
"job_count": 1,
|
||||||
|
"send_status_counts": {"outcome_unknown": 1},
|
||||||
|
"attempt_status_counts": {"outcome_unknown": 1},
|
||||||
|
"unfinished_attempt_count": 0,
|
||||||
|
}
|
||||||
|
if recovered_state != expected_recovered:
|
||||||
|
raise AcceptanceError("Redelivered task did not freeze the unfinished attempt")
|
||||||
|
|
||||||
|
protocol = endpoint.evidence()
|
||||||
|
expected_protocol = {
|
||||||
|
"connection_count": 1,
|
||||||
|
"accepted_rcpt_commands": 1,
|
||||||
|
"refused_rcpt_commands": 0,
|
||||||
|
"data_transactions": 1,
|
||||||
|
}
|
||||||
|
if protocol != expected_protocol:
|
||||||
|
raise AcceptanceError("Broker redelivery caused an unexpected SMTP transaction")
|
||||||
|
broker_after = _wait_for_broker_drained(
|
||||||
|
redis_url,
|
||||||
|
timeout_seconds=settings.provider_timeout_seconds,
|
||||||
|
)
|
||||||
|
first_received = first_worker.received_task_ids()
|
||||||
|
replacement_received = replacement_worker.received_task_ids()
|
||||||
|
if first_received != (first_task_id,) or replacement_received != (
|
||||||
|
redelivered_task_id,
|
||||||
|
):
|
||||||
|
raise AcceptanceError(
|
||||||
|
"Celery workers did not each receive the broker task exactly once"
|
||||||
|
)
|
||||||
|
report = _report_evidence(
|
||||||
|
_expect(
|
||||||
|
client.get(
|
||||||
|
f"/api/v1/campaigns/{prepared.campaign_id}/report",
|
||||||
|
headers=dict(headers),
|
||||||
|
params={"version_id": prepared.version_id},
|
||||||
|
),
|
||||||
|
200,
|
||||||
|
"Celery-redelivery Campaign report",
|
||||||
|
)
|
||||||
|
)
|
||||||
|
if report["send_status_counts"] != {"outcome_unknown": 1}:
|
||||||
|
raise AcceptanceError("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("Celery-redelivery Campaign audit evidence is incomplete")
|
||||||
|
|
||||||
|
return {
|
||||||
|
**prepared.public_evidence(),
|
||||||
|
"queue": queue_evidence,
|
||||||
|
"interrupted_durable_state": interrupted_state,
|
||||||
|
"recovered_durable_state": recovered_state,
|
||||||
|
"protocol": protocol,
|
||||||
|
"report": report,
|
||||||
|
"audit_actions": audit_actions,
|
||||||
|
"broker": {
|
||||||
|
"transport": "redis",
|
||||||
|
"same_task_identity_redelivered": first_task_id
|
||||||
|
== redelivered_task_id,
|
||||||
|
"first_worker_received_count": len(first_received),
|
||||||
|
"replacement_worker_received_count": len(replacement_received),
|
||||||
|
**broker_after.as_dict(),
|
||||||
|
},
|
||||||
|
"supervision": {
|
||||||
|
"first_worker_killed_after_complete_data": True,
|
||||||
|
"first_worker_forced_exit": first_exit_code != 0,
|
||||||
|
"replacement_worker_started": True,
|
||||||
|
"replacement_worker_completed_redelivery": True,
|
||||||
|
},
|
||||||
|
}
|
||||||
|
finally:
|
||||||
|
endpoint.release_held_connection()
|
||||||
|
_stop_worker(first_worker, timeout_seconds=5)
|
||||||
|
if replacement_worker is not None:
|
||||||
|
_stop_worker(replacement_worker, timeout_seconds=5)
|
||||||
|
|
||||||
|
|
||||||
|
def _runtime_module_versions(registry: Any) -> dict[str, str]:
|
||||||
|
versions = {"core": _core_package_version()}
|
||||||
|
versions.update(
|
||||||
|
{
|
||||||
|
manifest.id: manifest.version
|
||||||
|
for manifest in registry.manifests()
|
||||||
|
if manifest.id in {"access", "audit", "campaigns", "mail"}
|
||||||
|
}
|
||||||
|
)
|
||||||
|
return versions
|
||||||
|
|
||||||
|
|
||||||
|
def _create_runtime_root() -> Path:
|
||||||
|
"""Create the isolated runtime under the platform-selected temp root."""
|
||||||
|
|
||||||
|
return Path(tempfile.mkdtemp(prefix="govoplan-campaign-celery-redelivery-"))
|
||||||
|
|
||||||
|
|
||||||
|
def _bootstrap_and_run(
|
||||||
|
*,
|
||||||
|
settings: TestbedSettings,
|
||||||
|
fixture_path: Path,
|
||||||
|
redis_url: str,
|
||||||
|
visibility_timeout_seconds: int,
|
||||||
|
) -> dict[str, Any]:
|
||||||
|
runtime_root = _create_runtime_root()
|
||||||
|
database = None
|
||||||
|
try:
|
||||||
|
os.environ.update(
|
||||||
|
{
|
||||||
|
"APP_ENV": "test",
|
||||||
|
"DATABASE_URL": f"sqlite:///{runtime_root / 'acceptance.db'}",
|
||||||
|
"FILE_STORAGE_BACKEND": "local",
|
||||||
|
"FILE_STORAGE_LOCAL_ROOT": str(runtime_root / "files"),
|
||||||
|
"MOCK_MAILBOX_DIR": str(runtime_root / "mock-mailbox"),
|
||||||
|
"DEV_BOOTSTRAP_ENABLED": "false",
|
||||||
|
"CELERY_ENABLED": "true",
|
||||||
|
"REDIS_URL": redis_url,
|
||||||
|
"GOVOPLAN_CAMPAIGN_TEST_REDIS_VISIBILITY_TIMEOUT_SECONDS": str(
|
||||||
|
visibility_timeout_seconds
|
||||||
|
),
|
||||||
|
"GOVOPLAN_CONNECTOR_ALLOW_PRIVATE_NETWORKS": "true",
|
||||||
|
}
|
||||||
|
)
|
||||||
|
from fastapi.testclient import TestClient
|
||||||
|
|
||||||
|
from govoplan_core.db.base import Base
|
||||||
|
from govoplan_core.db.bootstrap import bootstrap_dev_data
|
||||||
|
from govoplan_core.db.session import configure_database, set_database
|
||||||
|
from govoplan_core.settings import Settings, settings as core_settings
|
||||||
|
from govoplan_core.tenancy.scope import create_scope_tables
|
||||||
|
|
||||||
|
isolated_settings = Settings()
|
||||||
|
for field_name in Settings.model_fields:
|
||||||
|
setattr(core_settings, field_name, getattr(isolated_settings, field_name))
|
||||||
|
database = configure_database(os.environ["DATABASE_URL"])
|
||||||
|
set_database(database)
|
||||||
|
|
||||||
|
from govoplan_core.server.app import app
|
||||||
|
|
||||||
|
create_scope_tables(database.engine)
|
||||||
|
Base.metadata.create_all(bind=database.engine)
|
||||||
|
with database.SessionLocal() as session:
|
||||||
|
bootstrap_dev_data(
|
||||||
|
session,
|
||||||
|
api_key_secret="celery-redelivery-unused-api-key",
|
||||||
|
user_password="celery-redelivery-admin",
|
||||||
|
)
|
||||||
|
|
||||||
|
def snapshot_probe(
|
||||||
|
version_id: str,
|
||||||
|
) -> tuple[Mapping[str, Any], Mapping[str, Any]]:
|
||||||
|
from govoplan_campaign.backend.db.models import CampaignVersion
|
||||||
|
|
||||||
|
with database.SessionLocal() as session:
|
||||||
|
version = session.get(CampaignVersion, version_id)
|
||||||
|
if version is None:
|
||||||
|
raise AcceptanceError("Campaign execution snapshot is unavailable")
|
||||||
|
raw = version.raw_json if isinstance(version.raw_json, dict) else {}
|
||||||
|
snapshot = (
|
||||||
|
version.execution_snapshot
|
||||||
|
if isinstance(version.execution_snapshot, dict)
|
||||||
|
else {}
|
||||||
|
)
|
||||||
|
return raw, snapshot
|
||||||
|
|
||||||
|
def audit_probe(campaign_id: str, version_id: str) -> Mapping[str, int]:
|
||||||
|
from govoplan_audit.backend.db.models import AuditLog
|
||||||
|
|
||||||
|
with database.SessionLocal() as session:
|
||||||
|
actions = [
|
||||||
|
row[0]
|
||||||
|
for row in session.query(AuditLog.action)
|
||||||
|
.filter(AuditLog.object_id.in_([campaign_id, version_id]))
|
||||||
|
.all()
|
||||||
|
if row[0] in EXPECTED_AUDIT_ACTIONS
|
||||||
|
]
|
||||||
|
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
|
||||||
|
),
|
||||||
|
}
|
||||||
|
|
||||||
|
with TestClient(app) as client:
|
||||||
|
login = _expect(
|
||||||
|
client.post(
|
||||||
|
"/api/v1/auth/login",
|
||||||
|
json={
|
||||||
|
"email": "admin@example.local",
|
||||||
|
"password": "celery-redelivery-admin",
|
||||||
|
},
|
||||||
|
),
|
||||||
|
200,
|
||||||
|
"Acceptance login",
|
||||||
|
)
|
||||||
|
access_token = str(login.get("access_token") or "")
|
||||||
|
if not access_token:
|
||||||
|
raise AcceptanceError("Acceptance login returned no access token")
|
||||||
|
from govoplan_core.core.runtime import get_registry
|
||||||
|
|
||||||
|
registry = get_registry()
|
||||||
|
if registry is None:
|
||||||
|
raise AcceptanceError("The GovOPlaN module registry is unavailable")
|
||||||
|
composition_versions = required_composition_versions(
|
||||||
|
fixture_path,
|
||||||
|
_runtime_module_versions(registry),
|
||||||
|
)
|
||||||
|
with smtp_fault_endpoint("post_data_hold") as endpoint:
|
||||||
|
scenario = execute_redelivery_scenario(
|
||||||
|
client,
|
||||||
|
{"Authorization": f"Bearer {access_token}"},
|
||||||
|
fixture_path=fixture_path,
|
||||||
|
settings=settings,
|
||||||
|
endpoint=endpoint,
|
||||||
|
redis_url=redis_url,
|
||||||
|
runtime_root=runtime_root,
|
||||||
|
snapshot_probe=snapshot_probe,
|
||||||
|
audit_probe=audit_probe,
|
||||||
|
delivery_probe=delivery_probe,
|
||||||
|
)
|
||||||
|
|
||||||
|
evidence = {
|
||||||
|
"schema_version": EVIDENCE_SCHEMA,
|
||||||
|
"generated_at": datetime.now(timezone.utc).isoformat(),
|
||||||
|
"fixture_sha256": hashlib.sha256(fixture_path.read_bytes()).hexdigest(),
|
||||||
|
"declared_module_versions": composition_versions,
|
||||||
|
"target": {
|
||||||
|
"kind": "local_redis_celery_controlled_smtp",
|
||||||
|
"isolated_temporary_database": True,
|
||||||
|
"redis_started_by_runner": True,
|
||||||
|
"worker_pool": "solo",
|
||||||
|
"worker_prefetch_multiplier": 1,
|
||||||
|
"task_acks_late": True,
|
||||||
|
"task_reject_on_worker_lost": True,
|
||||||
|
"visibility_timeout_seconds": visibility_timeout_seconds,
|
||||||
|
},
|
||||||
|
"scenario": scenario,
|
||||||
|
"coverage": {
|
||||||
|
"redis_broker_delivery": True,
|
||||||
|
"celery_worker_processes": True,
|
||||||
|
"forced_worker_loss_after_complete_data": True,
|
||||||
|
"same_task_broker_redelivery": True,
|
||||||
|
"durable_outcome_unknown_recovery": True,
|
||||||
|
"duplicate_smtp_transaction_prevented": True,
|
||||||
|
"production_daemon_supervisor": False,
|
||||||
|
"target_provider": False,
|
||||||
|
"source_artifact_provenance": False,
|
||||||
|
},
|
||||||
|
}
|
||||||
|
_assert_evidence_safe(evidence, settings=settings)
|
||||||
|
rendered = json.dumps(
|
||||||
|
evidence,
|
||||||
|
ensure_ascii=False,
|
||||||
|
indent=2,
|
||||||
|
sort_keys=True,
|
||||||
|
).encode("utf-8") + b"\n"
|
||||||
|
if len(rendered) > MAX_EVIDENCE_BYTES:
|
||||||
|
raise AcceptanceError("Celery-redelivery evidence exceeds its size limit")
|
||||||
|
return evidence
|
||||||
|
finally:
|
||||||
|
if database is not None:
|
||||||
|
database.engine.dispose()
|
||||||
|
shutil.rmtree(runtime_root, ignore_errors=True)
|
||||||
|
|
||||||
|
|
||||||
|
def main(argv: list[str] | None = None) -> int:
|
||||||
|
parser = argparse.ArgumentParser(description=__doc__)
|
||||||
|
parser.add_argument("--fixture", type=Path, default=DEFAULT_FIXTURE)
|
||||||
|
parser.add_argument("--compose-file", type=Path, default=DEFAULT_COMPOSE_FILE)
|
||||||
|
parser.add_argument("--redis-port", type=int)
|
||||||
|
parser.add_argument(
|
||||||
|
"--visibility-timeout-seconds",
|
||||||
|
type=lambda value: _positive_int(value, label="visibility timeout"),
|
||||||
|
default=os.environ.get(
|
||||||
|
"GOVOPLAN_CAMPAIGN_TEST_REDIS_VISIBILITY_TIMEOUT_SECONDS",
|
||||||
|
"3",
|
||||||
|
),
|
||||||
|
)
|
||||||
|
parser.add_argument(
|
||||||
|
"--timeout-seconds",
|
||||||
|
type=lambda value: _positive_int(value, label="timeout"),
|
||||||
|
default=60,
|
||||||
|
)
|
||||||
|
parser.add_argument("--evidence", type=Path)
|
||||||
|
args = parser.parse_args(argv)
|
||||||
|
|
||||||
|
try:
|
||||||
|
settings = TestbedSettings.from_environment()
|
||||||
|
settings.assert_local_testbed()
|
||||||
|
settings = replace(
|
||||||
|
settings,
|
||||||
|
provider_timeout_seconds=args.timeout_seconds,
|
||||||
|
)
|
||||||
|
with isolated_redis_broker(
|
||||||
|
compose_file=args.compose_file.resolve(),
|
||||||
|
timeout_seconds=args.timeout_seconds,
|
||||||
|
requested_port=args.redis_port,
|
||||||
|
) as redis_url:
|
||||||
|
evidence = _bootstrap_and_run(
|
||||||
|
settings=settings,
|
||||||
|
fixture_path=args.fixture.resolve(),
|
||||||
|
redis_url=redis_url,
|
||||||
|
visibility_timeout_seconds=args.visibility_timeout_seconds,
|
||||||
|
)
|
||||||
|
rendered = json.dumps(
|
||||||
|
evidence,
|
||||||
|
ensure_ascii=False,
|
||||||
|
indent=2,
|
||||||
|
sort_keys=True,
|
||||||
|
) + "\n"
|
||||||
|
if args.evidence:
|
||||||
|
args.evidence.parent.mkdir(parents=True, exist_ok=True)
|
||||||
|
args.evidence.write_text(rendered, encoding="utf-8")
|
||||||
|
else:
|
||||||
|
sys.stdout.write(rendered)
|
||||||
|
return 0
|
||||||
|
except AcceptanceError as exc:
|
||||||
|
print(f"Campaign Celery-redelivery acceptance failed: {exc}", file=sys.stderr)
|
||||||
|
return 1
|
||||||
|
except Exception as exc:
|
||||||
|
print(
|
||||||
|
"Campaign Celery-redelivery acceptance failed unexpectedly "
|
||||||
|
f"({type(exc).__name__}); inspect local service logs.",
|
||||||
|
file=sys.stderr,
|
||||||
|
)
|
||||||
|
return 1
|
||||||
|
|
||||||
|
|
||||||
|
if __name__ == "__main__":
|
||||||
|
raise SystemExit(main())
|
||||||
@@ -104,6 +104,36 @@ bed where possible:
|
|||||||
- IMAP append failure after SMTP acceptance.
|
- IMAP append failure after SMTP acceptance.
|
||||||
- Worker restart with queued, claimed, and sending jobs.
|
- Worker restart with queued, claimed, and sending jobs.
|
||||||
|
|
||||||
|
For the maintained loopback baseline, run
|
||||||
|
`dev/mail-testbed/run_campaign_acceptance.py`. It proves the public Campaign
|
||||||
|
path for SMTP acceptance, IMAP append, repeat-send blocking, an SMTP connection
|
||||||
|
failure before transmission, an explicit SMTP authentication rejection, an
|
||||||
|
explicit temporary `451` response after DATA, partial RCPT refusal, a
|
||||||
|
connection loss after complete DATA, and an IMAP authentication rejection
|
||||||
|
after SMTP acceptance. Its evidence is an allowlisted classification/count
|
||||||
|
projection; raw provider diagnostics and transport/account identifiers are
|
||||||
|
deliberately excluded.
|
||||||
|
|
||||||
|
The runner also terminates a dedicated OS process executing the registered
|
||||||
|
Campaign send task after complete DATA, then invokes the task in a fresh
|
||||||
|
process. The unfinished durable attempt must become `outcome_unknown` and the
|
||||||
|
endpoint must observe no second connection or DATA transaction. This covers
|
||||||
|
the worker task/process boundary but not a broker or daemon.
|
||||||
|
|
||||||
|
Run `dev/mail-testbed/run_celery_redelivery_acceptance.py` for the maintained
|
||||||
|
Redis/Celery delivery and broker redelivery boundary. It starts an isolated
|
||||||
|
Redis Compose service and real Celery workers, kills the first solo worker after complete DATA while the
|
||||||
|
late-ack task is unacknowledged, and requires the same task identity to reach a
|
||||||
|
replacement worker after Redis visibility recovery. Passing evidence also
|
||||||
|
requires durable `outcome_unknown`, an empty broker queue/unacked set, and
|
||||||
|
exactly one SMTP connection and DATA transaction. Raw worker logs and task,
|
||||||
|
database, endpoint, and credential identifiers are never retained.
|
||||||
|
|
||||||
|
That second runner proves local runner-supervised process replacement, not the
|
||||||
|
production process manager. Repeat the worker-loss drill under the selected
|
||||||
|
systemd, container, Kubernetes, or other production supervisor and the target
|
||||||
|
Redis/SMTP infrastructure before deployment approval.
|
||||||
|
|
||||||
## Reporting Checks
|
## Reporting Checks
|
||||||
|
|
||||||
- Partial delivery must show accepted, failed, and unknown counts separately.
|
- Partial delivery must show accepted, failed, and unknown counts separately.
|
||||||
|
|||||||
@@ -144,10 +144,14 @@ paths, storage keys, worker claim tokens, and raw provider diagnostics require
|
|||||||
the dedicated diagnostic permission and must not leak through ordinary campaign,
|
the dedicated diagnostic permission and must not leak through ordinary campaign,
|
||||||
version, job, or report responses.
|
version, job, or report responses.
|
||||||
|
|
||||||
The current Campaign Report Web UI additionally requires recipient-read access
|
Campaign now provides a separate aggregate **Reports** surface for readers with
|
||||||
and does not yet hide every action control that the actor lacks. The server
|
`campaigns:report:read` and access to the campaign. It loads only the safe
|
||||||
still authorizes each action, but an aggregate-only reader UI remains open
|
aggregate projections, applies small-cell suppression, and offers no recipient
|
||||||
work; do not promise that experience from `campaigns:report:read` alone.
|
rows, drill-down, filtering, export, or delivery actions. The recipient-aware
|
||||||
|
**Campaign Report** still requires recipient-read access and does not yet hide
|
||||||
|
every action control that the actor lacks. The server authorizes each action,
|
||||||
|
but permission-aware action visibility on that detailed surface remains open
|
||||||
|
work; do not confuse it with the aggregate reader experience.
|
||||||
|
|
||||||
### Deliver and resolve outcomes
|
### Deliver and resolve outcomes
|
||||||
|
|
||||||
|
|||||||
@@ -10,7 +10,11 @@ scenario catalogue lives in `examples/README.md`; committed fixture files should
|
|||||||
be added under `examples/` only when they validate against the current campaign
|
be added under `examples/` only when they validate against the current campaign
|
||||||
schema and are safe to run in non-production environments.
|
schema and are safe to run in non-production environments.
|
||||||
|
|
||||||
- simple announcement with one active recipient and no attachments
|
- [`simple-announcement`](../examples/simple-announcement/campaign.json), a
|
||||||
|
credential-free campaign with one active recipient and no attachments; its
|
||||||
|
automated acceptance check physically blocks Mail and Files imports, denies
|
||||||
|
network connections, and validates/builds from an unrelated temporary
|
||||||
|
workspace
|
||||||
- multi-recipient message with To, CC, BCC, Reply-To, bounce, and disposition
|
- multi-recipient message with To, CC, BCC, Reply-To, bounce, and disposition
|
||||||
notification fields
|
notification fields
|
||||||
- campaign with global attachments and recipient-specific attachment rules
|
- campaign with global attachments and recipient-specific attachment rules
|
||||||
@@ -22,7 +26,9 @@ schema and are safe to run in non-production environments.
|
|||||||
- campaign with blocked recipients or attachment errors that must not be sent
|
- campaign with blocked recipients or attachment errors that must not be sent
|
||||||
- mock delivery campaign that captures SMTP and IMAP append messages in the mail
|
- mock delivery campaign that captures SMTP and IMAP append messages in the mail
|
||||||
development mailbox
|
development mailbox
|
||||||
- real non-production delivery campaign against the GreenMail test bed
|
- [`greenmail-delivery`](../examples/greenmail-delivery/campaign.json), a
|
||||||
|
credential-free real-delivery Campaign materialized with a temporary
|
||||||
|
Mail-owned profile by the loopback acceptance runner
|
||||||
|
|
||||||
## Fixture Rules
|
## Fixture Rules
|
||||||
|
|
||||||
@@ -41,12 +47,31 @@ Before tagging a campaign release:
|
|||||||
|
|
||||||
- Review `examples/README.md` and update the scenario catalogue when a release
|
- Review `examples/README.md` and update the scenario catalogue when a release
|
||||||
adds or removes delivery behavior.
|
adds or removes delivery behavior.
|
||||||
|
- Run `python -m unittest discover -s tests -p 'test_example_campaigns.py'` and
|
||||||
|
retain its isolated validate/build result as release evidence.
|
||||||
- Run core module permutation tests with campaign installed both with and
|
- Run core module permutation tests with campaign installed both with and
|
||||||
without files/mail.
|
without files/mail.
|
||||||
- Validate and build each maintained example campaign.
|
- Validate and build each maintained example campaign.
|
||||||
- Run the mock delivery example when the mail development mailbox capability is
|
- Run the mock delivery example when the mail development mailbox capability is
|
||||||
enabled.
|
enabled.
|
||||||
- Run the GreenMail SMTP/IMAP smoke for a non-production real delivery path.
|
- Run the GreenMail SMTP/IMAP smoke for a non-production real delivery path.
|
||||||
|
- Run `dev/mail-testbed/run_campaign_acceptance.py` and retain its bounded JSON
|
||||||
|
projection. It must show one SMTP acceptance, one IMAP append, no duplicate
|
||||||
|
effect from a repeated ordinary send, matching Campaign report/audit state,
|
||||||
|
and no resolved transport material in Campaign JSON or its execution
|
||||||
|
snapshot. Its controlled endpoint evidence must also show explicit SMTP 451,
|
||||||
|
partial RCPT refusal, post-DATA ambiguity, and task-process interruption
|
||||||
|
classifications without retaining addresses or provider diagnostics.
|
||||||
|
- Treat the task-process restart proof separately from the still-open
|
||||||
|
Redis/Celery broker redelivery and daemon-supervision check; the coverage
|
||||||
|
projection must keep `celery_broker_redelivery` false.
|
||||||
|
- Run `dev/mail-testbed/run_celery_redelivery_acceptance.py` as a separate
|
||||||
|
destructive worker-loss check. Retain its bounded evidence only when the same
|
||||||
|
broker task is observed at both workers, the durable state is
|
||||||
|
`outcome_unknown`, broker queue/unacked counts are zero, and the controlled
|
||||||
|
SMTP endpoint observed one connection and one DATA transaction. This closes
|
||||||
|
local Redis/Celery redelivery coverage, while production supervisor and
|
||||||
|
target-provider coverage remain false until separately tested.
|
||||||
- Confirm reusable mail profile selection is revalidated after campaign owner
|
- Confirm reusable mail profile selection is revalidated after campaign owner
|
||||||
transfer.
|
transfer.
|
||||||
- Confirm every inline SMTP/IMAP field is rejected on import/write, omitted
|
- Confirm every inline SMTP/IMAP field is rejected on import/write, omitted
|
||||||
|
|||||||
@@ -9,7 +9,7 @@ campaign schema and do not require production data.
|
|||||||
|
|
||||||
| Scenario | Required Modules | Release Check |
|
| Scenario | Required Modules | Release Check |
|
||||||
| --- | --- | --- |
|
| --- | --- | --- |
|
||||||
| `simple-announcement` | core, access, campaigns | Validate and build one active recipient without attachments. |
|
| [`simple-announcement`](simple-announcement/campaign.json) | core, access, campaigns | Validate and build one active recipient without attachments while Mail and Files are absent. |
|
||||||
| `addressing-matrix` | core, access, campaigns | Exercise To, CC, BCC, Reply-To, bounce, and disposition-notification fields. |
|
| `addressing-matrix` | core, access, campaigns | Exercise To, CC, BCC, Reply-To, bounce, and disposition-notification fields. |
|
||||||
| `global-attachment` | core, access, campaigns; optional files | Build one deterministic attachment and verify evidence. |
|
| `global-attachment` | core, access, campaigns; optional files | Build one deterministic attachment and verify evidence. |
|
||||||
| `recipient-attachment-rules` | core, access, campaigns; optional files | Match recipient-specific attachment rules and verify per-recipient evidence. |
|
| `recipient-attachment-rules` | core, access, campaigns; optional files | Match recipient-specific attachment rules and verify per-recipient evidence. |
|
||||||
@@ -19,7 +19,7 @@ campaign schema and do not require production data.
|
|||||||
| `warnings-review` | core, access, campaigns | Require explicit review before queueing jobs with warnings. |
|
| `warnings-review` | core, access, campaigns | Require explicit review before queueing jobs with warnings. |
|
||||||
| `blocked-send` | core, access, campaigns | Confirm blocked recipients or missing attachments cannot be queued. |
|
| `blocked-send` | core, access, campaigns | Confirm blocked recipients or missing attachments cannot be queued. |
|
||||||
| `mock-delivery` | core, access, campaigns, mail with dev capability | Capture messages in the development mailbox. |
|
| `mock-delivery` | core, access, campaigns, mail with dev capability | Capture messages in the development mailbox. |
|
||||||
| `greenmail-delivery` | core, access, campaigns, mail | Send no-attachment, normal attachment, and ZIP attachment variants through `dev/mail-testbed`. |
|
| [`greenmail-delivery`](greenmail-delivery/campaign.json) | core, access, audit, campaigns, mail | Run a credential-free Campaign through a Mail-owned profile, GreenMail SMTP/IMAP, report/audit checks, repeat-send protection, and bounded failure drills. |
|
||||||
|
|
||||||
## Fixture Rules
|
## Fixture Rules
|
||||||
|
|
||||||
@@ -37,9 +37,14 @@ campaign schema and do not require production data.
|
|||||||
Before a release tag:
|
Before a release tag:
|
||||||
|
|
||||||
1. Run module permutation startup checks from core.
|
1. Run module permutation startup checks from core.
|
||||||
2. Validate every committed example fixture against the current campaign schema.
|
2. Run `python -m unittest discover -s tests -p 'test_example_campaigns.py'`
|
||||||
3. Build exact messages for each fixture.
|
from this repository. The acceptance test copies each maintained fixture to
|
||||||
4. Run the mock-delivery example when the dev mailbox capability is enabled.
|
an unrelated temporary workspace before using Campaign's public loader,
|
||||||
5. Run `dev/mail-testbed/run_transport_smoke.py`.
|
validator, and message builder.
|
||||||
6. Execute the delivery checklist in
|
3. Validate every committed example fixture against the current campaign schema.
|
||||||
|
4. Build exact messages for each fixture.
|
||||||
|
5. Run the mock-delivery example when the dev mailbox capability is enabled.
|
||||||
|
6. Run `dev/mail-testbed/run_transport_smoke.py` for low-level transport and attachment variants.
|
||||||
|
7. Run `dev/mail-testbed/run_campaign_acceptance.py` for the Campaign journey and bounded evidence.
|
||||||
|
8. Execute the delivery checklist in
|
||||||
`docs/EXAMPLE_CAMPAIGNS_AND_RELEASE_CHECKLIST.md`.
|
`docs/EXAMPLE_CAMPAIGNS_AND_RELEASE_CHECKLIST.md`.
|
||||||
|
|||||||
88
examples/greenmail-delivery/campaign.json
Normal file
88
examples/greenmail-delivery/campaign.json
Normal file
@@ -0,0 +1,88 @@
|
|||||||
|
{
|
||||||
|
"version": "1.0",
|
||||||
|
"campaign": {
|
||||||
|
"id": "greenmail-delivery",
|
||||||
|
"name": "GreenMail delivery acceptance",
|
||||||
|
"description": "Credential-free Campaign fixture for the local SMTP/IMAP acceptance test bed.",
|
||||||
|
"mode": "test"
|
||||||
|
},
|
||||||
|
"fields": [
|
||||||
|
{
|
||||||
|
"name": "display_name",
|
||||||
|
"type": "string",
|
||||||
|
"label": "Display name",
|
||||||
|
"required": true
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"name": "acceptance_run",
|
||||||
|
"type": "string",
|
||||||
|
"label": "Acceptance run",
|
||||||
|
"required": true
|
||||||
|
}
|
||||||
|
],
|
||||||
|
"server": {
|
||||||
|
"mail_profile_id": "00000000-0000-4000-8000-000000000001"
|
||||||
|
},
|
||||||
|
"recipients": {
|
||||||
|
"from": [
|
||||||
|
{
|
||||||
|
"email": "campaign-test@govoplan.test",
|
||||||
|
"name": "GovOPlaN acceptance",
|
||||||
|
"type": "to"
|
||||||
|
}
|
||||||
|
],
|
||||||
|
"allow_individual_to": true
|
||||||
|
},
|
||||||
|
"template": {
|
||||||
|
"subject": "[GovOPlaN acceptance ${acceptance_run}] Campaign delivery",
|
||||||
|
"text": "Hello ${display_name},\n\nThis is an isolated GovOPlaN Campaign SMTP/IMAP acceptance message.\n",
|
||||||
|
"body_mode": "text"
|
||||||
|
},
|
||||||
|
"attachments": {
|
||||||
|
"base_path": ".",
|
||||||
|
"send_without_attachments_behavior": "continue",
|
||||||
|
"global": []
|
||||||
|
},
|
||||||
|
"entries": {
|
||||||
|
"inline": [
|
||||||
|
{
|
||||||
|
"id": "greenmail-recipient",
|
||||||
|
"to": [
|
||||||
|
{
|
||||||
|
"email": "campaign-test@govoplan.test",
|
||||||
|
"name": "GreenMail recipient",
|
||||||
|
"type": "to"
|
||||||
|
}
|
||||||
|
],
|
||||||
|
"fields": {
|
||||||
|
"display_name": "GreenMail recipient",
|
||||||
|
"acceptance_run": "fixture"
|
||||||
|
}
|
||||||
|
}
|
||||||
|
]
|
||||||
|
},
|
||||||
|
"validation_policy": {
|
||||||
|
"missing_email": "block",
|
||||||
|
"template_error": "block"
|
||||||
|
},
|
||||||
|
"delivery": {
|
||||||
|
"rate_limit": {
|
||||||
|
"messages_per_minute": 60
|
||||||
|
},
|
||||||
|
"retry": {
|
||||||
|
"max_attempts": 3,
|
||||||
|
"backoff_seconds": [
|
||||||
|
1,
|
||||||
|
5,
|
||||||
|
30
|
||||||
|
]
|
||||||
|
},
|
||||||
|
"imap_append_sent": {
|
||||||
|
"enabled": true,
|
||||||
|
"folder": "Sent"
|
||||||
|
}
|
||||||
|
},
|
||||||
|
"status_tracking": {
|
||||||
|
"enabled": true
|
||||||
|
}
|
||||||
|
}
|
||||||
22
examples/greenmail-delivery/fixture.json
Normal file
22
examples/greenmail-delivery/fixture.json
Normal file
@@ -0,0 +1,22 @@
|
|||||||
|
{
|
||||||
|
"scenario": "greenmail-delivery",
|
||||||
|
"campaign_file": "campaign.json",
|
||||||
|
"required_modules": [
|
||||||
|
"core",
|
||||||
|
"access",
|
||||||
|
"audit",
|
||||||
|
"campaigns",
|
||||||
|
"mail"
|
||||||
|
],
|
||||||
|
"required_capabilities": [
|
||||||
|
"mail.campaign_delivery"
|
||||||
|
],
|
||||||
|
"transport": "local GreenMail SMTP/IMAP test bed",
|
||||||
|
"credentials": "local environment only; never copied into Campaign JSON or evidence",
|
||||||
|
"expected": {
|
||||||
|
"entries_count": 1,
|
||||||
|
"built_count": 1,
|
||||||
|
"smtp_accepted_count": 1,
|
||||||
|
"imap_appended_count": 1
|
||||||
|
}
|
||||||
|
}
|
||||||
54
examples/simple-announcement/campaign.json
Normal file
54
examples/simple-announcement/campaign.json
Normal file
@@ -0,0 +1,54 @@
|
|||||||
|
{
|
||||||
|
"version": "1.0",
|
||||||
|
"campaign": {
|
||||||
|
"id": "simple-announcement",
|
||||||
|
"name": "Simple announcement",
|
||||||
|
"description": "Credential-free release fixture for Campaign validation and message building.",
|
||||||
|
"mode": "test"
|
||||||
|
},
|
||||||
|
"fields": [
|
||||||
|
{
|
||||||
|
"name": "display_name",
|
||||||
|
"type": "string",
|
||||||
|
"label": "Display name",
|
||||||
|
"required": true
|
||||||
|
}
|
||||||
|
],
|
||||||
|
"recipients": {
|
||||||
|
"from": [
|
||||||
|
{
|
||||||
|
"email": "announcements@example.test",
|
||||||
|
"name": "GovOPlaN Example",
|
||||||
|
"type": "to"
|
||||||
|
}
|
||||||
|
],
|
||||||
|
"allow_individual_to": true
|
||||||
|
},
|
||||||
|
"template": {
|
||||||
|
"subject": "Planned service maintenance for ${display_name}",
|
||||||
|
"text": "Hello ${display_name},\n\nThe example service will be unavailable during the announced maintenance window.\n\nThis message was built locally and was not sent.\n",
|
||||||
|
"body_mode": "text"
|
||||||
|
},
|
||||||
|
"attachments": {
|
||||||
|
"base_path": ".",
|
||||||
|
"send_without_attachments_behavior": "continue",
|
||||||
|
"global": []
|
||||||
|
},
|
||||||
|
"entries": {
|
||||||
|
"inline": [
|
||||||
|
{
|
||||||
|
"id": "example-recipient",
|
||||||
|
"to": [
|
||||||
|
{
|
||||||
|
"email": "recipient@example.test",
|
||||||
|
"name": "Example Recipient",
|
||||||
|
"type": "to"
|
||||||
|
}
|
||||||
|
],
|
||||||
|
"fields": {
|
||||||
|
"display_name": "Example Recipient"
|
||||||
|
}
|
||||||
|
}
|
||||||
|
]
|
||||||
|
}
|
||||||
|
}
|
||||||
23
examples/simple-announcement/fixture.json
Normal file
23
examples/simple-announcement/fixture.json
Normal file
@@ -0,0 +1,23 @@
|
|||||||
|
{
|
||||||
|
"schema_version": 1,
|
||||||
|
"id": "simple-announcement",
|
||||||
|
"campaign_file": "campaign.json",
|
||||||
|
"required_modules": [
|
||||||
|
"core",
|
||||||
|
"access",
|
||||||
|
"campaigns"
|
||||||
|
],
|
||||||
|
"absent_optional_modules": [
|
||||||
|
"files",
|
||||||
|
"mail"
|
||||||
|
],
|
||||||
|
"external_effects": "forbidden",
|
||||||
|
"expected": {
|
||||||
|
"campaign_id": "simple-announcement",
|
||||||
|
"entries_count": 1,
|
||||||
|
"built_count": 1,
|
||||||
|
"queueable_count": 1,
|
||||||
|
"attachment_count": 0,
|
||||||
|
"subject": "Planned service maintenance for Example Recipient"
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -1,6 +1,6 @@
|
|||||||
{
|
{
|
||||||
"name": "@govoplan/campaign-webui",
|
"name": "@govoplan/campaign-webui",
|
||||||
"version": "0.1.10",
|
"version": "0.1.11",
|
||||||
"private": true,
|
"private": true,
|
||||||
"type": "module",
|
"type": "module",
|
||||||
"main": "webui/src/index.ts",
|
"main": "webui/src/index.ts",
|
||||||
|
|||||||
@@ -4,7 +4,7 @@ build-backend = "setuptools.build_meta"
|
|||||||
|
|
||||||
[project]
|
[project]
|
||||||
name = "govoplan-campaign"
|
name = "govoplan-campaign"
|
||||||
version = "0.1.10"
|
version = "0.1.11"
|
||||||
description = "GovOPlaN campaigns module with backend and WebUI integration."
|
description = "GovOPlaN campaigns module with backend and WebUI integration."
|
||||||
readme = "README.md"
|
readme = "README.md"
|
||||||
requires-python = ">=3.12"
|
requires-python = ">=3.12"
|
||||||
|
|||||||
@@ -279,6 +279,12 @@ class CampaignDeliveryTaskService(CampaignDeliveryTaskProvider):
|
|||||||
|
|
||||||
|
|
||||||
def delivery_tasks_capability(context: object) -> CampaignDeliveryTaskService:
|
def delivery_tasks_capability(context: object) -> CampaignDeliveryTaskService:
|
||||||
|
from govoplan_campaign.backend.runtime import configure_runtime
|
||||||
|
|
||||||
|
configure_runtime(
|
||||||
|
registry=getattr(context, "registry", None),
|
||||||
|
settings=getattr(context, "settings", None),
|
||||||
|
)
|
||||||
return CampaignDeliveryTaskService()
|
return CampaignDeliveryTaskService()
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
@@ -156,7 +156,7 @@ def _campaigns_router(context: ModuleContext):
|
|||||||
manifest = ModuleManifest(
|
manifest = ModuleManifest(
|
||||||
id="campaigns",
|
id="campaigns",
|
||||||
name="Campaigns",
|
name="Campaigns",
|
||||||
version="0.1.10",
|
version="0.1.11",
|
||||||
required_capabilities=(CAPABILITY_AUTH_PRINCIPAL_RESOLVER, CAPABILITY_AUTH_PERMISSION_EVALUATOR),
|
required_capabilities=(CAPABILITY_AUTH_PRINCIPAL_RESOLVER, CAPABILITY_AUTH_PERMISSION_EVALUATOR),
|
||||||
optional_dependencies=("files", "mail", "notifications", "addresses"),
|
optional_dependencies=("files", "mail", "notifications", "addresses"),
|
||||||
provides_interfaces=(
|
provides_interfaces=(
|
||||||
|
|||||||
@@ -3608,6 +3608,10 @@ def send_campaign_now_endpoint(
|
|||||||
)
|
)
|
||||||
return SendCampaignNowResponse(result=response_result)
|
return SendCampaignNowResponse(result=response_result)
|
||||||
except SynchronousSendRejected as exc:
|
except SynchronousSendRejected as exc:
|
||||||
|
# A synchronous request stages queue state before the all-message
|
||||||
|
# preflight can run. Rejecting that preflight must not leave work
|
||||||
|
# eligible for a background worker when no provider effect occurred.
|
||||||
|
session.rollback()
|
||||||
audit_from_principal(
|
audit_from_principal(
|
||||||
session,
|
session,
|
||||||
principal,
|
principal,
|
||||||
|
|||||||
@@ -622,6 +622,7 @@ def _persist_campaign_queue(
|
|||||||
version: CampaignVersion,
|
version: CampaignVersion,
|
||||||
queued: list[CampaignJob],
|
queued: list[CampaignJob],
|
||||||
delivery_mode: str,
|
delivery_mode: str,
|
||||||
|
commit: bool = True,
|
||||||
) -> None:
|
) -> None:
|
||||||
if queued:
|
if queued:
|
||||||
previous_status = campaign.status
|
previous_status = campaign.status
|
||||||
@@ -640,6 +641,7 @@ def _persist_campaign_queue(
|
|||||||
version_id=version.id,
|
version_id=version.id,
|
||||||
)
|
)
|
||||||
session.add(campaign)
|
session.add(campaign)
|
||||||
|
if commit:
|
||||||
session.commit()
|
session.commit()
|
||||||
|
|
||||||
|
|
||||||
@@ -661,6 +663,7 @@ def queue_campaign_jobs(
|
|||||||
include_warnings: bool = True,
|
include_warnings: bool = True,
|
||||||
dry_run: bool = False,
|
dry_run: bool = False,
|
||||||
delivery_mode: str | None = None,
|
delivery_mode: str | None = None,
|
||||||
|
commit_queue: bool = True,
|
||||||
) -> QueueCampaignResult:
|
) -> QueueCampaignResult:
|
||||||
"""Move queueable DB jobs to QUEUED and optionally enqueue Celery tasks."""
|
"""Move queueable DB jobs to QUEUED and optionally enqueue Celery tasks."""
|
||||||
|
|
||||||
@@ -694,6 +697,7 @@ def queue_campaign_jobs(
|
|||||||
version=version,
|
version=version,
|
||||||
queued=queued,
|
queued=queued,
|
||||||
delivery_mode=selected_delivery_mode,
|
delivery_mode=selected_delivery_mode,
|
||||||
|
commit=commit_queue,
|
||||||
)
|
)
|
||||||
enqueued_count = _enqueue_campaign_jobs(
|
enqueued_count = _enqueue_campaign_jobs(
|
||||||
queued,
|
queued,
|
||||||
@@ -766,6 +770,7 @@ def send_campaign_now(
|
|||||||
enqueue_celery=False,
|
enqueue_celery=False,
|
||||||
dry_run=dry_run,
|
dry_run=dry_run,
|
||||||
delivery_mode=DELIVERY_MODE_SYNCHRONOUS,
|
delivery_mode=DELIVERY_MODE_SYNCHRONOUS,
|
||||||
|
commit_queue=False,
|
||||||
)
|
)
|
||||||
if dry_run:
|
if dry_run:
|
||||||
return SendCampaignNowResult(
|
return SendCampaignNowResult(
|
||||||
@@ -802,6 +807,11 @@ def send_campaign_now(
|
|||||||
jobs=jobs,
|
jobs=jobs,
|
||||||
policy=synchronous_policy,
|
policy=synchronous_policy,
|
||||||
)
|
)
|
||||||
|
# Queue state and its inbox notification become durable only after every
|
||||||
|
# message and the selected transport revision have passed preflight. This
|
||||||
|
# preserves late-ack recovery without leaving rejected work eligible for a
|
||||||
|
# background worker.
|
||||||
|
session.commit()
|
||||||
|
|
||||||
results: list[dict[str, Any]] = []
|
results: list[dict[str, Any]] = []
|
||||||
sent_count = 0
|
sent_count = 0
|
||||||
|
|||||||
314
tests/test_celery_redelivery_acceptance.py
Normal file
314
tests/test_celery_redelivery_acceptance.py
Normal file
@@ -0,0 +1,314 @@
|
|||||||
|
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
|
||||||
@@ -1,6 +1,7 @@
|
|||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
||||||
from dataclasses import dataclass
|
from dataclasses import dataclass
|
||||||
|
from pathlib import Path
|
||||||
from types import SimpleNamespace
|
from types import SimpleNamespace
|
||||||
from unittest.mock import patch
|
from unittest.mock import patch
|
||||||
|
|
||||||
@@ -364,6 +365,18 @@ def test_aggregate_report_task_never_implies_recipient_detail_or_export_authorit
|
|||||||
assert "export" in topic.metadata["verification"].lower()
|
assert "export" in topic.metadata["verification"].lower()
|
||||||
|
|
||||||
|
|
||||||
|
def test_handbook_distinguishes_shipped_aggregate_reports_from_detailed_report_gaps() -> None:
|
||||||
|
handbook = " ".join(
|
||||||
|
(
|
||||||
|
Path(__file__).resolve().parents[1] / "docs" / "CAMPAIGN_HANDBOOK.md"
|
||||||
|
).read_text(encoding="utf-8").lower().split()
|
||||||
|
)
|
||||||
|
|
||||||
|
assert "aggregate-only reader ui remains open" not in handbook
|
||||||
|
assert "separate aggregate **reports** surface" in handbook
|
||||||
|
assert "permission-aware action visibility on that detailed surface remains open work" in handbook
|
||||||
|
|
||||||
|
|
||||||
def test_static_campaign_handbook_has_unique_ids_help_contexts_and_no_planned_resend_claim() -> None:
|
def test_static_campaign_handbook_has_unique_ids_help_contexts_and_no_planned_resend_claim() -> None:
|
||||||
from govoplan_campaign.backend.manifest import get_manifest
|
from govoplan_campaign.backend.manifest import get_manifest
|
||||||
|
|
||||||
|
|||||||
207
tests/test_example_campaigns.py
Normal file
207
tests/test_example_campaigns.py
Normal file
@@ -0,0 +1,207 @@
|
|||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import json
|
||||||
|
import shutil
|
||||||
|
import subprocess
|
||||||
|
import sys
|
||||||
|
import tempfile
|
||||||
|
import unittest
|
||||||
|
from pathlib import Path
|
||||||
|
|
||||||
|
|
||||||
|
REPOSITORY_ROOT = Path(__file__).resolve().parents[1]
|
||||||
|
FIXTURE_ROOT = REPOSITORY_ROOT / "examples" / "simple-announcement"
|
||||||
|
|
||||||
|
|
||||||
|
_ISOLATED_ACCEPTANCE_PROGRAM = r"""
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import hashlib
|
||||||
|
import importlib.abc
|
||||||
|
import json
|
||||||
|
import socket
|
||||||
|
import sys
|
||||||
|
from email import policy
|
||||||
|
from email.parser import BytesParser
|
||||||
|
from pathlib import Path
|
||||||
|
|
||||||
|
|
||||||
|
source_root = Path(sys.argv[1]).resolve()
|
||||||
|
fixture_root = Path(sys.argv[2]).resolve()
|
||||||
|
sys.path.insert(0, str(source_root))
|
||||||
|
|
||||||
|
|
||||||
|
class AbsentOptionalModuleFinder(importlib.abc.MetaPathFinder):
|
||||||
|
absent_roots = {"govoplan_files", "govoplan_mail"}
|
||||||
|
|
||||||
|
def find_spec(self, fullname, path=None, target=None):
|
||||||
|
if fullname.partition(".")[0] in self.absent_roots:
|
||||||
|
raise ModuleNotFoundError(
|
||||||
|
f"{fullname} is intentionally absent in the Campaign fixture check",
|
||||||
|
name=fullname,
|
||||||
|
)
|
||||||
|
return None
|
||||||
|
|
||||||
|
|
||||||
|
sys.meta_path.insert(0, AbsentOptionalModuleFinder())
|
||||||
|
|
||||||
|
|
||||||
|
def deny_network(*args, **kwargs):
|
||||||
|
raise AssertionError("the Campaign validate/build fixture must not open a network connection")
|
||||||
|
|
||||||
|
|
||||||
|
class NoNetworkSocket(socket.socket):
|
||||||
|
def connect(self, address):
|
||||||
|
deny_network(address)
|
||||||
|
|
||||||
|
def connect_ex(self, address):
|
||||||
|
deny_network(address)
|
||||||
|
|
||||||
|
|
||||||
|
socket.create_connection = deny_network
|
||||||
|
socket.socket = NoNetworkSocket
|
||||||
|
|
||||||
|
from govoplan_campaign.backend.campaign import ( # noqa: E402
|
||||||
|
load_campaign_config,
|
||||||
|
load_campaign_json,
|
||||||
|
validate_campaign_config,
|
||||||
|
)
|
||||||
|
from govoplan_campaign.backend.messages import build_campaign_messages # noqa: E402
|
||||||
|
|
||||||
|
|
||||||
|
metadata = json.loads((fixture_root / "fixture.json").read_text(encoding="utf-8"))
|
||||||
|
assert metadata["required_modules"] == ["core", "access", "campaigns"]
|
||||||
|
assert metadata["absent_optional_modules"] == ["files", "mail"]
|
||||||
|
assert metadata["external_effects"] == "forbidden"
|
||||||
|
|
||||||
|
campaign_file = fixture_root / metadata["campaign_file"]
|
||||||
|
raw_campaign = load_campaign_json(campaign_file)
|
||||||
|
|
||||||
|
|
||||||
|
def iter_keys(value):
|
||||||
|
if isinstance(value, dict):
|
||||||
|
for key, nested in value.items():
|
||||||
|
yield key.casefold()
|
||||||
|
yield from iter_keys(nested)
|
||||||
|
elif isinstance(value, list):
|
||||||
|
for nested in value:
|
||||||
|
yield from iter_keys(nested)
|
||||||
|
|
||||||
|
|
||||||
|
forbidden_key_fragments = ("credential", "imap", "mail_profile", "password", "secret", "smtp")
|
||||||
|
assert not [
|
||||||
|
key
|
||||||
|
for key in iter_keys(raw_campaign)
|
||||||
|
if any(fragment in key for fragment in forbidden_key_fragments)
|
||||||
|
]
|
||||||
|
|
||||||
|
config = load_campaign_config(campaign_file)
|
||||||
|
assert config.server.mail_profile_id is None
|
||||||
|
assert not config.attachments.global_
|
||||||
|
|
||||||
|
validation = validate_campaign_config(config, campaign_file=campaign_file, check_files=True)
|
||||||
|
assert validation.ok
|
||||||
|
assert validation.error_count == 0
|
||||||
|
assert validation.warning_count == 0
|
||||||
|
assert validation.entries_count == metadata["expected"]["entries_count"]
|
||||||
|
|
||||||
|
|
||||||
|
def build(output_name):
|
||||||
|
return build_campaign_messages(
|
||||||
|
config,
|
||||||
|
campaign_file=campaign_file,
|
||||||
|
output_dir=fixture_root.parent / output_name,
|
||||||
|
write_eml=True,
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
first = build("build-first")
|
||||||
|
second = build("build-second")
|
||||||
|
expected = metadata["expected"]
|
||||||
|
for result in (first, second):
|
||||||
|
assert result.report.campaign_id == expected["campaign_id"]
|
||||||
|
assert result.report.entries_count == expected["entries_count"]
|
||||||
|
assert result.report.built_count == expected["built_count"]
|
||||||
|
assert result.report.build_failed_count == 0
|
||||||
|
assert result.report.queueable_count == expected["queueable_count"]
|
||||||
|
assert len(result.built_messages) == 1
|
||||||
|
|
||||||
|
built = result.built_messages[0]
|
||||||
|
assert built.mime is not None
|
||||||
|
assert built.draft.subject == expected["subject"]
|
||||||
|
assert built.draft.validation_status.value == "ready"
|
||||||
|
assert built.draft.send_status.value == "draft"
|
||||||
|
assert built.draft.imap_status.value == "not_requested"
|
||||||
|
assert built.draft.attachment_count == expected["attachment_count"]
|
||||||
|
assert not built.draft.attachments
|
||||||
|
assert not built.draft.issues
|
||||||
|
assert built.draft.from_ is not None
|
||||||
|
assert built.draft.from_.email == "announcements@example.test"
|
||||||
|
assert [address.email for address in built.draft.to] == ["recipient@example.test"]
|
||||||
|
assert built.mime["Subject"] == expected["subject"]
|
||||||
|
assert "Hello Example Recipient" in built.mime.get_content()
|
||||||
|
assert list(built.mime.iter_attachments()) == []
|
||||||
|
|
||||||
|
|
||||||
|
def normalized_eml(path_value):
|
||||||
|
message = BytesParser(policy=policy.default).parsebytes(Path(path_value).read_bytes())
|
||||||
|
del message["Date"]
|
||||||
|
del message["Message-ID"]
|
||||||
|
return message.as_bytes(policy=policy.default)
|
||||||
|
|
||||||
|
|
||||||
|
first_eml = normalized_eml(first.report.messages[0].eml_path)
|
||||||
|
second_eml = normalized_eml(second.report.messages[0].eml_path)
|
||||||
|
assert first_eml == second_eml
|
||||||
|
assert not any(
|
||||||
|
name == root or name.startswith(root + ".")
|
||||||
|
for name in sys.modules
|
||||||
|
for root in ("govoplan_files", "govoplan_mail")
|
||||||
|
)
|
||||||
|
|
||||||
|
print(json.dumps({
|
||||||
|
"campaign_id": first.report.campaign_id,
|
||||||
|
"built_count": first.report.built_count,
|
||||||
|
"queueable_count": first.report.queueable_count,
|
||||||
|
"normalized_eml_sha256": hashlib.sha256(first_eml).hexdigest(),
|
||||||
|
}, sort_keys=True))
|
||||||
|
"""
|
||||||
|
|
||||||
|
|
||||||
|
class CampaignExampleAcceptanceTests(unittest.TestCase):
|
||||||
|
def test_simple_announcement_validates_and_builds_without_mail_or_files(self) -> None:
|
||||||
|
self.assertTrue((FIXTURE_ROOT / "campaign.json").is_file())
|
||||||
|
self.assertTrue((FIXTURE_ROOT / "fixture.json").is_file())
|
||||||
|
|
||||||
|
with tempfile.TemporaryDirectory(prefix="govoplan-campaign-acceptance-", dir="/tmp") as temp_dir:
|
||||||
|
workspace = Path(temp_dir).resolve()
|
||||||
|
self.assertNotIn(REPOSITORY_ROOT, workspace.parents)
|
||||||
|
isolated_fixture = workspace / "simple-announcement"
|
||||||
|
shutil.copytree(FIXTURE_ROOT, isolated_fixture)
|
||||||
|
|
||||||
|
completed = subprocess.run(
|
||||||
|
[
|
||||||
|
sys.executable,
|
||||||
|
"-I",
|
||||||
|
"-c",
|
||||||
|
_ISOLATED_ACCEPTANCE_PROGRAM,
|
||||||
|
str(REPOSITORY_ROOT / "src"),
|
||||||
|
str(isolated_fixture),
|
||||||
|
],
|
||||||
|
cwd=workspace,
|
||||||
|
check=False,
|
||||||
|
capture_output=True,
|
||||||
|
text=True,
|
||||||
|
timeout=30,
|
||||||
|
)
|
||||||
|
|
||||||
|
self.assertEqual(completed.returncode, 0, completed.stderr or completed.stdout)
|
||||||
|
evidence = json.loads(completed.stdout)
|
||||||
|
self.assertEqual(evidence["campaign_id"], "simple-announcement")
|
||||||
|
self.assertEqual(evidence["built_count"], 1)
|
||||||
|
self.assertEqual(evidence["queueable_count"], 1)
|
||||||
|
self.assertRegex(evidence["normalized_eml_sha256"], r"^[0-9a-f]{64}$")
|
||||||
|
|
||||||
|
|
||||||
|
if __name__ == "__main__":
|
||||||
|
unittest.main()
|
||||||
480
tests/test_mail_testbed_acceptance.py
Normal file
480
tests/test_mail_testbed_acceptance.py
Normal file
@@ -0,0 +1,480 @@
|
|||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import importlib.util
|
||||||
|
import json
|
||||||
|
import smtplib
|
||||||
|
import sys
|
||||||
|
from pathlib import Path
|
||||||
|
from typing import Any
|
||||||
|
|
||||||
|
import pytest
|
||||||
|
|
||||||
|
from govoplan_campaign.backend.campaign import load_campaign_config
|
||||||
|
|
||||||
|
|
||||||
|
REPOSITORY_ROOT = Path(__file__).resolve().parents[1]
|
||||||
|
RUNNER_PATH = REPOSITORY_ROOT / "dev" / "mail-testbed" / "run_campaign_acceptance.py"
|
||||||
|
FIXTURE_PATH = REPOSITORY_ROOT / "examples" / "greenmail-delivery" / "campaign.json"
|
||||||
|
|
||||||
|
|
||||||
|
def _load_runner():
|
||||||
|
spec = importlib.util.spec_from_file_location("govoplan_campaign_greenmail_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[str, Any]) -> None:
|
||||||
|
self.status_code = status_code
|
||||||
|
self._payload = payload
|
||||||
|
|
||||||
|
def json(self) -> dict[str, Any]:
|
||||||
|
return self._payload
|
||||||
|
|
||||||
|
|
||||||
|
class _AcceptanceClient:
|
||||||
|
def __init__(self) -> None:
|
||||||
|
self.campaign_json: dict[str, Any] | None = None
|
||||||
|
self.send_calls = 0
|
||||||
|
|
||||||
|
def post(self, path: str, **kwargs: Any) -> _Response:
|
||||||
|
if path == "/api/v1/campaigns":
|
||||||
|
self.campaign_json = kwargs["json"]["config"]
|
||||||
|
return _Response(
|
||||||
|
200,
|
||||||
|
{
|
||||||
|
"campaign": {"id": "campaign-internal"},
|
||||||
|
"version": {"id": "version-internal"},
|
||||||
|
},
|
||||||
|
)
|
||||||
|
if path.endswith("/validate"):
|
||||||
|
return _Response(200, {"ok": True, "error_count": 0, "warning_count": 0})
|
||||||
|
if path.endswith("/build"):
|
||||||
|
return _Response(
|
||||||
|
200,
|
||||||
|
{
|
||||||
|
"built_count": 1,
|
||||||
|
"build_failed_count": 0,
|
||||||
|
"queueable_count": 1,
|
||||||
|
},
|
||||||
|
)
|
||||||
|
if path.endswith("/send-now"):
|
||||||
|
self.send_calls += 1
|
||||||
|
if self.send_calls == 2:
|
||||||
|
return _Response(422, {"detail": "Already accepted"})
|
||||||
|
return _Response(
|
||||||
|
200,
|
||||||
|
{
|
||||||
|
"result": {
|
||||||
|
"attempted_count": 1,
|
||||||
|
"sent_count": 1,
|
||||||
|
"failed_count": 0,
|
||||||
|
"outcome_unknown_count": 0,
|
||||||
|
"skipped_count": 0,
|
||||||
|
"delivery_mode": "synchronous",
|
||||||
|
"results": [{"job_id": "not-retained", "status": "smtp_accepted"}],
|
||||||
|
}
|
||||||
|
},
|
||||||
|
)
|
||||||
|
if path.endswith("/append-sent"):
|
||||||
|
return _Response(
|
||||||
|
200,
|
||||||
|
{
|
||||||
|
"result": {
|
||||||
|
"pending_count": 1,
|
||||||
|
"processed_count": 1,
|
||||||
|
"appended_count": 1,
|
||||||
|
"failed_count": 0,
|
||||||
|
"skipped_count": 0,
|
||||||
|
"results": [{"job_id": "not-retained", "status": "appended"}],
|
||||||
|
}
|
||||||
|
},
|
||||||
|
)
|
||||||
|
raise AssertionError(f"unexpected POST {path}")
|
||||||
|
|
||||||
|
def get(self, path: str, **kwargs: Any) -> _Response:
|
||||||
|
if path.endswith("/report"):
|
||||||
|
return _Response(
|
||||||
|
200,
|
||||||
|
{
|
||||||
|
"cards": {
|
||||||
|
"jobs_total": 1,
|
||||||
|
"sent": 1,
|
||||||
|
"smtp_accepted": 1,
|
||||||
|
"failed": 0,
|
||||||
|
"outcome_unknown": 0,
|
||||||
|
"retryable": 0,
|
||||||
|
"needs_attention": 0,
|
||||||
|
"imap_appended": 1,
|
||||||
|
"imap_failed": 0,
|
||||||
|
},
|
||||||
|
"status_counts": {
|
||||||
|
"send": {"smtp_accepted": 1},
|
||||||
|
"imap": {"appended": 1},
|
||||||
|
},
|
||||||
|
},
|
||||||
|
)
|
||||||
|
raise AssertionError(f"unexpected GET {path}")
|
||||||
|
|
||||||
|
|
||||||
|
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,
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.parametrize("host", ["localhost", "mail.test", "192.168.1.20", "8.8.8.8"])
|
||||||
|
def test_testbed_rejects_hostnames_and_non_loopback_addresses(host: str) -> None:
|
||||||
|
settings = _settings()
|
||||||
|
rejected = runner.TestbedSettings(
|
||||||
|
smtp_host=host,
|
||||||
|
smtp_port=settings.smtp_port,
|
||||||
|
imap_host=settings.imap_host,
|
||||||
|
imap_port=settings.imap_port,
|
||||||
|
username=settings.username,
|
||||||
|
password=settings.password,
|
||||||
|
sender=settings.sender,
|
||||||
|
recipient=settings.recipient,
|
||||||
|
sent_folder=settings.sent_folder,
|
||||||
|
provider_timeout_seconds=settings.provider_timeout_seconds,
|
||||||
|
)
|
||||||
|
|
||||||
|
with pytest.raises(runner.AcceptanceError, match="literal loopback|restricted to the loopback"):
|
||||||
|
rejected.assert_local_testbed()
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.parametrize("host", ["127.0.0.1", "127.8.9.10", "::1"])
|
||||||
|
def test_testbed_accepts_literal_loopback_and_preserves_it_in_profile(host: str) -> None:
|
||||||
|
settings = _settings()
|
||||||
|
accepted = runner.TestbedSettings(
|
||||||
|
smtp_host=host,
|
||||||
|
smtp_port=settings.smtp_port,
|
||||||
|
imap_host=host,
|
||||||
|
imap_port=settings.imap_port,
|
||||||
|
username=settings.username,
|
||||||
|
password=settings.password,
|
||||||
|
sender=settings.sender,
|
||||||
|
recipient=settings.recipient,
|
||||||
|
sent_folder=settings.sent_folder,
|
||||||
|
provider_timeout_seconds=settings.provider_timeout_seconds,
|
||||||
|
)
|
||||||
|
|
||||||
|
accepted.assert_local_testbed()
|
||||||
|
profile = runner._profile_payload(accepted, name="Literal loopback")
|
||||||
|
assert profile["smtp"]["host"] == host
|
||||||
|
assert profile["imap"]["host"] == host
|
||||||
|
|
||||||
|
|
||||||
|
def test_campaign_acceptance_orchestration_retains_only_profile_reference_and_safe_evidence() -> None:
|
||||||
|
client = _AcceptanceClient()
|
||||||
|
|
||||||
|
def snapshot_probe(_version_id: str):
|
||||||
|
assert client.campaign_json is not None
|
||||||
|
return client.campaign_json, {
|
||||||
|
"mail_profile_id": "profile-1",
|
||||||
|
"smtp_transport_revision": "opaque-smtp-revision",
|
||||||
|
"imap_transport_revision": "opaque-imap-revision",
|
||||||
|
"delivery": {"imap_append_sent": {"enabled": True, "folder": "Sent"}},
|
||||||
|
}
|
||||||
|
|
||||||
|
audit = {
|
||||||
|
"campaign.created": 1,
|
||||||
|
"campaign.validated": 1,
|
||||||
|
"campaign.messages_built": 1,
|
||||||
|
"campaign.sent_now": 1,
|
||||||
|
"campaign.send_now_rejected": 1,
|
||||||
|
"campaign.append_sent_enqueued": 1,
|
||||||
|
}
|
||||||
|
evidence, subject = runner.execute_campaign_scenario(
|
||||||
|
client,
|
||||||
|
{"Authorization": "not-retained"},
|
||||||
|
fixture_path=FIXTURE_PATH,
|
||||||
|
profile_id="profile-1",
|
||||||
|
settings=_settings(),
|
||||||
|
scenario="success",
|
||||||
|
snapshot_probe=snapshot_probe,
|
||||||
|
audit_probe=lambda _campaign_id, _version_id: audit,
|
||||||
|
append_sent=True,
|
||||||
|
repeat_send=True,
|
||||||
|
)
|
||||||
|
evidence["provider_verification"] = {
|
||||||
|
"inbox_increment": 1,
|
||||||
|
"sent_increment": 1,
|
||||||
|
"unique_subject_matches_in_inbox": 1,
|
||||||
|
"unique_subject_matches_in_sent": 1,
|
||||||
|
}
|
||||||
|
|
||||||
|
runner._assert_success_evidence(evidence)
|
||||||
|
runner._assert_evidence_safe(
|
||||||
|
{
|
||||||
|
"schema_version": runner.EVIDENCE_SCHEMA,
|
||||||
|
"coverage": {
|
||||||
|
"smtp_acceptance": True,
|
||||||
|
"partial_envelope_refusal": False,
|
||||||
|
"post_data_connection_loss_outcome_unknown": False,
|
||||||
|
"source_artifact_provenance": False,
|
||||||
|
"worker_restart_interruption": False,
|
||||||
|
},
|
||||||
|
"success": evidence,
|
||||||
|
},
|
||||||
|
settings=_settings(),
|
||||||
|
)
|
||||||
|
assert subject.startswith("[GovOPlaN acceptance ")
|
||||||
|
assert client.send_calls == 2
|
||||||
|
assert client.campaign_json is not None
|
||||||
|
assert client.campaign_json["server"] == {"mail_profile_id": "profile-1"}
|
||||||
|
assert "credentials" not in json.dumps(client.campaign_json).casefold()
|
||||||
|
serialized = json.dumps(evidence, sort_keys=True)
|
||||||
|
assert "not-retained" not in serialized
|
||||||
|
assert "local-test-password" not in serialized
|
||||||
|
|
||||||
|
|
||||||
|
def test_campaign_boundary_rejects_resolved_transport_material() -> None:
|
||||||
|
with pytest.raises(runner.AcceptanceError, match="forbidden transport material"):
|
||||||
|
runner.assert_campaign_boundary(
|
||||||
|
{"server": {"mail_profile_id": "profile-1"}},
|
||||||
|
{
|
||||||
|
"mail_profile_id": "profile-1",
|
||||||
|
"smtp_transport_revision": "smtp-revision",
|
||||||
|
"imap_transport_revision": "imap-revision",
|
||||||
|
"smtp": {"host": "should-not-be-here"},
|
||||||
|
},
|
||||||
|
profile_id="profile-1",
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def test_success_projection_fails_closed_on_inconsistent_campaign_report() -> None:
|
||||||
|
evidence = {
|
||||||
|
"send": {
|
||||||
|
"attempted_count": 1,
|
||||||
|
"sent_count": 1,
|
||||||
|
"failed_count": 0,
|
||||||
|
"outcome_unknown_count": 0,
|
||||||
|
"skipped_count": 0,
|
||||||
|
"delivery_mode": "synchronous",
|
||||||
|
"statuses": {"smtp_accepted": 1},
|
||||||
|
},
|
||||||
|
"append_sent": {
|
||||||
|
"pending_count": 1,
|
||||||
|
"processed_count": 1,
|
||||||
|
"appended_count": 1,
|
||||||
|
"failed_count": 0,
|
||||||
|
"skipped_count": 0,
|
||||||
|
"statuses": {"appended": 1},
|
||||||
|
},
|
||||||
|
"report": {
|
||||||
|
"cards": {
|
||||||
|
"jobs_total": 1,
|
||||||
|
"sent": 0,
|
||||||
|
"smtp_accepted": 0,
|
||||||
|
"failed": 0,
|
||||||
|
"outcome_unknown": 0,
|
||||||
|
"needs_attention": 0,
|
||||||
|
"imap_appended": 0,
|
||||||
|
"imap_failed": 0,
|
||||||
|
},
|
||||||
|
"send_status_counts": {},
|
||||||
|
"imap_status_counts": {},
|
||||||
|
},
|
||||||
|
"provider_verification": {
|
||||||
|
"inbox_increment": 1,
|
||||||
|
"sent_increment": 1,
|
||||||
|
"unique_subject_matches_in_inbox": 1,
|
||||||
|
"unique_subject_matches_in_sent": 1,
|
||||||
|
},
|
||||||
|
"campaign_mail_boundary": {
|
||||||
|
"profile_reference_only": True,
|
||||||
|
"smtp_revision_frozen": True,
|
||||||
|
"imap_revision_frozen": True,
|
||||||
|
"resolved_transport_material_present": False,
|
||||||
|
},
|
||||||
|
}
|
||||||
|
|
||||||
|
with pytest.raises(runner.AcceptanceError, match="report does not agree"):
|
||||||
|
runner._assert_success_evidence(evidence)
|
||||||
|
|
||||||
|
|
||||||
|
def test_evidence_projection_rejects_unknown_status_keys() -> None:
|
||||||
|
with pytest.raises(runner.AcceptanceError, match="unsupported status"):
|
||||||
|
runner._send_evidence(
|
||||||
|
{
|
||||||
|
"delivery_mode": "synchronous",
|
||||||
|
"results": [{"status": "provider diagnostic: recipient@example.test"}],
|
||||||
|
}
|
||||||
|
)
|
||||||
|
|
||||||
|
with pytest.raises(runner.AcceptanceError, match="durable attempt status"):
|
||||||
|
runner._durable_state_evidence(
|
||||||
|
{
|
||||||
|
"job_count": 1,
|
||||||
|
"send_status_counts": {"sending": 1},
|
||||||
|
"attempt_status_counts": {"provider-secret": 1},
|
||||||
|
"unfinished_attempt_count": 1,
|
||||||
|
}
|
||||||
|
)
|
||||||
|
|
||||||
|
with pytest.raises(runner.AcceptanceError, match="synchronous delivery mode"):
|
||||||
|
runner._send_evidence(
|
||||||
|
{
|
||||||
|
"delivery_mode": "provider diagnostic: recipient@example.test",
|
||||||
|
"results": [],
|
||||||
|
}
|
||||||
|
)
|
||||||
|
|
||||||
|
with pytest.raises(runner.AcceptanceError, match="unsupported"):
|
||||||
|
runner._report_evidence(
|
||||||
|
{
|
||||||
|
"cards": {},
|
||||||
|
"status_counts": {
|
||||||
|
"send": {"smtp_accepted": 1, "provider-secret": 1},
|
||||||
|
"imap": {},
|
||||||
|
},
|
||||||
|
}
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def _open_fault_smtp(endpoint):
|
||||||
|
client = smtplib.SMTP(endpoint.host, endpoint.port, timeout=5)
|
||||||
|
client.ehlo()
|
||||||
|
client.login("acceptance-user", "acceptance-password")
|
||||||
|
return client
|
||||||
|
|
||||||
|
|
||||||
|
def test_explicit_temporary_smtp_response_occurs_after_complete_data() -> None:
|
||||||
|
with runner.smtp_fault_endpoint("temporary_data_response") as endpoint:
|
||||||
|
client = _open_fault_smtp(endpoint)
|
||||||
|
try:
|
||||||
|
with pytest.raises(smtplib.SMTPDataError) as captured:
|
||||||
|
client.sendmail(
|
||||||
|
"sender@example.test",
|
||||||
|
["recipient@example.test"],
|
||||||
|
b"Subject: temporary\r\n\r\nmessage",
|
||||||
|
)
|
||||||
|
finally:
|
||||||
|
client.close()
|
||||||
|
assert captured.value.smtp_code == 451
|
||||||
|
assert endpoint.evidence() == {
|
||||||
|
"connection_count": 1,
|
||||||
|
"accepted_rcpt_commands": 1,
|
||||||
|
"refused_rcpt_commands": 0,
|
||||||
|
"data_transactions": 1,
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
|
def test_partial_recipient_refusal_retains_one_accepted_envelope() -> None:
|
||||||
|
with runner.smtp_fault_endpoint("partial_recipient_refusal") as endpoint:
|
||||||
|
client = _open_fault_smtp(endpoint)
|
||||||
|
try:
|
||||||
|
refused = client.sendmail(
|
||||||
|
"sender@example.test",
|
||||||
|
["accepted@example.test", "refused@example.test"],
|
||||||
|
b"Subject: partial\r\n\r\nmessage",
|
||||||
|
)
|
||||||
|
finally:
|
||||||
|
client.quit()
|
||||||
|
assert set(refused) == {"refused@example.test"}
|
||||||
|
assert endpoint.evidence() == {
|
||||||
|
"connection_count": 1,
|
||||||
|
"accepted_rcpt_commands": 1,
|
||||||
|
"refused_rcpt_commands": 1,
|
||||||
|
"data_transactions": 1,
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
|
def test_post_data_disconnect_is_a_real_ambiguous_protocol_boundary() -> None:
|
||||||
|
with runner.smtp_fault_endpoint("post_data_disconnect") as endpoint:
|
||||||
|
client = _open_fault_smtp(endpoint)
|
||||||
|
try:
|
||||||
|
with pytest.raises(smtplib.SMTPServerDisconnected):
|
||||||
|
client.sendmail(
|
||||||
|
"sender@example.test",
|
||||||
|
["recipient@example.test"],
|
||||||
|
b"Subject: ambiguous\r\n\r\nmessage",
|
||||||
|
)
|
||||||
|
finally:
|
||||||
|
client.close()
|
||||||
|
assert endpoint.wait_for_data(1)
|
||||||
|
assert endpoint.evidence()["data_transactions"] == 1
|
||||||
|
|
||||||
|
|
||||||
|
def test_partial_refusal_fixture_adds_a_second_distinct_recipient() -> None:
|
||||||
|
raw, _subject = runner.materialize_campaign_fixture(
|
||||||
|
FIXTURE_PATH,
|
||||||
|
profile_id="profile-1",
|
||||||
|
settings=_settings(),
|
||||||
|
scenario="partial_envelope_refusal",
|
||||||
|
run_token="0123456789ab",
|
||||||
|
additional_envelope_recipient=True,
|
||||||
|
)
|
||||||
|
|
||||||
|
recipients = raw["entries"]["inline"][0]["to"]
|
||||||
|
assert len(recipients) == 2
|
||||||
|
assert recipients[0]["email"] != recipients[1]["email"]
|
||||||
|
assert recipients[1]["email"].endswith("@govoplan.test")
|
||||||
|
|
||||||
|
def test_fixture_contains_no_transport_credentials() -> None:
|
||||||
|
raw = json.loads(FIXTURE_PATH.read_text(encoding="utf-8"))
|
||||||
|
runner._assert_no_forbidden_campaign_keys(raw)
|
||||||
|
assert raw["server"] == {"mail_profile_id": "00000000-0000-4000-8000-000000000001"}
|
||||||
|
config = load_campaign_config(FIXTURE_PATH)
|
||||||
|
assert config.server.mail_profile_id == "00000000-0000-4000-8000-000000000001"
|
||||||
|
|
||||||
|
|
||||||
|
def test_fixture_composition_versions_are_complete_and_exact() -> None:
|
||||||
|
versions = {
|
||||||
|
"core": "0.1.13",
|
||||||
|
"access": "0.1.11",
|
||||||
|
"audit": "0.1.8",
|
||||||
|
"campaigns": "0.1.11",
|
||||||
|
"mail": "0.1.10",
|
||||||
|
"files": "0.1.9",
|
||||||
|
}
|
||||||
|
|
||||||
|
assert runner.required_composition_versions(FIXTURE_PATH, versions) == {
|
||||||
|
"core": "0.1.13",
|
||||||
|
"access": "0.1.11",
|
||||||
|
"audit": "0.1.8",
|
||||||
|
"campaigns": "0.1.11",
|
||||||
|
"mail": "0.1.10",
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
|
def test_fixture_composition_fails_closed_when_a_required_version_is_missing() -> None:
|
||||||
|
with pytest.raises(runner.AcceptanceError, match="versions are unavailable"):
|
||||||
|
runner.required_composition_versions(
|
||||||
|
FIXTURE_PATH,
|
||||||
|
{
|
||||||
|
"core": "0.1.13",
|
||||||
|
"access": "0.1.11",
|
||||||
|
"campaigns": "0.1.11",
|
||||||
|
"mail": "0.1.10",
|
||||||
|
},
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
def test_testbed_documentation_distinguishes_proven_and_open_failure_drills() -> 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 "second ordinary send must be rejected before another provider effect" in testbed
|
||||||
|
assert "connection loss after complete DATA is frozen" in testbed
|
||||||
|
assert "dedicated OS process" in testbed
|
||||||
|
assert "Redis/Celery delivery" in runbook
|
||||||
|
assert "broker redelivery" in runbook
|
||||||
|
assert "celery_broker_redelivery" in testbed
|
||||||
|
assert "raw provider diagnostics" in runbook
|
||||||
@@ -10,6 +10,7 @@ from govoplan_campaign.backend.db.models import (
|
|||||||
JobSendStatus,
|
JobSendStatus,
|
||||||
JobValidationStatus,
|
JobValidationStatus,
|
||||||
)
|
)
|
||||||
|
from govoplan_campaign.backend.capabilities import delivery_tasks_capability
|
||||||
from govoplan_campaign.backend.sending.jobs import (
|
from govoplan_campaign.backend.sending.jobs import (
|
||||||
SendJobResult,
|
SendJobResult,
|
||||||
_queue_validation_statuses,
|
_queue_validation_statuses,
|
||||||
@@ -48,6 +49,17 @@ def _job(entry_id: str, **overrides):
|
|||||||
|
|
||||||
|
|
||||||
class CampaignQueueSelectionTests(unittest.TestCase):
|
class CampaignQueueSelectionTests(unittest.TestCase):
|
||||||
|
def test_delivery_task_capability_configures_module_runtime_for_worker_processes(self):
|
||||||
|
registry = object()
|
||||||
|
settings = object()
|
||||||
|
context = SimpleNamespace(registry=registry, settings=settings)
|
||||||
|
|
||||||
|
with patch("govoplan_campaign.backend.runtime.configure_runtime") as configure:
|
||||||
|
capability = delivery_tasks_capability(context)
|
||||||
|
|
||||||
|
self.assertIsNotNone(capability)
|
||||||
|
configure.assert_called_once_with(registry=registry, settings=settings)
|
||||||
|
|
||||||
def test_selects_queueable_jobs_without_reclassifying_retry_states(self):
|
def test_selects_queueable_jobs_without_reclassifying_retry_states(self):
|
||||||
skipped_send = _job("1", send_status=JobSendStatus.FAILED_TEMPORARY.value)
|
skipped_send = _job("1", send_status=JobSendStatus.FAILED_TEMPORARY.value)
|
||||||
skipped_queue = _job("2", queue_status=JobQueueStatus.PAUSED.value)
|
skipped_queue = _job("2", queue_status=JobQueueStatus.PAUSED.value)
|
||||||
|
|||||||
@@ -195,7 +195,10 @@ def test_post_queue_growth_is_rejected_before_batch_or_provider_preflight() -> N
|
|||||||
patch("govoplan_campaign.backend.sending.jobs._ensure_campaign_execution_snapshot"),
|
patch("govoplan_campaign.backend.sending.jobs._ensure_campaign_execution_snapshot"),
|
||||||
patch("govoplan_campaign.backend.sending.jobs.effective_synchronous_send_policy", return_value=policy),
|
patch("govoplan_campaign.backend.sending.jobs.effective_synchronous_send_policy", return_value=policy),
|
||||||
patch("govoplan_campaign.backend.sending.jobs._campaign_jobs_for_queue", return_value=initial_jobs),
|
patch("govoplan_campaign.backend.sending.jobs._campaign_jobs_for_queue", return_value=initial_jobs),
|
||||||
patch("govoplan_campaign.backend.sending.jobs.queue_campaign_jobs", return_value=queued),
|
patch(
|
||||||
|
"govoplan_campaign.backend.sending.jobs.queue_campaign_jobs",
|
||||||
|
return_value=queued,
|
||||||
|
) as queue,
|
||||||
patch("govoplan_campaign.backend.sending.jobs._campaign_jobs_for_version", return_value=post_queue_jobs),
|
patch("govoplan_campaign.backend.sending.jobs._campaign_jobs_for_version", return_value=post_queue_jobs),
|
||||||
patch("govoplan_campaign.backend.sending.jobs._preflight_synchronous_send_batch") as batch_preflight,
|
patch("govoplan_campaign.backend.sending.jobs._preflight_synchronous_send_batch") as batch_preflight,
|
||||||
):
|
):
|
||||||
@@ -207,6 +210,7 @@ def test_post_queue_growth_is_rejected_before_batch_or_provider_preflight() -> N
|
|||||||
)
|
)
|
||||||
|
|
||||||
assert rejected.value.eligible_count == 3
|
assert rejected.value.eligible_count == 3
|
||||||
|
assert queue.call_args.kwargs["commit_queue"] is False
|
||||||
batch_preflight.assert_not_called()
|
batch_preflight.assert_not_called()
|
||||||
|
|
||||||
|
|
||||||
@@ -251,6 +255,7 @@ def test_asynchronous_mode_matches_actual_worker_availability(
|
|||||||
assert result.worker_queue_available is workers_available
|
assert result.worker_queue_available is workers_available
|
||||||
assert result.enqueued_count == expected_enqueued
|
assert result.enqueued_count == expected_enqueued
|
||||||
assert persist.call_args.kwargs["delivery_mode"] == expected_mode
|
assert persist.call_args.kwargs["delivery_mode"] == expected_mode
|
||||||
|
assert persist.call_args.kwargs["commit"] is True
|
||||||
assert enqueue.call_args.kwargs["enabled"] is workers_available
|
assert enqueue.call_args.kwargs["enabled"] is workers_available
|
||||||
|
|
||||||
|
|
||||||
@@ -356,3 +361,45 @@ def test_batch_preflight_checks_every_message_before_provider_effects() -> None:
|
|||||||
assert state_preflight.call_count == 2
|
assert state_preflight.call_count == 2
|
||||||
assert input_preflight.call_count == 2
|
assert input_preflight.call_count == 2
|
||||||
provider.send_campaign_email_bytes.assert_not_called()
|
provider.send_campaign_email_bytes.assert_not_called()
|
||||||
|
|
||||||
|
|
||||||
|
def test_rejected_synchronous_preflight_rolls_back_staged_queue_before_audit() -> None:
|
||||||
|
session = Mock()
|
||||||
|
campaign = SimpleNamespace(id="campaign-1", current_version_id="version-1")
|
||||||
|
version = SimpleNamespace(
|
||||||
|
id="version-1",
|
||||||
|
raw_json={},
|
||||||
|
locked_at=object(),
|
||||||
|
validation_summary={"ok": True},
|
||||||
|
build_summary={"built_count": 1},
|
||||||
|
)
|
||||||
|
rejection = SynchronousSendRejected(
|
||||||
|
"Preflight rejected the staged send.",
|
||||||
|
reason="batch_preflight_failed",
|
||||||
|
eligible_count=1,
|
||||||
|
)
|
||||||
|
|
||||||
|
with (
|
||||||
|
patch.object(router, "_get_campaign_for_principal"),
|
||||||
|
patch.object(router, "_require_permission"),
|
||||||
|
patch.object(router, "_get_campaign_for_tenant", return_value=campaign),
|
||||||
|
patch.object(router, "_get_version_for_tenant", return_value=version),
|
||||||
|
patch.object(router, "_require_mail_profile_use_if_needed"),
|
||||||
|
patch.object(router, "is_user_locked_version", return_value=False),
|
||||||
|
patch.object(router, "send_campaign_now", side_effect=rejection),
|
||||||
|
patch.object(router, "audit_from_principal") as audit,
|
||||||
|
pytest.raises(HTTPException) as rejected,
|
||||||
|
):
|
||||||
|
router.send_campaign_now_endpoint(
|
||||||
|
"campaign-1",
|
||||||
|
session=session,
|
||||||
|
principal=_Principal(
|
||||||
|
"campaigns:campaign:send", "campaigns:recipient:read"
|
||||||
|
), # type: ignore[arg-type]
|
||||||
|
)
|
||||||
|
|
||||||
|
assert rejected.value.status_code == 422
|
||||||
|
session.rollback.assert_called_once_with()
|
||||||
|
audit.assert_called_once()
|
||||||
|
assert audit.call_args.kwargs["action"] == "campaign.send_now_rejected"
|
||||||
|
assert audit.call_args.kwargs["commit"] is True
|
||||||
|
|||||||
@@ -1,6 +1,6 @@
|
|||||||
{
|
{
|
||||||
"name": "@govoplan/campaign-webui",
|
"name": "@govoplan/campaign-webui",
|
||||||
"version": "0.1.10",
|
"version": "0.1.11",
|
||||||
"private": true,
|
"private": true,
|
||||||
"type": "module",
|
"type": "module",
|
||||||
"main": "src/index.ts",
|
"main": "src/index.ts",
|
||||||
|
|||||||
Reference in New Issue
Block a user