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