| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314 |
- from datetime import UTC, datetime, timedelta
- from fastapi.testclient import TestClient
- from zbt.commands.seed import build_seed_manifest
- from zbt.core.config import Settings
- from zbt.core.passwords import PasswordService
- from zbt.domains.catalog.repository import InMemoryCatalogRepository
- from zbt.domains.enrollment.models import OutboxEvent
- from zbt.domains.enrollment.repository import InMemoryEnrollmentRepository
- from zbt.domains.identity.repository import InMemoryIdentityRepository
- from zbt.domains.surrender.monitoring import (
- InMemoryWorkerRegistry,
- RefundOperationsMonitor,
- WorkerInstance,
- )
- from zbt.infrastructure.redis.worker_registry import RedisWorkerRegistry
- from zbt.main import create_app
- NOW = datetime(2026, 7, 29, 10, 0, tzinfo=UTC)
- def event(
- event_id: str,
- *,
- status: str,
- attempt_count: int = 0,
- last_error: str | None = None,
- ) -> OutboxEvent:
- return OutboxEvent(
- id=event_id,
- aggregate_type="SURRENDER_REQUEST",
- aggregate_id=f"request-{event_id}",
- event_type="SURRENDER_REFUND_REQUESTED",
- payload={},
- status=status,
- created_at=NOW - timedelta(minutes=attempt_count + 1),
- attempt_count=attempt_count,
- available_at=NOW + timedelta(seconds=30) if status == "RETRY" else None,
- last_error=last_error,
- )
- def test_refund_operations_snapshot_aggregates_workers_queue_and_failures() -> None:
- enrollments = InMemoryEnrollmentRepository()
- enrollments.save_outbox_event(event("pending", status="PENDING"))
- enrollments.save_outbox_event(
- event(
- "retry",
- status="RETRY",
- attempt_count=2,
- last_error="支付渠道连接超时",
- )
- )
- enrollments.save_outbox_event(
- event(
- "dead",
- status="DEAD",
- attempt_count=5,
- last_error="超过最大重试次数",
- )
- )
- enrollments.save_outbox_event(event("published", status="PUBLISHED", attempt_count=1))
- registry = InMemoryWorkerRegistry()
- registry.start(
- WorkerInstance(
- instance_id="worker-01",
- worker_name="refund",
- mode="parallel",
- host="node-a",
- pid=31001,
- status="RUNNING",
- started_at=NOW - timedelta(minutes=5),
- last_heartbeat_at=NOW,
- last_poll_at=NOW,
- last_processed_at=NOW - timedelta(seconds=10),
- total_processed=12,
- last_error=None,
- )
- )
- snapshot = RefundOperationsMonitor(
- enrollments,
- registry,
- clock=lambda: NOW,
- ).snapshot()
- assert snapshot["workers"]["online"] == 1
- assert snapshot["workers"]["items"][0]["instance_id"] == "worker-01"
- assert snapshot["queue"] == {
- "total": 4,
- "pending": 1,
- "processing": 0,
- "retrying": 1,
- "dead": 1,
- "cancelled": 0,
- "completed": 1,
- }
- assert snapshot["failures"][0]["event_id"] == "dead"
- assert snapshot["failures"][0]["request_id"] == "request-dead"
- assert snapshot["failures"][0]["last_error"] == "超过最大重试次数"
- assert snapshot["failures"][1]["event_id"] == "retry"
- def test_monitor_only_counts_refund_events_and_limits_failure_details() -> None:
- enrollments = InMemoryEnrollmentRepository()
- for index in range(25):
- enrollments.save_outbox_event(
- event(
- f"failed-{index:02d}",
- status="RETRY",
- attempt_count=index,
- last_error=f"错误-{index}",
- )
- )
- enrollments.save_outbox_event(
- OutboxEvent(
- id="unrelated",
- aggregate_type="ORDER",
- aggregate_id="order-1",
- event_type="POLICY_ISSUE_REQUESTED",
- payload={},
- status="PENDING",
- created_at=NOW,
- )
- )
- snapshot = RefundOperationsMonitor(
- enrollments,
- InMemoryWorkerRegistry(),
- clock=lambda: NOW,
- ).snapshot(failure_limit=20)
- assert snapshot["workers"]["online"] == 0
- assert snapshot["queue"]["total"] == 25
- assert len(snapshot["failures"]) == 20
- assert snapshot["failures"][0]["event_id"] == "failed-24"
- class FakeWorkerRegistryRedis:
- def __init__(self) -> None:
- self.hashes: dict[str, dict[str, str]] = {}
- self.expirations: dict[str, int] = {}
- def hset(self, name: str, key: str, value: str) -> int:
- values = self.hashes.setdefault(name, {})
- created = int(key not in values)
- values[key] = value
- return created
- def expire(self, name: str, seconds: int) -> bool:
- self.expirations[name] = seconds
- return True
- def scan_iter(self, match: str) -> list[str]:
- prefix = match.removesuffix("*")
- return [name for name in self.hashes if name.startswith(prefix)]
- def hgetall(self, name: str) -> dict[str, str]:
- return self.hashes.get(name, {})
- class Redis32WorkerRegistryRedis:
- """只实现 Redis 3.2 支持的单字段 HSET 形式。"""
- def __init__(self) -> None:
- self.hashes: dict[str, dict[str, str]] = {}
- self.expirations: dict[str, int] = {}
- def hset(self, name: str, key: str, value: str) -> int:
- values = self.hashes.setdefault(name, {})
- created = int(key not in values)
- values[key] = value
- return created
- def expire(self, name: str, seconds: int) -> bool:
- self.expirations[name] = seconds
- return True
- def scan_iter(self, match: str) -> list[str]:
- prefix = match.removesuffix("*")
- return [name for name in self.hashes if name.startswith(prefix)]
- def hgetall(self, name: str) -> dict[str, str]:
- return self.hashes.get(name, {})
- def test_redis_worker_registry_supports_redis_32_single_field_hset() -> None:
- redis = Redis32WorkerRegistryRedis()
- registry = RedisWorkerRegistry(
- redis,
- prefix="ins:s3:test:redis32:",
- ttl_seconds=30,
- )
- registry.start(
- WorkerInstance(
- instance_id="worker-redis32-01",
- worker_name="refund",
- mode="singleton",
- host="legacy-node",
- pid=32032,
- status="RUNNING",
- started_at=NOW,
- last_heartbeat_at=NOW,
- last_poll_at=None,
- last_processed_at=None,
- total_processed=0,
- last_error=None,
- )
- )
- instances = registry.list_instances("refund")
- assert instances[0].instance_id == "worker-redis32-01"
- assert instances[0].status == "RUNNING"
- assert (
- redis.expirations[
- "ins:s3:test:redis32:worker:instance:refund:worker-redis32-01"
- ]
- == 30
- )
- def test_redis_worker_registry_refreshes_ttl_and_accumulates_processed_count() -> None:
- redis = FakeWorkerRegistryRedis()
- registry = RedisWorkerRegistry(
- redis,
- prefix="ins:s3:test:monitor:",
- ttl_seconds=30,
- )
- registry.start(
- WorkerInstance(
- instance_id="worker-redis-01",
- worker_name="refund",
- mode="singleton",
- host="node-a",
- pid=32001,
- status="RUNNING",
- started_at=NOW,
- last_heartbeat_at=NOW,
- last_poll_at=None,
- last_processed_at=None,
- total_processed=0,
- last_error=None,
- )
- )
- registry.record_poll(
- "worker-redis-01",
- occurred_at=NOW + timedelta(seconds=2),
- processed=3,
- )
- registry.heartbeat("worker-redis-01", NOW + timedelta(seconds=10))
- instances = registry.list_instances("refund")
- assert instances[0].total_processed == 3
- assert instances[0].last_processed_at == NOW + timedelta(seconds=2)
- assert instances[0].last_heartbeat_at == NOW + timedelta(seconds=10)
- assert redis.expirations["ins:s3:test:monitor:worker:instance:refund:worker-redis-01"] == 30
- def test_admin_can_read_refund_worker_operations_but_cannot_run_worker_via_api() -> None:
- identities = InMemoryIdentityRepository()
- for admin in build_seed_manifest(PasswordService().hash("zaq1XSW@")).admin_users:
- identities.save_admin_user(admin)
- enrollments = InMemoryEnrollmentRepository()
- enrollments.save_outbox_event(event("pending-api", status="PENDING"))
- registry = InMemoryWorkerRegistry()
- registry.start(instance := WorkerInstance(
- instance_id="worker-api-01",
- worker_name="refund",
- mode="singleton",
- host="node-a",
- pid=34001,
- status="RUNNING",
- started_at=NOW,
- last_heartbeat_at=NOW,
- last_poll_at=NOW,
- last_processed_at=None,
- total_processed=0,
- last_error=None,
- ))
- app = create_app(
- settings=Settings(
- app_env="test",
- jwt_access_secret="a" * 32,
- jwt_refresh_secret="b" * 32,
- field_encryption_key="c" * 32,
- ),
- identity_repository=identities,
- catalog_repository=InMemoryCatalogRepository(),
- enrollment_repository=enrollments,
- worker_registry=registry,
- clock=lambda: NOW,
- )
- with TestClient(app) as client:
- login = client.post(
- "/api/v1/admin/auth/login",
- json={"username": "admin", "password": "zaq1XSW@"},
- )
- token = login.json()["data"]["tokens"]["access_token"]
- headers = {"Authorization": f"Bearer {token}"}
- snapshot = client.get("/api/v1/admin/refund-operations", headers=headers)
- removed_trigger = client.post(
- "/api/v1/dev/refund-worker/run-once",
- headers=headers,
- )
- assert instance.instance_id == "worker-api-01"
- assert snapshot.status_code == 200
- assert snapshot.json()["data"]["workers"]["online"] == 1
- assert snapshot.json()["data"]["queue"]["pending"] == 1
- assert removed_trigger.status_code == 404
|