from __future__ import annotations from copy import deepcopy from datetime import UTC, datetime import pytest OWNER = "01900000-0000-7000-8000-000000069001" OPERATOR = "01900000-0000-7000-8000-000000069002" ASSET = "01900000-0000-7000-8000-000000069101" GRANT = "01900000-0000-7000-8000-000000069201" class MemorySubscriptionRepository: def __init__(self): self.users = {OWNER, OPERATOR} self.subscriptions = {} self.deliveries = {} self.events = [] self.incident_links = [] self.grant = {"uid": GRANT, "asset_uid": ASSET, "purpose": "quality_review", "status": "active", "expires_at": "2026-08-13T08:00:00+00:00"} def users_available(self, values): return set(values) & self.users def create_subscription(self, record): self.subscriptions[record["uid"]] = deepcopy(record) return deepcopy(record) def get_grant_for_subscription(self, uid): return deepcopy(self.grant) if uid == self.grant["uid"] else None def transition_subscription(self, uid, expected_status, new_status): value = self.subscriptions[uid] if value["status"] not in expected_status: raise RuntimeError("subscription transition conflict") value["status"] = new_status value["current_version"] += 1 return deepcopy(value) def enqueue_subscription_delivery(self, record): existing = next((item for item in self.deliveries.values() if item["idempotency_key"] == record["idempotency_key"]), None) if existing: if existing["request_digest"] != record["request_digest"]: raise RuntimeError("subscription idempotency conflict") return deepcopy(existing) self.deliveries[record["uid"]] = deepcopy(record) return deepcopy(record) def claim_subscription_delivery(self, uid, worker): value = self.deliveries[uid] if value["status"] not in {"pending", "processing"}: return None value.update(status="processing", lease_owner=worker, lease_fence=value["lease_fence"] + 1) return deepcopy(value) def record_subscription_attempt(self, uid, worker, fence, delivered): value = self.deliveries[uid] if value["lease_owner"] != worker or value["lease_fence"] != fence or value["status"] != "processing": return None value["lease_owner"] = None value["attempt_count"] += 1 if delivered: value["status"] = "delivered" elif value["attempt_count"] >= 3: value["status"] = "dead_letter" else: value["status"] = "pending" return deepcopy(value) def compensate_subscription_delivery(self, uid, reason_code, receipt_code): value = self.deliveries[uid] if value["status"] == "compensated": if (value["compensation_reason_code"], value["compensation_receipt_code"]) == (reason_code, receipt_code): return deepcopy(value) raise RuntimeError("subscription compensation conflict") if value["status"] != "dead_letter": raise RuntimeError("subscription delivery is not compensable") value.update(status="compensated", compensation_reason_code=reason_code, compensation_receipt_code=receipt_code) return deepcopy(value) def add_subscription_evidence(self, record): self.events.append(deepcopy(record)) def link_anomaly_incident(self, incident_uid, asset_uid, label): self.incident_links.append((incident_uid, asset_uid, label)) def subscriptions_due_for_reclaim(self, instant): return [deepcopy(value) for value in self.subscriptions.values() if value["status"] == "active"] def mark_subscription_reclaimed(self, uid, expected_version): value = self.subscriptions[uid] if value["current_version"] != expected_version: raise RuntimeError("subscription reclaim conflict") value.update(status="reclaimed", current_version=expected_version + 1) return deepcopy(value) @pytest.fixture() def subscriptions(): from app.core.system.trusted_delivery_subscriptions import ( TrustedDeliverySubscriptionService, ) now = datetime(2026, 8, 11, 8, 0, tzinfo=UTC) ids = iter(f"01900000-0000-7000-8000-{value:012d}" for value in range(6900, 7000)) return TrustedDeliverySubscriptionService( MemorySubscriptionRepository(), uid_factory=lambda: next(ids), now_factory=lambda: now ) def _payload(): return { "grant_uid": GRANT, "asset_uid": ASSET, "trigger": {"kind": "schedule", "schedule_ref": "hourly-quality-v1"}, "purpose": "quality_review", "expires_at": "2026-08-12T08:00:00+00:00", "idempotency_key": "wp06-subscription-1", } def test_subscription_state_machine_is_closed_and_default_denies_unknown_fields(subscriptions): subscription = subscriptions.create_subscription(_payload(), actor_uid=OWNER) assert subscription["status"] == "draft" assert subscriptions.activate_subscription(subscription["uid"], actor_uid=OPERATOR)["status"] == "active" assert subscriptions.pause_subscription(subscription["uid"], actor_uid=OPERATOR)["status"] == "paused" assert subscriptions.resume_subscription(subscription["uid"], actor_uid=OPERATOR)["status"] == "active" assert subscriptions.terminate_subscription(subscription["uid"], actor_uid=OPERATOR)["status"] == "terminated" with pytest.raises(ValueError, match="unsupported fields"): subscriptions.create_subscription({**_payload(), "url": "https://target.invalid"}, actor_uid=OWNER) def test_subscription_creation_binds_active_unexpired_grant_asset_and_purpose(subscriptions): with pytest.raises(PermissionError, match="grant"): subscriptions.create_subscription({**_payload(), "asset_uid": "01900000-0000-7000-8000-000000069102"}, actor_uid=OWNER) def test_delivery_is_fenced_retries_to_dead_letter_and_compensation_is_replay_safe(subscriptions): subscription = subscriptions.create_subscription(_payload(), actor_uid=OWNER) subscriptions.activate_subscription(subscription["uid"], actor_uid=OPERATOR) delivery = subscriptions.enqueue_delivery({"subscription_uid": subscription["uid"], "trigger_ref": "tick-1", "event_digest": "a" * 64}, actor_uid=OPERATOR) for _ in range(3): claimed = subscriptions.claim_delivery(delivery["uid"], worker_id="worker-a") assert subscriptions.record_delivery_attempt(delivery["uid"], worker_id="worker-b", lease_fence=claimed["lease_fence"], delivered=False, actor_uid=OPERATOR) is None result = subscriptions.record_delivery_attempt(delivery["uid"], worker_id="worker-a", lease_fence=claimed["lease_fence"], delivered=False, actor_uid=OPERATOR) assert result["status"] == "dead_letter" compensated = subscriptions.compensate_dead_letter(delivery["uid"], reason_code="operator_review", receipt_code="COMPENSATED", actor_uid=OPERATOR) assert compensated["status"] == "compensated" assert subscriptions.compensate_dead_letter(delivery["uid"], reason_code="operator_review", receipt_code="COMPENSATED", actor_uid=OPERATOR) == compensated def test_event_trigger_and_anomaly_reference_are_safe_and_incident_is_a_reference(subscriptions): subscription = subscriptions.create_subscription({ **_payload(), "idempotency_key": "wp06-subscription-event-1", "trigger": {"kind": "event", "event_type": "asset_changed"}, }, actor_uid=OWNER) assert subscription["trigger"]["kind"] == "event" anomaly = subscriptions.report_anomaly({ "kind": "purpose_breach", "asset_uid": ASSET, "subscription_uid": subscription["uid"], "incident_uid": "01900000-0000-7000-8000-000000069301", "summary": "purpose mismatch", "correlation_id": "01900000-0000-7000-8000-000000069302", }, actor_uid=OPERATOR) assert anomaly["incident_uid"].endswith("69301") assert subscriptions.repository.incident_links == [(anomaly["incident_uid"], ASSET, "trusted_delivery:purpose_breach")] with pytest.raises(ValueError, match="unsafe"): subscriptions.report_anomaly({**anomaly, "summary": "https://leak.invalid"}, actor_uid=OPERATOR) def test_expiry_marks_subscription_reclaimed_only_after_reuse_of_reclaim_callback(subscriptions): subscription = subscriptions.create_subscription(_payload(), actor_uid=OWNER) subscriptions.activate_subscription(subscription["uid"], actor_uid=OPERATOR) calls = [] assert subscriptions.expire_and_reclaim( as_of="2026-08-13T00:00:00+00:00", actor_uid=OPERATOR, reclaim=lambda grant_uid, key: calls.append((grant_uid, key)) or {"status": "reclaimed"}, )[0]["status"] == "reclaimed" assert calls[0][0] == GRANT