Files
govoplan-campaign/tests/test_artifact_reconciliation.py

339 lines
11 KiB
Python

from __future__ import annotations
from datetime import datetime, timedelta, timezone
from unittest.mock import patch
import pytest
from sqlalchemy import create_engine, select
from sqlalchemy.orm import sessionmaker
from sqlalchemy.pool import StaticPool
from govoplan_campaign.backend.artifact_reconciliation import (
campaign_artifact_inventory,
reconcile_campaign_artifacts,
)
from govoplan_core.core.object_storage import (
StorageBackendError,
StorageObjectInfo,
StorageObjectMissing,
StorageObjectPage,
)
from govoplan_core.core.recovery import RecoveryCheckpoint, RecoveryOperation
from govoplan_core.core.recovery_runtime import RecoveryOperationBusy
from govoplan_core.core.runtime_coordination import (
DistributedLease,
RuntimeIdentity,
acquire_lease,
)
NOW = datetime(2026, 8, 3, 12, 0, tzinfo=timezone.utc)
PREFIX = "campaign-artifacts/tenant-1/campaign-1/version-1/"
class _MemoryStorage:
name = "memory"
def __init__(self) -> None:
self.objects: dict[str, bytes] = {}
self.modified_at: dict[str, datetime | None] = {}
self.delete_failures: set[str] = set()
self.exists_failures: set[str] = set()
def add(
self,
key: str,
*,
payload: bytes = b"artifact",
modified_at: datetime | None = NOW - timedelta(days=2),
) -> None:
self.objects[key] = payload
self.modified_at[key] = modified_at
def put_bytes(self, key: str, data: bytes, **_kwargs) -> None:
self.add(key, payload=data, modified_at=NOW)
def get_bytes(self, key: str) -> bytes:
try:
return self.objects[key]
except KeyError as exc:
raise StorageObjectMissing("missing") from exc
def iter_bytes(self, key: str, *, chunk_size: int = 1024 * 1024):
del chunk_size
yield self.get_bytes(key)
def delete(self, key: str) -> None:
if key in self.delete_failures:
raise StorageBackendError("delete unavailable")
self.objects.pop(key, None)
self.modified_at.pop(key, None)
def exists(self, key: str) -> bool:
if key in self.exists_failures:
raise StorageBackendError("probe unavailable")
return key in self.objects
def stat(self, key: str) -> StorageObjectInfo:
if key not in self.objects:
raise StorageObjectMissing("missing")
return StorageObjectInfo(
key=key,
size_bytes=len(self.objects[key]),
modified_at=self.modified_at[key],
)
def list_objects(
self,
*,
prefix: str,
after: str | None = None,
limit: int = 500,
) -> StorageObjectPage:
keys = [
key
for key in sorted(self.objects)
if key.startswith(prefix) and (after is None or key > after)
]
selected = keys[:limit]
return StorageObjectPage(
objects=tuple(self.stat(key) for key in selected),
next_cursor=(selected[-1] if len(keys) > len(selected) else None),
)
@pytest.fixture
def recovery_session_factory():
engine = create_engine(
"sqlite://",
connect_args={"check_same_thread": False},
poolclass=StaticPool,
)
DistributedLease.__table__.create(engine)
RecoveryOperation.__table__.create(engine)
RecoveryCheckpoint.__table__.create(engine)
factory = sessionmaker(bind=engine, expire_on_commit=False)
try:
yield factory
finally:
engine.dispose()
def _identity(*, node_id: str = "node-1", incarnation: str = "run-1"):
return RuntimeIdentity(
installation_id="campaign-artifact-tests",
node_id=node_id,
incarnation=incarnation,
role="worker",
software_version="test",
composition_hash="a" * 64,
)
def _without_domain_references():
return (
patch(
"govoplan_campaign.backend.artifact_reconciliation._referenced_artifact_keys",
return_value=set(),
),
patch(
"govoplan_campaign.backend.artifact_reconciliation._active_build_ids",
return_value=set(),
),
)
def test_inventory_classifies_reference_grace_active_build_and_unknown_age() -> None:
storage = _MemoryStorage()
orphan = f"{PREFIX}build-orphan/message.eml"
referenced = f"{PREFIX}build-referenced/message.eml"
active = f"{PREFIX}build-active/message.eml"
young = f"{PREFIX}build-young/message.eml"
unknown_age = f"{PREFIX}build-unknown/message.eml"
malformed = "campaign-artifacts/tenant-1/not-a-build-object"
storage.add(orphan)
storage.add(referenced)
storage.add(active)
storage.add(young, modified_at=NOW - timedelta(hours=1))
storage.add(unknown_age, modified_at=None)
storage.add(malformed)
with (
patch(
"govoplan_campaign.backend.artifact_reconciliation._referenced_artifact_keys",
return_value={referenced},
),
patch(
"govoplan_campaign.backend.artifact_reconciliation._active_build_ids",
return_value={"build-active"},
),
):
inventory = campaign_artifact_inventory(
object(), # type: ignore[arg-type]
storage=storage,
tenant_id="tenant-1",
grace_period=timedelta(hours=24),
page_size=20,
now=NOW,
)
assert [candidate.key for candidate in inventory.candidates] == [orphan]
assert inventory.referenced_count == 1
assert inventory.active_build_count == 1
assert inventory.young_count == 1
assert inventory.unknown_age_count == 1
assert inventory.invalid_shape_count == 1
def test_process_loss_orphan_is_deleted_once_and_same_request_replays(
recovery_session_factory,
) -> None:
storage = _MemoryStorage()
orphan = f"{PREFIX}lost-build/message.eml"
storage.add(orphan)
reference_patch, active_patch = _without_domain_references()
with reference_patch, active_patch:
dry_run = reconcile_campaign_artifacts(
recovery_session_factory,
storage=storage,
identity=_identity(),
tenant_id="tenant-1",
now=NOW,
)
applied = reconcile_campaign_artifacts(
recovery_session_factory,
storage=storage,
identity=_identity(),
tenant_id="tenant-1",
apply=True,
idempotency_key="cleanup-lost-build",
now=NOW,
)
replayed = reconcile_campaign_artifacts(
recovery_session_factory,
storage=storage,
identity=_identity(),
tenant_id="tenant-1",
apply=True,
idempotency_key="cleanup-lost-build",
now=NOW,
)
repeated = reconcile_campaign_artifacts(
recovery_session_factory,
storage=storage,
identity=_identity(),
tenant_id="tenant-1",
apply=True,
idempotency_key="cleanup-empty-page",
now=NOW,
)
assert dry_run["candidate_count"] == 1
assert dry_run["deleted_count"] == 0
assert applied["status"] == "applied"
assert applied["deleted_count"] == 1
assert orphan not in storage.objects
assert replayed["status"] == "already_completed"
assert repeated["candidate_count"] == 0
def test_competing_node_cannot_acquire_cleanup_authority(
recovery_session_factory,
) -> None:
lease_observed_at = datetime.now(timezone.utc)
with recovery_session_factory() as session:
claim = acquire_lease(
session,
installation_id="campaign-artifact-tests",
resource_key="campaign:artifact-reconcile:tenant-1",
holder_node_id="node-other",
holder_incarnation="run-other",
ttl_seconds=900,
now=lease_observed_at,
)
assert claim is not None
session.commit()
storage = _MemoryStorage()
reference_patch, active_patch = _without_domain_references()
with reference_patch, active_patch, pytest.raises(RecoveryOperationBusy):
reconcile_campaign_artifacts(
recovery_session_factory,
storage=storage,
identity=_identity(),
tenant_id="tenant-1",
apply=True,
idempotency_key="competing-cleanup",
now=NOW,
)
def test_partial_storage_outage_remains_visible_and_retryable(
recovery_session_factory,
) -> None:
storage = _MemoryStorage()
removed = f"{PREFIX}build-a/message.eml"
retained = f"{PREFIX}build-b/message.eml"
storage.add(removed)
storage.add(retained)
storage.delete_failures.add(retained)
reference_patch, active_patch = _without_domain_references()
with reference_patch, active_patch:
partial = reconcile_campaign_artifacts(
recovery_session_factory,
storage=storage,
identity=_identity(),
tenant_id="tenant-1",
apply=True,
idempotency_key="partial-cleanup",
now=NOW,
)
storage.delete_failures.clear()
retry = reconcile_campaign_artifacts(
recovery_session_factory,
storage=storage,
identity=_identity(),
tenant_id="tenant-1",
apply=True,
idempotency_key="partial-cleanup-retry",
now=NOW,
)
assert partial["status"] == "recovery_required"
assert partial["deleted_count"] == 1
assert partial["failure_count"] == 1
assert removed not in storage.objects
assert retry["status"] == "applied"
assert retry["deleted_count"] == 1
assert storage.objects == {}
with recovery_session_factory() as session:
states = session.execute(
select(RecoveryOperation.status).order_by(RecoveryOperation.created_at)
).scalars().all()
assert "recovery_required" in states
assert states[-1] == "succeeded"
def test_unverifiable_delete_is_outcome_unknown(
recovery_session_factory,
) -> None:
storage = _MemoryStorage()
orphan = f"{PREFIX}build-unknown/message.eml"
storage.add(orphan)
storage.exists_failures.add(orphan)
reference_patch, active_patch = _without_domain_references()
with reference_patch, active_patch:
result = reconcile_campaign_artifacts(
recovery_session_factory,
storage=storage,
identity=_identity(),
tenant_id="tenant-1",
apply=True,
idempotency_key="unknown-cleanup",
now=NOW,
)
assert result["status"] == "outcome_unknown"
assert result["failure_count"] == 1