test_refund_worker_monitoring.py 9.8 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314
  1. from datetime import UTC, datetime, timedelta
  2. from fastapi.testclient import TestClient
  3. from zbt.commands.seed import build_seed_manifest
  4. from zbt.core.config import Settings
  5. from zbt.core.passwords import PasswordService
  6. from zbt.domains.catalog.repository import InMemoryCatalogRepository
  7. from zbt.domains.enrollment.models import OutboxEvent
  8. from zbt.domains.enrollment.repository import InMemoryEnrollmentRepository
  9. from zbt.domains.identity.repository import InMemoryIdentityRepository
  10. from zbt.domains.surrender.monitoring import (
  11. InMemoryWorkerRegistry,
  12. RefundOperationsMonitor,
  13. WorkerInstance,
  14. )
  15. from zbt.infrastructure.redis.worker_registry import RedisWorkerRegistry
  16. from zbt.main import create_app
  17. NOW = datetime(2026, 7, 29, 10, 0, tzinfo=UTC)
  18. def event(
  19. event_id: str,
  20. *,
  21. status: str,
  22. attempt_count: int = 0,
  23. last_error: str | None = None,
  24. ) -> OutboxEvent:
  25. return OutboxEvent(
  26. id=event_id,
  27. aggregate_type="SURRENDER_REQUEST",
  28. aggregate_id=f"request-{event_id}",
  29. event_type="SURRENDER_REFUND_REQUESTED",
  30. payload={},
  31. status=status,
  32. created_at=NOW - timedelta(minutes=attempt_count + 1),
  33. attempt_count=attempt_count,
  34. available_at=NOW + timedelta(seconds=30) if status == "RETRY" else None,
  35. last_error=last_error,
  36. )
  37. def test_refund_operations_snapshot_aggregates_workers_queue_and_failures() -> None:
  38. enrollments = InMemoryEnrollmentRepository()
  39. enrollments.save_outbox_event(event("pending", status="PENDING"))
  40. enrollments.save_outbox_event(
  41. event(
  42. "retry",
  43. status="RETRY",
  44. attempt_count=2,
  45. last_error="支付渠道连接超时",
  46. )
  47. )
  48. enrollments.save_outbox_event(
  49. event(
  50. "dead",
  51. status="DEAD",
  52. attempt_count=5,
  53. last_error="超过最大重试次数",
  54. )
  55. )
  56. enrollments.save_outbox_event(event("published", status="PUBLISHED", attempt_count=1))
  57. registry = InMemoryWorkerRegistry()
  58. registry.start(
  59. WorkerInstance(
  60. instance_id="worker-01",
  61. worker_name="refund",
  62. mode="parallel",
  63. host="node-a",
  64. pid=31001,
  65. status="RUNNING",
  66. started_at=NOW - timedelta(minutes=5),
  67. last_heartbeat_at=NOW,
  68. last_poll_at=NOW,
  69. last_processed_at=NOW - timedelta(seconds=10),
  70. total_processed=12,
  71. last_error=None,
  72. )
  73. )
  74. snapshot = RefundOperationsMonitor(
  75. enrollments,
  76. registry,
  77. clock=lambda: NOW,
  78. ).snapshot()
  79. assert snapshot["workers"]["online"] == 1
  80. assert snapshot["workers"]["items"][0]["instance_id"] == "worker-01"
  81. assert snapshot["queue"] == {
  82. "total": 4,
  83. "pending": 1,
  84. "processing": 0,
  85. "retrying": 1,
  86. "dead": 1,
  87. "cancelled": 0,
  88. "completed": 1,
  89. }
  90. assert snapshot["failures"][0]["event_id"] == "dead"
  91. assert snapshot["failures"][0]["request_id"] == "request-dead"
  92. assert snapshot["failures"][0]["last_error"] == "超过最大重试次数"
  93. assert snapshot["failures"][1]["event_id"] == "retry"
  94. def test_monitor_only_counts_refund_events_and_limits_failure_details() -> None:
  95. enrollments = InMemoryEnrollmentRepository()
  96. for index in range(25):
  97. enrollments.save_outbox_event(
  98. event(
  99. f"failed-{index:02d}",
  100. status="RETRY",
  101. attempt_count=index,
  102. last_error=f"错误-{index}",
  103. )
  104. )
  105. enrollments.save_outbox_event(
  106. OutboxEvent(
  107. id="unrelated",
  108. aggregate_type="ORDER",
  109. aggregate_id="order-1",
  110. event_type="POLICY_ISSUE_REQUESTED",
  111. payload={},
  112. status="PENDING",
  113. created_at=NOW,
  114. )
  115. )
  116. snapshot = RefundOperationsMonitor(
  117. enrollments,
  118. InMemoryWorkerRegistry(),
  119. clock=lambda: NOW,
  120. ).snapshot(failure_limit=20)
  121. assert snapshot["workers"]["online"] == 0
  122. assert snapshot["queue"]["total"] == 25
  123. assert len(snapshot["failures"]) == 20
  124. assert snapshot["failures"][0]["event_id"] == "failed-24"
  125. class FakeWorkerRegistryRedis:
  126. def __init__(self) -> None:
  127. self.hashes: dict[str, dict[str, str]] = {}
  128. self.expirations: dict[str, int] = {}
  129. def hset(self, name: str, key: str, value: str) -> int:
  130. values = self.hashes.setdefault(name, {})
  131. created = int(key not in values)
  132. values[key] = value
  133. return created
  134. def expire(self, name: str, seconds: int) -> bool:
  135. self.expirations[name] = seconds
  136. return True
  137. def scan_iter(self, match: str) -> list[str]:
  138. prefix = match.removesuffix("*")
  139. return [name for name in self.hashes if name.startswith(prefix)]
  140. def hgetall(self, name: str) -> dict[str, str]:
  141. return self.hashes.get(name, {})
  142. class Redis32WorkerRegistryRedis:
  143. """只实现 Redis 3.2 支持的单字段 HSET 形式。"""
  144. def __init__(self) -> None:
  145. self.hashes: dict[str, dict[str, str]] = {}
  146. self.expirations: dict[str, int] = {}
  147. def hset(self, name: str, key: str, value: str) -> int:
  148. values = self.hashes.setdefault(name, {})
  149. created = int(key not in values)
  150. values[key] = value
  151. return created
  152. def expire(self, name: str, seconds: int) -> bool:
  153. self.expirations[name] = seconds
  154. return True
  155. def scan_iter(self, match: str) -> list[str]:
  156. prefix = match.removesuffix("*")
  157. return [name for name in self.hashes if name.startswith(prefix)]
  158. def hgetall(self, name: str) -> dict[str, str]:
  159. return self.hashes.get(name, {})
  160. def test_redis_worker_registry_supports_redis_32_single_field_hset() -> None:
  161. redis = Redis32WorkerRegistryRedis()
  162. registry = RedisWorkerRegistry(
  163. redis,
  164. prefix="ins:s3:test:redis32:",
  165. ttl_seconds=30,
  166. )
  167. registry.start(
  168. WorkerInstance(
  169. instance_id="worker-redis32-01",
  170. worker_name="refund",
  171. mode="singleton",
  172. host="legacy-node",
  173. pid=32032,
  174. status="RUNNING",
  175. started_at=NOW,
  176. last_heartbeat_at=NOW,
  177. last_poll_at=None,
  178. last_processed_at=None,
  179. total_processed=0,
  180. last_error=None,
  181. )
  182. )
  183. instances = registry.list_instances("refund")
  184. assert instances[0].instance_id == "worker-redis32-01"
  185. assert instances[0].status == "RUNNING"
  186. assert (
  187. redis.expirations[
  188. "ins:s3:test:redis32:worker:instance:refund:worker-redis32-01"
  189. ]
  190. == 30
  191. )
  192. def test_redis_worker_registry_refreshes_ttl_and_accumulates_processed_count() -> None:
  193. redis = FakeWorkerRegistryRedis()
  194. registry = RedisWorkerRegistry(
  195. redis,
  196. prefix="ins:s3:test:monitor:",
  197. ttl_seconds=30,
  198. )
  199. registry.start(
  200. WorkerInstance(
  201. instance_id="worker-redis-01",
  202. worker_name="refund",
  203. mode="singleton",
  204. host="node-a",
  205. pid=32001,
  206. status="RUNNING",
  207. started_at=NOW,
  208. last_heartbeat_at=NOW,
  209. last_poll_at=None,
  210. last_processed_at=None,
  211. total_processed=0,
  212. last_error=None,
  213. )
  214. )
  215. registry.record_poll(
  216. "worker-redis-01",
  217. occurred_at=NOW + timedelta(seconds=2),
  218. processed=3,
  219. )
  220. registry.heartbeat("worker-redis-01", NOW + timedelta(seconds=10))
  221. instances = registry.list_instances("refund")
  222. assert instances[0].total_processed == 3
  223. assert instances[0].last_processed_at == NOW + timedelta(seconds=2)
  224. assert instances[0].last_heartbeat_at == NOW + timedelta(seconds=10)
  225. assert redis.expirations["ins:s3:test:monitor:worker:instance:refund:worker-redis-01"] == 30
  226. def test_admin_can_read_refund_worker_operations_but_cannot_run_worker_via_api() -> None:
  227. identities = InMemoryIdentityRepository()
  228. for admin in build_seed_manifest(PasswordService().hash("zaq1XSW@")).admin_users:
  229. identities.save_admin_user(admin)
  230. enrollments = InMemoryEnrollmentRepository()
  231. enrollments.save_outbox_event(event("pending-api", status="PENDING"))
  232. registry = InMemoryWorkerRegistry()
  233. registry.start(instance := WorkerInstance(
  234. instance_id="worker-api-01",
  235. worker_name="refund",
  236. mode="singleton",
  237. host="node-a",
  238. pid=34001,
  239. status="RUNNING",
  240. started_at=NOW,
  241. last_heartbeat_at=NOW,
  242. last_poll_at=NOW,
  243. last_processed_at=None,
  244. total_processed=0,
  245. last_error=None,
  246. ))
  247. app = create_app(
  248. settings=Settings(
  249. app_env="test",
  250. jwt_access_secret="a" * 32,
  251. jwt_refresh_secret="b" * 32,
  252. field_encryption_key="c" * 32,
  253. ),
  254. identity_repository=identities,
  255. catalog_repository=InMemoryCatalogRepository(),
  256. enrollment_repository=enrollments,
  257. worker_registry=registry,
  258. clock=lambda: NOW,
  259. )
  260. with TestClient(app) as client:
  261. login = client.post(
  262. "/api/v1/admin/auth/login",
  263. json={"username": "admin", "password": "zaq1XSW@"},
  264. )
  265. token = login.json()["data"]["tokens"]["access_token"]
  266. headers = {"Authorization": f"Bearer {token}"}
  267. snapshot = client.get("/api/v1/admin/refund-operations", headers=headers)
  268. removed_trigger = client.post(
  269. "/api/v1/dev/refund-worker/run-once",
  270. headers=headers,
  271. )
  272. assert instance.instance_id == "worker-api-01"
  273. assert snapshot.status_code == 200
  274. assert snapshot.json()["data"]["workers"]["online"] == 1
  275. assert snapshot.json()["data"]["queue"]["pending"] == 1
  276. assert removed_trigger.status_code == 404