test_trusted_delivery_subscriptions.py 8.6 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177
  1. from __future__ import annotations
  2. from copy import deepcopy
  3. from datetime import UTC, datetime
  4. import pytest
  5. OWNER = "01900000-0000-7000-8000-000000069001"
  6. OPERATOR = "01900000-0000-7000-8000-000000069002"
  7. ASSET = "01900000-0000-7000-8000-000000069101"
  8. GRANT = "01900000-0000-7000-8000-000000069201"
  9. class MemorySubscriptionRepository:
  10. def __init__(self):
  11. self.users = {OWNER, OPERATOR}
  12. self.subscriptions = {}
  13. self.deliveries = {}
  14. self.events = []
  15. self.incident_links = []
  16. self.grant = {"uid": GRANT, "asset_uid": ASSET, "purpose": "quality_review", "status": "active", "expires_at": "2026-08-13T08:00:00+00:00"}
  17. def users_available(self, values):
  18. return set(values) & self.users
  19. def create_subscription(self, record):
  20. self.subscriptions[record["uid"]] = deepcopy(record)
  21. return deepcopy(record)
  22. def get_grant_for_subscription(self, uid):
  23. return deepcopy(self.grant) if uid == self.grant["uid"] else None
  24. def transition_subscription(self, uid, expected_status, new_status):
  25. value = self.subscriptions[uid]
  26. if value["status"] not in expected_status:
  27. raise RuntimeError("subscription transition conflict")
  28. value["status"] = new_status
  29. value["current_version"] += 1
  30. return deepcopy(value)
  31. def enqueue_subscription_delivery(self, record):
  32. existing = next((item for item in self.deliveries.values() if item["idempotency_key"] == record["idempotency_key"]), None)
  33. if existing:
  34. if existing["request_digest"] != record["request_digest"]:
  35. raise RuntimeError("subscription idempotency conflict")
  36. return deepcopy(existing)
  37. self.deliveries[record["uid"]] = deepcopy(record)
  38. return deepcopy(record)
  39. def claim_subscription_delivery(self, uid, worker):
  40. value = self.deliveries[uid]
  41. if value["status"] not in {"pending", "processing"}:
  42. return None
  43. value.update(status="processing", lease_owner=worker, lease_fence=value["lease_fence"] + 1)
  44. return deepcopy(value)
  45. def record_subscription_attempt(self, uid, worker, fence, delivered):
  46. value = self.deliveries[uid]
  47. if value["lease_owner"] != worker or value["lease_fence"] != fence or value["status"] != "processing":
  48. return None
  49. value["lease_owner"] = None
  50. value["attempt_count"] += 1
  51. if delivered:
  52. value["status"] = "delivered"
  53. elif value["attempt_count"] >= 3:
  54. value["status"] = "dead_letter"
  55. else:
  56. value["status"] = "pending"
  57. return deepcopy(value)
  58. def compensate_subscription_delivery(self, uid, reason_code, receipt_code):
  59. value = self.deliveries[uid]
  60. if value["status"] == "compensated":
  61. if (value["compensation_reason_code"], value["compensation_receipt_code"]) == (reason_code, receipt_code):
  62. return deepcopy(value)
  63. raise RuntimeError("subscription compensation conflict")
  64. if value["status"] != "dead_letter":
  65. raise RuntimeError("subscription delivery is not compensable")
  66. value.update(status="compensated", compensation_reason_code=reason_code, compensation_receipt_code=receipt_code)
  67. return deepcopy(value)
  68. def add_subscription_evidence(self, record):
  69. self.events.append(deepcopy(record))
  70. def link_anomaly_incident(self, incident_uid, asset_uid, label):
  71. self.incident_links.append((incident_uid, asset_uid, label))
  72. def subscriptions_due_for_reclaim(self, instant):
  73. return [deepcopy(value) for value in self.subscriptions.values() if value["status"] == "active"]
  74. def mark_subscription_reclaimed(self, uid, expected_version):
  75. value = self.subscriptions[uid]
  76. if value["current_version"] != expected_version:
  77. raise RuntimeError("subscription reclaim conflict")
  78. value.update(status="reclaimed", current_version=expected_version + 1)
  79. return deepcopy(value)
  80. @pytest.fixture()
  81. def subscriptions():
  82. from app.core.system.trusted_delivery_subscriptions import (
  83. TrustedDeliverySubscriptionService,
  84. )
  85. now = datetime(2026, 8, 11, 8, 0, tzinfo=UTC)
  86. ids = iter(f"01900000-0000-7000-8000-{value:012d}" for value in range(6900, 7000))
  87. return TrustedDeliverySubscriptionService(
  88. MemorySubscriptionRepository(), uid_factory=lambda: next(ids), now_factory=lambda: now
  89. )
  90. def _payload():
  91. return {
  92. "grant_uid": GRANT,
  93. "asset_uid": ASSET,
  94. "trigger": {"kind": "schedule", "schedule_ref": "hourly-quality-v1"},
  95. "purpose": "quality_review",
  96. "expires_at": "2026-08-12T08:00:00+00:00",
  97. "idempotency_key": "wp06-subscription-1",
  98. }
  99. def test_subscription_state_machine_is_closed_and_default_denies_unknown_fields(subscriptions):
  100. subscription = subscriptions.create_subscription(_payload(), actor_uid=OWNER)
  101. assert subscription["status"] == "draft"
  102. assert subscriptions.activate_subscription(subscription["uid"], actor_uid=OPERATOR)["status"] == "active"
  103. assert subscriptions.pause_subscription(subscription["uid"], actor_uid=OPERATOR)["status"] == "paused"
  104. assert subscriptions.resume_subscription(subscription["uid"], actor_uid=OPERATOR)["status"] == "active"
  105. assert subscriptions.terminate_subscription(subscription["uid"], actor_uid=OPERATOR)["status"] == "terminated"
  106. with pytest.raises(ValueError, match="unsupported fields"):
  107. subscriptions.create_subscription({**_payload(), "url": "https://target.invalid"}, actor_uid=OWNER)
  108. def test_subscription_creation_binds_active_unexpired_grant_asset_and_purpose(subscriptions):
  109. with pytest.raises(PermissionError, match="grant"):
  110. subscriptions.create_subscription({**_payload(), "asset_uid": "01900000-0000-7000-8000-000000069102"}, actor_uid=OWNER)
  111. def test_delivery_is_fenced_retries_to_dead_letter_and_compensation_is_replay_safe(subscriptions):
  112. subscription = subscriptions.create_subscription(_payload(), actor_uid=OWNER)
  113. subscriptions.activate_subscription(subscription["uid"], actor_uid=OPERATOR)
  114. delivery = subscriptions.enqueue_delivery({"subscription_uid": subscription["uid"], "trigger_ref": "tick-1", "event_digest": "a" * 64}, actor_uid=OPERATOR)
  115. for _ in range(3):
  116. claimed = subscriptions.claim_delivery(delivery["uid"], worker_id="worker-a")
  117. assert subscriptions.record_delivery_attempt(delivery["uid"], worker_id="worker-b", lease_fence=claimed["lease_fence"], delivered=False, actor_uid=OPERATOR) is None
  118. result = subscriptions.record_delivery_attempt(delivery["uid"], worker_id="worker-a", lease_fence=claimed["lease_fence"], delivered=False, actor_uid=OPERATOR)
  119. assert result["status"] == "dead_letter"
  120. compensated = subscriptions.compensate_dead_letter(delivery["uid"], reason_code="operator_review", receipt_code="COMPENSATED", actor_uid=OPERATOR)
  121. assert compensated["status"] == "compensated"
  122. assert subscriptions.compensate_dead_letter(delivery["uid"], reason_code="operator_review", receipt_code="COMPENSATED", actor_uid=OPERATOR) == compensated
  123. def test_event_trigger_and_anomaly_reference_are_safe_and_incident_is_a_reference(subscriptions):
  124. subscription = subscriptions.create_subscription({
  125. **_payload(), "idempotency_key": "wp06-subscription-event-1",
  126. "trigger": {"kind": "event", "event_type": "asset_changed"},
  127. }, actor_uid=OWNER)
  128. assert subscription["trigger"]["kind"] == "event"
  129. anomaly = subscriptions.report_anomaly({
  130. "kind": "purpose_breach", "asset_uid": ASSET, "subscription_uid": subscription["uid"],
  131. "incident_uid": "01900000-0000-7000-8000-000000069301", "summary": "purpose mismatch", "correlation_id": "01900000-0000-7000-8000-000000069302",
  132. }, actor_uid=OPERATOR)
  133. assert anomaly["incident_uid"].endswith("69301")
  134. assert subscriptions.repository.incident_links == [(anomaly["incident_uid"], ASSET, "trusted_delivery:purpose_breach")]
  135. with pytest.raises(ValueError, match="unsafe"):
  136. subscriptions.report_anomaly({**anomaly, "summary": "https://leak.invalid"}, actor_uid=OPERATOR)
  137. def test_expiry_marks_subscription_reclaimed_only_after_reuse_of_reclaim_callback(subscriptions):
  138. subscription = subscriptions.create_subscription(_payload(), actor_uid=OWNER)
  139. subscriptions.activate_subscription(subscription["uid"], actor_uid=OPERATOR)
  140. calls = []
  141. assert subscriptions.expire_and_reclaim(
  142. as_of="2026-08-13T00:00:00+00:00", actor_uid=OPERATOR,
  143. reclaim=lambda grant_uid, key: calls.append((grant_uid, key)) or {"status": "reclaimed"},
  144. )[0]["status"] == "reclaimed"
  145. assert calls[0][0] == GRANT