| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177 |
- 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
|