| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899 |
- """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,
- }
|