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.celery_runtime_tasks", "app.tasks.credit_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["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["credit-maintenance-every-minute"] = { "task": CeleryTaskName.CREDIT_MAINTENANCE.value, "schedule": 60, "options": {"queue": CeleryQueue.GEN_CREDIT_MAINTENANCE.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}, CeleryTaskName.CREDIT_MAINTENANCE.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.CREDIT_MAINTENANCE.value: {"queue": CeleryQueue.GEN_CREDIT_MAINTENANCE.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 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, ) logger.info( "启动容灾恢复协调任务已投递。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()