from __future__ import annotations import logging from typing import TYPE_CHECKING from apscheduler.schedulers.background import BackgroundScheduler from apscheduler.triggers.cron import CronTrigger from apscheduler.triggers.interval import IntervalTrigger from supply_infra.config import get_infra_settings from supply_infra.pipeline.run_service import submit_pipeline_run from supply_infra.scheduler.constants import ( AIGC_PUBLISH_JOB_ID, AIGC_PUBLISH_JOB_NAME, SUPPLY_PIPELINE_JOB_ID, SUPPLY_PIPELINE_JOB_NAME, ) if TYPE_CHECKING: from apscheduler.schedulers.base import BaseScheduler logger = logging.getLogger(__name__) _scheduler: BackgroundScheduler | None = None def submit_scheduled_pipeline() -> dict: """Cron callback: create a durable batch and return immediately.""" return submit_pipeline_run( trigger_type="cron", trigger_source="scheduler", ).to_dict() def submit_aigc_publish() -> dict: """Interval callback: publish any unpublished primary discovery candidates.""" from supply_infra.scheduler.jobs.publish_videos_from_discovery import ( publish_videos_from_discovery, ) result = publish_videos_from_discovery( skip_published=True, any_biz_dt=True, ) logger.info( "AIGC publish job finished: success=%s candidates=%s batches=%s failed=%s", result.get("success"), result.get("candidate_count"), result.get("batch_count"), result.get("failed_batch_count"), ) return result def create_scheduler() -> BackgroundScheduler: """Create a Scheduler for the daily pipeline and AIGC publish polling.""" settings = get_infra_settings() tz = settings.scheduler_timezone scheduler = BackgroundScheduler(timezone=tz) scheduler.add_job( submit_scheduled_pipeline, trigger=CronTrigger( hour=settings.scheduler_cron_hour, minute=settings.scheduler_cron_minute, timezone=tz, ), id=SUPPLY_PIPELINE_JOB_ID, name=SUPPLY_PIPELINE_JOB_NAME, replace_existing=True, max_instances=1, coalesce=True, misfire_grace_time=3600, ) scheduler.add_job( submit_aigc_publish, trigger=IntervalTrigger( minutes=settings.scheduler_aigc_publish_interval_minutes, timezone=tz, ), id=AIGC_PUBLISH_JOB_ID, name=AIGC_PUBLISH_JOB_NAME, replace_existing=True, max_instances=1, coalesce=True, misfire_grace_time=300, ) logger.info( "Scheduler configured with %d job(s) | timezone=%s | pipeline_cron=%02d:%02d " "| aigc_publish_interval=%dm", len(scheduler.get_jobs()), tz, settings.scheduler_cron_hour, settings.scheduler_cron_minute, settings.scheduler_aigc_publish_interval_minutes, ) return scheduler def _log_next_runs(scheduler: BaseScheduler) -> None: for job in scheduler.get_jobs(): # apscheduler 3.x: 在 scheduler.start() 之前,Job.next_run_time 访问会抛 AttributeError # 这里用 getattr 兜底,保证启动阶段不因日志而中断。 next_run_time = getattr(job, "next_run_time", None) logger.info(" - %s | next run: %s", job.name, next_run_time) def start_scheduler() -> BackgroundScheduler | None: """Start the independent scheduler process (idempotent).""" global _scheduler settings = get_infra_settings() if not settings.scheduler_enabled: logger.warning("Scheduler is disabled (SCHEDULER_ENABLED=false)") return None if _scheduler is not None and _scheduler.running: return _scheduler _scheduler = create_scheduler() logger.info("Starting scheduler...") _scheduler.start() _log_next_runs(_scheduler) return _scheduler def stop_scheduler() -> None: """Shut down the background scheduler if running.""" global _scheduler if _scheduler is None or not _scheduler.running: return logger.info("Stopping scheduler...") _scheduler.shutdown(wait=False) _scheduler = None logger.info("Scheduler stopped.")