Files
root 0c511f3451 1、增加调用 AI 视频生成能力和虚拟素材库管理的对外api
2、增加后台apikkey管理
3、增加apikey单独的模型定价
4、增加apikey调用情况
5、完善所有数据的注释增加
2026-08-06 13:13:28 +08:00

540 lines
23 KiB
Python

import logging
from celery import Celery
from celery.signals import (
celeryd_init,
heartbeat_sent,
worker_init,
worker_process_init,
worker_process_shutdown,
worker_ready,
worker_shutdown,
)
from app.config import settings
from app.enums.celery_queue import CeleryQueue, CeleryTaskName
from app.models.base import engine
from app.tasks.async_runner import close_loop, run_async
logger = logging.getLogger("video_gen")
# 显式注册所有 Celery 任务模块,避免新增任务文件后 worker 启动时未注册任务。
# 不再依赖 app.tasks.__init__ 内部 import,也不再依赖 autodiscover_tasks。
CELERY_TASK_IMPORTS = (
"app.tasks.generation_create_tasks",
"app.tasks.generation_poll_tasks",
"app.tasks.generation_download_tasks",
"app.tasks.generation_recovery_tasks",
"app.tasks.video_upscale_tasks",
"app.tasks.hot_opening_replicate_tasks",
"app.tasks.shot_replicate_tasks",
"app.tasks.shot_replicate_flow_tasks",
"app.tasks.module_async_recovery_tasks",
"app.tasks.module_generation_v2_tasks",
"app.tasks.private_portrait_asset_tasks",
"app.tasks.vp_v3_asset_tasks",
"app.tasks.celery_runtime_tasks",
"app.tasks.api_generation_tasks",
"app.tasks.api_recovery_tasks",
"app.tasks.api_upscale_tasks",
)
RECOVERY_QUEUE = settings.CELERY_RECOVERY_QUEUE or CeleryQueue.GEN_RECOVERY.value
def _derive_redis_db(url: str, db_no: int) -> str:
if not url:
return url
import re
if re.search(r"/\d+$", url):
return re.sub(r"/\d+$", f"/{db_no}", url)
return url.rstrip("/") + f"/{db_no}"
def _beat_schedule() -> dict:
schedule: dict = {}
if bool(getattr(settings, "POLL_DUE_DISPATCH_ENABLED", True)):
schedule["dispatch-due-poll-tasks-every-minute"] = {
"task": CeleryTaskName.DISPATCH_DUE_POLL.value,
"schedule": max(1, int(settings.POLL_DUE_DISPATCH_INTERVAL_SECONDS or 60)),
"options": {
"queue": RECOVERY_QUEUE,
"priority": settings.DOWNLOAD_TASK_PRIORITY_RECOVER,
},
}
schedule["generation-create-recovery"] = {
"task": CeleryTaskName.RECOVER_CREATE.value,
"schedule": max(1, int(settings.GENERATION_CREATE_RECOVERY_INTERVAL_SECONDS or 60)),
"options": {
"queue": RECOVERY_QUEUE,
"priority": settings.DOWNLOAD_TASK_PRIORITY_RECOVER,
},
}
schedule["module-async-recovery"] = {
"task": CeleryTaskName.MODULE_ASYNC_RECOVERY.value,
"schedule": max(1, int(settings.MODULE_ASYNC_RECOVERY_INTERVAL_SECONDS or 60)),
"options": {
"queue": RECOVERY_QUEUE,
"priority": settings.DOWNLOAD_TASK_PRIORITY_RECOVER,
},
}
schedule["video-upscale-recovery-every-minute"] = {
"task": CeleryTaskName.VIDEO_UPSCALE_RECOVER.value,
"schedule": 60,
"options": {
"queue": RECOVERY_QUEUE,
"priority": settings.DOWNLOAD_TASK_PRIORITY_RECOVER,
},
}
schedule["api-generation-recovery-every-minute"] = {
"task": "api_generation.recover_tasks_once",
"schedule": 60,
"options": {
"queue": RECOVERY_QUEUE,
"priority": settings.DOWNLOAD_TASK_PRIORITY_RECOVER,
},
}
schedule["generation-download-recovery"] = {
"task": CeleryTaskName.RECOVER_DOWNLOAD.value,
"schedule": max(1, int(settings.DOWNLOAD_RECOVERY_INTERVAL_SECONDS or 60)),
"options": {"queue": RECOVERY_QUEUE, "priority": settings.DOWNLOAD_TASK_PRIORITY_RECOVER},
}
schedule["shot-split-recovery"] = {
"task": CeleryTaskName.SHOT_SPLIT_RECOVERY.value,
"schedule": 60,
"options": {"queue": RECOVERY_QUEUE, "priority": settings.DOWNLOAD_TASK_PRIORITY_RECOVER},
}
schedule["shot-analysis-recovery"] = {
"task": CeleryTaskName.SHOT_ANALYSIS_RECOVERY.value,
"schedule": 60,
"options": {"queue": RECOVERY_QUEUE, "priority": settings.DOWNLOAD_TASK_PRIORITY_RECOVER},
}
schedule["celery-runtime-reconcile"] = {
"task": CeleryTaskName.CELERY_RUNTIME_RECONCILE.value,
"schedule": max(60, int(settings.CELERY_RUNTIME_RECONCILE_INTERVAL_SECONDS or 300)),
"options": {"queue": RECOVERY_QUEUE, "priority": settings.DOWNLOAD_TASK_PRIORITY_RECOVER},
}
schedule["celery-runtime-registry-gc"] = {
"task": CeleryTaskName.CELERY_RUNTIME_GC.value,
"schedule": max(60, int(settings.CELERY_RUNTIME_GC_INTERVAL_SECONDS or 600)),
"options": {"queue": RECOVERY_QUEUE, "priority": settings.DOWNLOAD_TASK_PRIORITY_RECOVER},
}
schedule["private-portrait-sync-due-assets-every-minute"] = {
"task": CeleryTaskName.PRIVATE_PORTRAIT_SYNC_DUE_ASSETS.value,
"schedule": 60,
"options": {"queue": CeleryQueue.GEN_PRIVATE_PORTRAIT.value},
}
schedule["private-portrait-recover-remote-deletes-every-5-minutes"] = {
"task": CeleryTaskName.PRIVATE_PORTRAIT_RECOVER_REMOTE_DELETES.value,
"schedule": 300,
"options": {"queue": CeleryQueue.GEN_PRIVATE_PORTRAIT.value},
}
schedule["vp-v3-sync-due-assets-every-minute"] = {
"task": CeleryTaskName.VP_V3_SYNC_DUE_ASSETS.value,
"schedule": 60,
"options": {"queue": CeleryQueue.GEN_PRIVATE_PORTRAIT.value},
}
schedule["vp-v3-recover-remote-deletes-every-5-minutes"] = {
"task": CeleryTaskName.VP_V3_RECOVER_REMOTE_DELETES.value,
"schedule": 300,
"options": {"queue": CeleryQueue.GEN_PRIVATE_PORTRAIT.value},
}
return schedule
broker_url = settings.CELERY_BROKER_URL or (_derive_redis_db(settings.REDIS_URL, 1) if settings.REDIS_URL else "")
backend_url = settings.CELERY_RESULT_BACKEND or (_derive_redis_db(settings.REDIS_URL, 2) if settings.REDIS_URL else "")
if broker_url:
celery_app = Celery("videogen", include=CELERY_TASK_IMPORTS)
celery_app.conf.update(
broker_url=broker_url,
result_backend=backend_url or broker_url,
imports=CELERY_TASK_IMPORTS,
task_serializer="json",
accept_content=["json"],
result_serializer="json",
result_expires=max(60, int(settings.CELERY_RESULT_EXPIRES_SECONDS or 7200)),
task_store_errors_even_if_ignored=True,
timezone="Asia/Shanghai",
enable_utc=True,
task_soft_time_limit=600,
task_time_limit=900,
task_acks_late=True,
task_reject_on_worker_lost=True,
task_track_started=True,
beat_schedule=_beat_schedule(),
task_annotations={
# 生成链路任务以数据库状态为准,不依赖 Celery result backend。
# 这里忽略结果可避免任务误返回 ORM / 非 JSON 对象时触发结果序列化失败。
# "generation.chatapi_create_generation_task": {"ignore_result": True},
# "generation.poll_generation_task": {"ignore_result": True},
# "generation.download_generation_result_task": {"ignore_result": True},
"hot_opening.start_image_prompt_optimize": {"ignore_result": True},
"hot_opening.start_video_prompt_optimize": {"ignore_result": True},
"shot_replicate.analyze_original_video": {
"ignore_result": True,
"soft_time_limit": int(settings.SHOT_ANALYSIS_SOFT_TIME_LIMIT_SECONDS or 3720),
"time_limit": int(settings.SHOT_ANALYSIS_TIME_LIMIT_SECONDS or 3900),
},
"shot_replicate.analyze_custom_segment_video": {
"ignore_result": True,
"soft_time_limit": int(settings.SHOT_ANALYSIS_SOFT_TIME_LIMIT_SECONDS or 3720),
"time_limit": int(settings.SHOT_ANALYSIS_TIME_LIMIT_SECONDS or 3900),
},
"shot_replicate.split_one_segment": {"ignore_result": True},
"shot_replicate.start_image_prompt_optimize": {"ignore_result": True},
"shot_replicate.start_video_prompt_optimize": {"ignore_result": True},
"module_generation_v2.start_video_prompt_optimize": {"ignore_result": True},
CeleryTaskName.VIDEO_UPSCALE_EXECUTE_LOCAL.value: {
"ignore_result": True,
"soft_time_limit": max(60, int(settings.VIDEO_UPSCALE_LOCAL_TIMEOUT_SECONDS or 3600)) + 60,
"time_limit": max(60, int(settings.VIDEO_UPSCALE_LOCAL_TIMEOUT_SECONDS or 3600)) + 300,
},
CeleryTaskName.VIDEO_UPSCALE_SUBMIT_REMOTE.value: {"ignore_result": True},
CeleryTaskName.VIDEO_UPSCALE_POLL_REMOTE.value: {"ignore_result": True},
CeleryTaskName.VIDEO_UPSCALE_DOWNLOAD_REMOTE_RESULT.value: {"ignore_result": True},
CeleryTaskName.VIDEO_UPSCALE_FINALIZE.value: {"ignore_result": True},
CeleryTaskName.VIDEO_UPSCALE_RECOVER.value: {"ignore_result": True},
CeleryTaskName.RECOVER_CREATE.value: {"ignore_result": True},
CeleryTaskName.MODULE_ASYNC_RECOVERY.value: {"ignore_result": True},
CeleryTaskName.SHOT_ANALYSIS_RECOVERY.value: {"ignore_result": True},
CeleryTaskName.SHOT_SPLIT_RECOVERY.value: {"ignore_result": True},
},
worker_prefetch_multiplier=1,
worker_cancel_long_running_tasks_on_connection_loss=True,
broker_transport_options={
"visibility_timeout": max(
3600,
int(settings.VIDEO_UPSCALE_LOCAL_TIMEOUT_SECONDS or 3600) + 600,
int(settings.VIDEO_UPSCALE_REMOTE_RESULT_DOWNLOAD_TIMEOUT_SECONDS or 600) + 600,
int(settings.SHOT_ANALYSIS_TIME_LIMIT_SECONDS or 3900) + 600,
),
"queue_order_strategy": "priority",
"priority_steps": list(range(10)),
"sep": ":",
},
task_routes={
CeleryTaskName.CHATAPI_CREATE.value: {"queue": CeleryQueue.GEN_CHATAPI_CREATE.value},
CeleryTaskName.POLL_GENERATION.value: {"queue": CeleryQueue.GEN_PROVIDER_POLL.value},
CeleryTaskName.DOWNLOAD_GENERATION_RESULT.value: {"queue": CeleryQueue.GEN_RESULT_DOWNLOAD.value},
CeleryTaskName.VIDEO_UPSCALE_EXECUTE_LOCAL.value: {
"queue": settings.VIDEO_UPSCALE_LOCAL_QUEUE or CeleryQueue.GEN_VIDEO_UPSCALE_LOCAL.value
},
CeleryTaskName.VIDEO_UPSCALE_SUBMIT_REMOTE.value: {
"queue": settings.VIDEO_UPSCALE_REMOTE_QUEUE or CeleryQueue.GEN_VIDEO_UPSCALE_REMOTE.value
},
CeleryTaskName.VIDEO_UPSCALE_POLL_REMOTE.value: {
"queue": settings.VIDEO_UPSCALE_REMOTE_QUEUE or CeleryQueue.GEN_VIDEO_UPSCALE_REMOTE.value
},
CeleryTaskName.VIDEO_UPSCALE_DOWNLOAD_REMOTE_RESULT.value: {
"queue": settings.VIDEO_UPSCALE_REMOTE_QUEUE or CeleryQueue.GEN_VIDEO_UPSCALE_REMOTE.value
},
CeleryTaskName.VIDEO_UPSCALE_FINALIZE.value: {
"queue": settings.VIDEO_UPSCALE_LOCAL_QUEUE or CeleryQueue.GEN_VIDEO_UPSCALE_LOCAL.value
},
CeleryTaskName.VIDEO_UPSCALE_RECOVER.value: {"queue": RECOVERY_QUEUE},
CeleryTaskName.DISPATCH_DUE_POLL.value: {"queue": RECOVERY_QUEUE},
"hot_opening.start_image_prompt_optimize": {"queue": CeleryQueue.GEN_CHATAPI_CREATE.value},
"hot_opening.start_video_prompt_optimize": {"queue": CeleryQueue.GEN_CHATAPI_CREATE.value},
CeleryTaskName.SHOT_ANALYZE_ORIGINAL.value: {"queue": CeleryQueue.GEN_SHOT_ANALYSIS.value},
CeleryTaskName.SHOT_ANALYZE_CUSTOM_SEGMENT.value: {"queue": CeleryQueue.GEN_SHOT_ANALYSIS.value},
CeleryTaskName.SHOT_SPLIT_ONE.value: {"queue": CeleryQueue.GEN_SHOT_SPLIT.value},
"shot_replicate.start_image_prompt_optimize": {"queue": CeleryQueue.GEN_CHATAPI_CREATE.value},
"shot_replicate.start_video_prompt_optimize": {"queue": CeleryQueue.GEN_CHATAPI_CREATE.value},
"module_generation_v2.start_video_prompt_optimize": {"queue": CeleryQueue.GEN_CHATAPI_CREATE.value},
# 恢复扫描统一走独立队列,避免占用下载/轮询/创建业务 worker。
CeleryTaskName.STARTUP_RECOVERY.value: {"queue": RECOVERY_QUEUE},
CeleryTaskName.SHOT_SPLIT_RECOVERY.value: {"queue": RECOVERY_QUEUE},
CeleryTaskName.SHOT_ANALYSIS_RECOVERY.value: {"queue": RECOVERY_QUEUE},
CeleryTaskName.RECOVER_DOWNLOAD.value: {"queue": RECOVERY_QUEUE},
CeleryTaskName.RECOVER_GENERATION.value: {"queue": RECOVERY_QUEUE},
CeleryTaskName.RECOVER_CREATE.value: {"queue": RECOVERY_QUEUE},
CeleryTaskName.MODULE_ASYNC_RECOVERY.value: {"queue": RECOVERY_QUEUE},
CeleryTaskName.CELERY_RUNTIME_RECONCILE.value: {"queue": RECOVERY_QUEUE},
CeleryTaskName.CELERY_RUNTIME_GC.value: {"queue": RECOVERY_QUEUE},
CeleryTaskName.PRIVATE_PORTRAIT_POLL_ASSET.value: {"queue": CeleryQueue.GEN_PRIVATE_PORTRAIT.value},
CeleryTaskName.PRIVATE_PORTRAIT_SYNC_DUE_ASSETS.value: {"queue": CeleryQueue.GEN_PRIVATE_PORTRAIT.value},
CeleryTaskName.PRIVATE_PORTRAIT_DELETE_ASSET.value: {"queue": CeleryQueue.GEN_PRIVATE_PORTRAIT.value},
CeleryTaskName.PRIVATE_PORTRAIT_DELETE_GROUP.value: {"queue": CeleryQueue.GEN_PRIVATE_PORTRAIT.value},
CeleryTaskName.PRIVATE_PORTRAIT_DELETE_PROJECT.value: {"queue": CeleryQueue.GEN_PRIVATE_PORTRAIT.value},
CeleryTaskName.PRIVATE_PORTRAIT_RECOVER_REMOTE_DELETES.value: {"queue": CeleryQueue.GEN_PRIVATE_PORTRAIT.value},
CeleryTaskName.VP_V3_POLL_ASSET.value: {"queue": CeleryQueue.GEN_PRIVATE_PORTRAIT.value},
CeleryTaskName.VP_V3_SYNC_DUE_ASSETS.value: {"queue": CeleryQueue.GEN_PRIVATE_PORTRAIT.value},
CeleryTaskName.VP_V3_DELETE_ASSET.value: {"queue": CeleryQueue.GEN_PRIVATE_PORTRAIT.value},
CeleryTaskName.VP_V3_DELETE_PROJECT.value: {"queue": CeleryQueue.GEN_PRIVATE_PORTRAIT.value},
CeleryTaskName.VP_V3_RECOVER_REMOTE_DELETES.value: {"queue": CeleryQueue.GEN_PRIVATE_PORTRAIT.value},
},
)
else:
celery_app = None
async def _try_acquire_startup_recovery_lock() -> bool:
"""任意 worker 启动时都可尝试抢恢复投递锁,避免依赖 hostname 命名。"""
from app.services.redis_registry_service import redis_acquire_lock
token = await redis_acquire_lock(
lock_key=settings.CELERY_STARTUP_RECOVERY_LOCK_KEY,
ttl_seconds=int(settings.CELERY_STARTUP_RECOVERY_LOCK_TTL_SECONDS or 120),
log_context="celery_startup_recovery",
)
return bool(token)
def _worker_name_from_sender(sender=None, **kwargs) -> str:
candidates = (
getattr(sender, "hostname", None),
getattr(sender, "name", None),
kwargs.get("hostname"),
kwargs.get("nodename"),
)
instance = kwargs.get("instance")
if instance is not None:
candidates += (
getattr(instance, "hostname", None),
getattr(instance, "name", None),
)
for value in candidates:
normalized = str(value or "").strip()
if normalized:
return normalized
return ""
def _worker_runtime_metadata(sender=None) -> dict:
queues: list[str] = []
try:
consumer = getattr(sender, "consumer", None)
task_consumer = getattr(consumer, "task_consumer", None)
for queue in list(getattr(task_consumer, "queues", None) or []):
name = str(getattr(queue, "name", queue) or "").strip()
if name and name not in queues:
queues.append(name)
except Exception:
pass
pool = getattr(sender, "pool", None)
pool_type = type(pool).__name__ if pool is not None else None
configured_concurrency = None
for value in (
getattr(pool, "limit", None),
getattr(sender, "concurrency", None),
):
try:
parsed = int(value)
except (TypeError, ValueError):
continue
if parsed > 0:
configured_concurrency = parsed
break
return {
"queues": queues,
"pool_type": pool_type,
"configured_concurrency": configured_concurrency,
}
@celeryd_init.connect
def on_celeryd_init(sender=None, instance=None, **kwargs):
"""尽早生成 Worker 主实例 token,确保 prefork 子进程继承。"""
try:
from app.services.celery_runtime.worker_service import initialize_worker_main_identity
initialize_worker_main_identity(
_worker_name_from_sender(sender, instance=instance, **kwargs) or None,
before_pool=True,
)
except Exception:
logger.exception("Celery Worker 主实例身份初始化失败。signal=celeryd_init")
@worker_init.connect
def on_worker_init(sender=None, **kwargs):
"""worker_init 幂等兜底,仍处于进程池创建之前。"""
try:
from app.services.celery_runtime.worker_service import initialize_worker_main_identity
initialize_worker_main_identity(
_worker_name_from_sender(sender, **kwargs) or None,
before_pool=True,
)
except Exception:
logger.exception("Celery Worker 主实例身份初始化失败。signal=worker_init")
@worker_ready.connect
def on_worker_ready(sender=None, **kwargs):
"""注册当前 Worker 主实例,并协调实例级与全局启动恢复。"""
if celery_app is None:
return
worker_name = _worker_name_from_sender(sender, **kwargs)
metadata = _worker_runtime_metadata(sender)
current_identity = None
try:
from app.services.celery_runtime.worker_service import (
initialize_worker_main_identity,
register_worker_instance,
)
# 正常情况下 token 已在 worker_init 前创建;这里仅做 late fallback。
initialize_worker_main_identity(worker_name or None, before_pool=False)
current_identity = run_async(
register_worker_instance(
worker_name=worker_name,
queues=metadata["queues"],
pool_type=metadata["pool_type"],
configured_concurrency=metadata["configured_concurrency"],
)
)
except Exception:
# Redis 或身份注册失败不能阻塞 Worker 启动,任务级执行锁仍会 fail-closed。
logger.exception("Celery Worker 主实例注册失败。worker_name=%s", worker_name)
if current_identity is not None:
try:
from app.services.celery_runtime.recovery_service import (
mark_stale_worker_instance_candidates,
)
run_async(
mark_stale_worker_instance_candidates(
worker_name=current_identity.worker_name,
current_worker_instance_id=current_identity.worker_instance_id,
supports_targeted_recovery=current_identity.supports_targeted_recovery,
)
)
except Exception:
logger.exception(
"Worker 旧主实例任务候选标记失败。worker_name=%s worker_instance_id=%s",
current_identity.worker_name,
current_identity.worker_instance_id,
)
if not bool(getattr(settings, "CELERY_STARTUP_RECOVERY_ENABLED", True)):
logger.info("启动容灾恢复已关闭。CELERY_STARTUP_RECOVERY_ENABLED=false")
return
try:
if not run_async(_try_acquire_startup_recovery_lock()):
return
except Exception:
logger.exception("启动容灾恢复锁获取失败,已跳过本次自动恢复投递")
return
try:
from app.services.celery_runtime.recovery_service import set_startup_barrier
from app.tasks.generation_recovery_tasks import startup_recovery_once
from app.tasks.api_recovery_tasks import api_generation_recover_tasks_once
run_async(set_startup_barrier())
countdown = max(0, int(settings.CELERY_STARTUP_RECOVERY_DELAY_SECONDS or 30))
startup_recovery_once.apply_async(
countdown=countdown,
queue=RECOVERY_QUEUE,
priority=settings.DOWNLOAD_TASK_PRIORITY_RECOVER,
)
# API v3 任务恢复(延迟 35 秒执行,避免与其他恢复任务冲突)
api_generation_recover_tasks_once.apply_async(
countdown=countdown + 5,
queue=RECOVERY_QUEUE,
priority=settings.DOWNLOAD_TASK_PRIORITY_RECOVER,
)
logger.info(
"启动容灾恢复协调任务已投递(含 API v3)。queue=%s countdown=%s",
RECOVERY_QUEUE,
countdown,
)
except Exception:
logger.exception("启动容灾恢复协调任务投递失败")
@heartbeat_sent.connect
def on_worker_heartbeat_sent(sender=None, **kwargs):
"""刷新主实例 heartbeat,并低频扫描已到期的旧主实例。"""
try:
from app.services.celery_runtime.worker_service import (
claim_worker_heartbeat_slot,
claim_worker_stale_scan_slot,
heartbeat_current_worker_instance,
registered_worker_identity,
)
heartbeat_ok = True
if claim_worker_heartbeat_slot():
heartbeat_ok = bool(run_async(heartbeat_current_worker_instance()))
if not heartbeat_ok or not claim_worker_stale_scan_slot():
return
identity = registered_worker_identity()
if identity is None:
return
from app.services.celery_runtime.recovery_service import (
mark_stale_worker_instance_candidates,
)
run_async(
mark_stale_worker_instance_candidates(
worker_name=identity.worker_name,
current_worker_instance_id=identity.worker_instance_id,
supports_targeted_recovery=identity.supports_targeted_recovery,
emit_duplicate_log=False,
)
)
except Exception:
logger.warning("Celery Worker 主实例 heartbeat/旧实例扫描失败", exc_info=True)
@worker_shutdown.connect
def on_worker_shutdown(sender=None, **kwargs):
"""优雅退出时撤销活跃 Worker key;旧任务集合保留给恢复流程。"""
try:
from app.services.celery_runtime.worker_service import unregister_current_worker_instance
run_async(unregister_current_worker_instance())
except Exception:
logger.debug("Celery Worker 主实例注销失败", exc_info=True)
finally:
close_loop()
@worker_process_init.connect
def on_worker_process_init(**kwargs):
"""prefork 子进程保留主 token,同时重建执行进程缓存和异步连接。"""
try:
from app.services.celery_runtime.worker_service import reset_process_identity_cache
reset_process_identity_cache()
except Exception:
pass
try:
run_async(engine.dispose())
except Exception:
pass
@worker_process_shutdown.connect
def on_worker_process_shutdown(**kwargs):
"""子进程退出前关闭连接池、Redis 注册表连接和 event loop。"""
try:
run_async(engine.dispose())
except Exception:
pass
try:
from app.services.redis_registry_service import close_registry_redis
run_async(close_registry_redis())
except Exception:
pass
finally:
close_loop()