Files
govoplan-campaign/tests/test_testbed_claim_recovery.py
T
zemion c51fc180fb
Module Package Release / publish-packages (push) Successful in 12s
Release govoplan-campaign v0.1.28: stabilize saving, review and delivery recovery
2026-09-08 01:32:26 +02:00

184 lines
9.7 KiB
Python

"""Disposable SQLite and in-process HTTP only; never run Docker or a mail provider."""
from __future__ import annotations
import sqlite3
from types import SimpleNamespace
from unittest.mock import Mock
import pytest
from fastapi import FastAPI
from fastapi.testclient import TestClient
from sqlalchemy import create_engine
from sqlalchemy.orm import sessionmaker
from govoplan_campaign.backend.db.models import Campaign, CampaignVersion, SendAttempt
from govoplan_campaign.backend.routes import delivery as routes
from govoplan_campaign.backend.sending import jobs
from govoplan_core.auth import ApiPrincipal, get_api_principal
from govoplan_core.core.access import PrincipalRef
from govoplan_core.core.recovery import RecoveryOperation, verify_recovery_evidence_chain
from govoplan_core.db.session import get_session
from test_mail_testbed_acceptance import runner, _settings, _Response, FIXTURE_PATH
from test_workerless_recovery import recovery, _stale_claim # noqa: F401
@pytest.fixture
def fixture_claim(recovery, tmp_path, monkeypatch):
root = tmp_path / "govoplan-campaign-greenmail-regression"
root.mkdir()
database_file = root / "acceptance.db"
with recovery.engine.connect() as source, sqlite3.connect(database_file) as target:
source.connection.driver_connection.backup(target)
engine = create_engine(f"sqlite:///{database_file}")
factory = sessionmaker(engine)
session = factory()
fixture = SimpleNamespace(
session=session, factory=factory, engine=engine, root=root,
campaign=session.get(Campaign, "campaign"), version=session.get(CampaignVersion, "version"),
snapshot=recovery.snapshot, provider=recovery.provider,
)
monkeypatch.setattr(jobs, "get_database", lambda: SimpleNamespace(SessionLocal=factory))
job, lease, node, operation_id = _stale_claim(fixture, expired=False, owner="active")
node.metadata_ = {"acceptance_worker_pid": 12345}
session.commit()
fixture.job, fixture.lease, fixture.node, fixture.operation_id = job, lease, node, operation_id
fixture.process = SimpleNamespace(pid=12345, returncode=-9, poll=lambda: -9)
monkeypatch.setenv("APP_ENV", "test")
monkeypatch.setenv("DATABASE_URL", f"sqlite:///{database_file}")
monkeypatch.setattr(runner.tempfile, "gettempdir", lambda: str(tmp_path))
yield fixture
session.close()
engine.dispose()
def _recover(fixture, client):
return runner.recover_stopped_fixture_claim(
client, {}, database=SimpleNamespace(engine=fixture.engine, SessionLocal=fixture.factory),
runtime_root=fixture.root, campaign_id="campaign", version_id="version", stopped_process=fixture.process,
)
def test_fixture_proof_calls_actual_fenced_http_action_without_replaying_smtp(fixture_claim, monkeypatch):
fixture = fixture_claim
app = FastAPI()
app.include_router(routes.router, prefix="/api/v1")
actor = SimpleNamespace(id="operator")
scopes = frozenset({"campaigns:campaign:reconcile", "campaigns:recipient:read"})
app.dependency_overrides[get_api_principal] = lambda: ApiPrincipal(
principal=PrincipalRef(account_id="operator", membership_id="operator", tenant_id="tenant", scopes=scopes),
user=actor, account=actor,
)
def session_dependency():
with fixture.factory() as session:
yield session
app.dependency_overrides[get_session] = session_dependency
monkeypatch.setattr(routes, "_get_campaign_for_principal", lambda session, *_args, **_kwargs: session.get(Campaign, "campaign"))
monkeypatch.setattr(routes, "audit_from_principal", lambda session, *_args, **_kwargs: session.commit())
with TestClient(app) as client:
evidence = _recover(fixture, client)
assert evidence == {"stopped_process_verified": True, "fixture_lease_expired": True, "explicit_fenced_recovery": True}
fixture.session.expire_all()
assert fixture.job.send_status == "outcome_unknown"
assert fixture.job.claim_token is None
assert fixture.session.query(SendAttempt).one().status == "outcome_unknown"
assert fixture.session.query(SendAttempt).one().finished_at is not None
assert fixture.session.get(RecoveryOperation, fixture.operation_id).status == "outcome_unknown"
assert verify_recovery_evidence_chain(fixture.session, fixture.operation_id)
assert fixture.lease.fencing_token > 1
fixture.provider.send_campaign_email_bytes.assert_not_called()
@pytest.mark.parametrize("unsafe", ["production", "foreign_database", "live_process", "wrong_pid", "different_incarnation"])
def test_fixture_claim_proof_rejects_unsafe_target_or_process_without_state_changes(fixture_claim, monkeypatch, unsafe):
fixture = fixture_claim
if unsafe == "production":
monkeypatch.setenv("APP_ENV", "production")
elif unsafe == "foreign_database":
monkeypatch.setenv("DATABASE_URL", "sqlite:////tmp/a-different-database.db")
elif unsafe == "live_process":
fixture.process = SimpleNamespace(pid=12345, returncode=None, poll=lambda: None)
elif unsafe == "wrong_pid":
fixture.process = SimpleNamespace(pid=22222, returncode=-9, poll=lambda: -9)
else:
fixture.node.incarnation = "new-worker-incarnation"
fixture.session.commit()
original_expiry = fixture.lease.expires_at
client = Mock()
with pytest.raises(runner.AcceptanceError):
_recover(fixture, client)
client.post.assert_not_called()
fixture.session.expire_all()
assert fixture.job.send_status == "sending"
assert fixture.job.claim_token == "stale-token"
assert fixture.node.state == "active"
assert fixture.lease.expires_at == original_expiry
assert fixture.session.get(RecoveryOperation, fixture.operation_id).status == "running"
@pytest.mark.parametrize("duplicate_mutates", [False, True])
def test_process_restart_requires_unchanged_claim_before_explicit_recovery(monkeypatch, duplicate_mutates):
interrupted = {"job_count": 1, "send_status_counts": {"sending": 1}, "attempt_status_counts": {"smtp_in_progress": 1}, "unfinished_attempt_count": 1}
recovered = {"job_count": 1, "send_status_counts": {"outcome_unknown": 1}, "attempt_status_counts": {"outcome_unknown": 1}, "unfinished_attempt_count": 0}
states = iter([interrupted, recovered if duplicate_mutates else interrupted, recovered])
first, second = Mock(), Mock()
first.poll.return_value = -9
workers = iter([first, second])
prepared = SimpleNamespace(campaign_id="campaign", version_id="version", public_evidence=lambda: {})
monkeypatch.setattr(runner, "prepare_campaign_scenario", lambda *_args, **_kwargs: prepared)
monkeypatch.setattr(runner, "_start_campaign_worker_task", lambda *_args: next(workers))
monkeypatch.setattr(runner, "_terminate_worker_process", lambda *_args: None)
monkeypatch.setattr(runner, "_wait_for_worker_process", lambda *_args, **_kwargs: None)
client = Mock()
client.post.return_value = _Response(200, {
"queued_count": 1, "skipped_count": 0, "blocked_count": 0, "enqueued_count": 0,
"delivery_mode": "database_queue", "worker_queue_available": False, "dry_run": False,
})
client.get.return_value = _Response(200, {"cards": {}, "status_counts": {"send": {"outcome_unknown": 1}, "imap": {}}})
endpoint = Mock()
endpoint.wait_for_data.return_value = True
endpoint.evidence.return_value = {"connection_count": 1, "accepted_rcpt_commands": 1, "refused_rcpt_commands": 0, "data_transactions": 1}
recover_claim = Mock(return_value={"explicit_fenced_recovery": True})
arguments = dict(
fixture_path=FIXTURE_PATH, profile_id="profile", settings=_settings(), endpoint=endpoint,
snapshot_probe=lambda _: ({}, {}), audit_probe=lambda *_: {"campaign.created": 1, "campaign.validated": 1, "campaign.messages_built": 1, "campaign.queued": 1},
delivery_probe=lambda *_: next(states), worker_job_probe=lambda _: "job", recover_claim=recover_claim,
)
if duplicate_mutates:
with pytest.raises(runner.AcceptanceError, match="without stopped-runtime proof"):
runner.execute_worker_interruption_scenario(client, {}, **arguments)
recover_claim.assert_not_called()
else:
evidence = runner.execute_worker_interruption_scenario(client, {}, **arguments)
assert evidence["interrupted_durable_state"] == evidence["restarted_durable_state"] == interrupted
assert evidence["recovered_durable_state"] == recovered
recover_claim.assert_called_once_with("campaign", "version", first)
def test_direct_task_bootstrap_registers_unique_fixture_owner_and_accepts_read_only_duplicate(monkeypatch):
import os
import sys
from unittest.mock import MagicMock
from govoplan_core import celery_app, db
from govoplan_core.core import runtime_coordination
identity = SimpleNamespace(node_id="fixture-owner")
bind = Mock(return_value=identity)
register = Mock()
session = MagicMock()
database = SimpleNamespace(SessionLocal=MagicMock())
database.SessionLocal.return_value.__enter__.return_value = session
task = SimpleNamespace(run=Mock(return_value={"status": "already_sending"}))
monkeypatch.setattr(celery_app, "_worker_runtime_identity", bind)
monkeypatch.setattr(celery_app, "send_email", task)
monkeypatch.setattr(runtime_coordination, "register_runtime_node", register)
monkeypatch.setattr(db.session, "get_database", lambda: database)
monkeypatch.setattr(sys, "argv", ["fixture-worker", "fixture-job"])
exec(compile(runner.WORKER_TASK_CODE, "<isolated-worker-test>", "exec"), {})
assert bind.call_args.args[0].hostname == f"campaign-acceptance-{os.getpid()}"
register.assert_called_once_with(session, identity, metadata={"acceptance_worker_pid": os.getpid()})
session.commit.assert_called_once()
task.run.assert_called_once_with("fixture-job")