"""Compatibility facade over the durable pipeline control plane.""" from __future__ import annotations from typing import Any from api.services.pipeline import get_run, pipeline_health from supply_infra.config import get_infra_settings from supply_infra.pipeline.dag import PIPELINE_STEPS from supply_infra.pipeline.dates import CHINA_TIMEZONE, next_schedule_china 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, ) from supply_infra.scheduler.jobs.publish_videos_from_discovery import ( publish_videos_from_discovery, ) def list_triggerable_jobs() -> list[dict[str, Any]]: settings = get_infra_settings() return [ { "id": SUPPLY_PIPELINE_JOB_ID, "name": SUPPLY_PIPELINE_JOB_NAME, "description": "提交严格门禁的 11 步持久化供给流水线(不含 AIGC 发布)", "accepts_biz_dt": True, "deprecated": False, "steps": [step.key for step in PIPELINE_STEPS], }, { "id": AIGC_PUBLISH_JOB_ID, "name": AIGC_PUBLISH_JOB_NAME, "description": ( "扫描未发布的 primary 候选并分发到 AIGC;" f"调度默认每 {settings.scheduler_aigc_publish_interval_minutes} 分钟执行" ), "accepts_biz_dt": True, "deprecated": False, "steps": [], }, ] def run_scheduler_job( job_id: str, *, biz_dt: str | None = None, ) -> dict[str, Any]: if job_id == SUPPLY_PIPELINE_JOB_ID: return submit_pipeline_run( biz_dt=biz_dt, trigger_type="api", trigger_source="legacy_scheduler_api", ).to_dict() if job_id == AIGC_PUBLISH_JOB_ID: return publish_videos_from_discovery( biz_dt=biz_dt, skip_published=True, any_biz_dt=biz_dt is None, ) raise KeyError(job_id) def get_scheduler_job_run(run_id: str) -> dict[str, Any] | None: return get_run(run_id) def run_supply_pipeline(*, biz_dt: str | None = None) -> dict[str, Any]: return run_scheduler_job(SUPPLY_PIPELINE_JOB_ID, biz_dt=biz_dt) def scheduler_status() -> dict[str, Any]: settings = get_infra_settings() health = pipeline_health() return { "enabled": settings.scheduler_enabled, "running": health["scheduler_running"], "embedded": False, "deprecated": True, "timezone": settings.scheduler_timezone, "jobs": [ { "id": SUPPLY_PIPELINE_JOB_ID, "name": SUPPLY_PIPELINE_JOB_NAME, "next_run_time": next_schedule_china(settings=settings) .replace(tzinfo=CHINA_TIMEZONE) .isoformat(), }, { "id": AIGC_PUBLISH_JOB_ID, "name": AIGC_PUBLISH_JOB_NAME, "interval_minutes": settings.scheduler_aigc_publish_interval_minutes, }, ], **health, }