xueyiming пре 3 дана
родитељ
комит
1a4da00b5e

+ 1 - 30
api/app.py

@@ -120,15 +120,6 @@ class TriggerSchedulerJobBody(BaseModel):
         pattern=r"^\d{8}$",
         description="业务日 YYYYMMDD;省略则取当天",
     )
-    partition_date: str | None = Field(
-        default=None,
-        pattern=r"^\d{8}$",
-        description="ODPS 分区日 YYYYMMDD;全局树默认取 biz_dt 前一日",
-    )
-    wait: bool = Field(
-        default=False,
-        description="是否同步等待执行完成;默认 false 表示后台执行",
-    )
 
 
 class AgentDocumentInjectionBody(BaseModel):
@@ -176,34 +167,14 @@ def trigger_scheduler_job(
         pattern=r"^\d{8}$",
         description="业务日 YYYYMMDD(与 body 二选一,body 优先)",
     ),
-    partition_date: str | None = Query(
-        default=None,
-        pattern=r"^\d{8}$",
-        description="ODPS 分区日 YYYYMMDD(与 body 二选一,body 优先)",
-    ),
-    wait: bool = Query(
-        default=False,
-        description="是否同步等待执行完成",
-    ),
 ) -> dict:
     """Manually trigger a scheduler job."""
     resolved_biz_dt = body.biz_dt if body and body.biz_dt is not None else biz_dt
-    resolved_partition_date = (
-        body.partition_date if body and body.partition_date is not None else partition_date
-    )
-    resolved_wait = body.wait if body is not None else wait
 
     try:
-        return run_scheduler_job(
-            job_id,
-            biz_dt=resolved_biz_dt,
-            partition_date=resolved_partition_date,
-            wait=resolved_wait,
-        )
+        return run_scheduler_job(job_id, biz_dt=resolved_biz_dt)
     except KeyError:
         raise HTTPException(status_code=404, detail=f"scheduler job not found: {job_id}") from None
-    except ValueError as exc:
-        raise HTTPException(status_code=409, detail=str(exc)) from None
 
 
 @app.get("/api/scheduler/runs/{run_id}")

+ 0 - 6
api/services/scheduler.py

@@ -18,7 +18,6 @@ def list_triggerable_jobs() -> list[dict[str, Any]]:
             "name": SUPPLY_PIPELINE_JOB_NAME,
             "description": "提交严格门禁的 12 步持久化供给流水线",
             "accepts_biz_dt": True,
-            "accepts_partition_date": False,
             "deprecated": False,
             "steps": [step.key for step in PIPELINE_STEPS],
         }
@@ -29,14 +28,9 @@ def run_scheduler_job(
     job_id: str,
     *,
     biz_dt: str | None = None,
-    partition_date: str | None = None,
-    wait: bool = False,
 ) -> dict[str, Any]:
-    del partition_date
     if job_id != SUPPLY_PIPELINE_JOB_ID:
         raise KeyError(job_id)
-    if wait:
-        raise ValueError("wait=true is disabled for the durable API; poll by run_id")
     return submit_pipeline_run(
         biz_dt=biz_dt,
         trigger_type="api",

+ 0 - 2
supply_infra/db/models/__init__.py

@@ -27,7 +27,6 @@ from supply_infra.db.models.pipeline_lock import PipelineLock
 from supply_infra.db.models.pipeline_outbox import PipelineOutbox
 from supply_infra.db.models.pipeline_run import PipelineRun
 from supply_infra.db.models.pipeline_step_run import PipelineStepRun
-from supply_infra.db.models.scheduler_job_execution import SchedulerJobExecution
 from supply_infra.db.models.video_discovery import (
     VideoDiscoveryCandidate,
     VideoDiscoveryRun,
@@ -58,7 +57,6 @@ __all__ = [
     "PipelineOutbox",
     "PipelineRun",
     "PipelineStepRun",
-    "SchedulerJobExecution",
     "VideoDiscoveryCandidate",
     "VideoDiscoveryRun",
     "VideoDiscoverySearch",

+ 0 - 3
supply_infra/db/models/pipeline_run.py

@@ -30,14 +30,11 @@ class PipelineRun(Base):
     biz_dt: Mapped[str] = mapped_column(String(8), nullable=False)
     trigger_type: Mapped[str] = mapped_column(String(24), nullable=False)
     trigger_source: Mapped[str | None] = mapped_column(String(64), nullable=True)
-    triggered_by: Mapped[str | None] = mapped_column(String(128), nullable=True)
     trigger_reason: Mapped[str | None] = mapped_column(Text, nullable=True)
-    parent_run_id: Mapped[str | None] = mapped_column(String(36), nullable=True)
     run_mode: Mapped[str] = mapped_column(String(24), nullable=False, default="full")
     dry_run: Mapped[bool] = mapped_column(Boolean, nullable=False, default=False)
     status: Mapped[str] = mapped_column(String(32), nullable=False, default="queued")
     current_step: Mapped[str | None] = mapped_column(String(64), nullable=True)
-    scheduled_for: Mapped[datetime | None] = mapped_column(DateTime, nullable=True)
     deadline_at: Mapped[datetime | None] = mapped_column(DateTime, nullable=True)
     started_at: Mapped[datetime | None] = mapped_column(DateTime, nullable=True)
     finished_at: Mapped[datetime | None] = mapped_column(DateTime, nullable=True)

+ 0 - 1
supply_infra/db/models/pipeline_step_run.py

@@ -68,7 +68,6 @@ class PipelineStepRun(Base):
     lease_until: Mapped[datetime | None] = mapped_column(DateTime, nullable=True)
     exit_code: Mapped[int | None] = mapped_column(Integer, nullable=True)
     result_summary_json: Mapped[dict[str, Any] | None] = mapped_column(JSON, nullable=True)
-    metrics_json: Mapped[dict[str, Any] | None] = mapped_column(JSON, nullable=True)
     error_code: Mapped[str | None] = mapped_column(String(64), nullable=True)
     error_message: Mapped[str | None] = mapped_column(Text, nullable=True)
     log_uri: Mapped[str | None] = mapped_column(String(1024), nullable=True)

+ 0 - 51
supply_infra/db/models/scheduler_job_execution.py

@@ -1,51 +0,0 @@
-from __future__ import annotations
-
-from datetime import datetime
-
-from sqlalchemy import BigInteger, Numeric, String, Text, func
-from sqlalchemy.orm import Mapped, mapped_column
-
-from supply_infra.db.base import Base
-
-
-class SchedulerJobExecution(Base):
-    """定时任务执行记录(开始/结束各记一条,同一 run_id 关联)。"""
-
-    __tablename__ = "scheduler_job_execution"
-
-    id: Mapped[int] = mapped_column(BigInteger, primary_key=True, autoincrement=True)
-    run_id: Mapped[str] = mapped_column(
-        String(36),
-        nullable=False,
-        index=True,
-        comment="同一次执行的唯一标识",
-    )
-    job_name: Mapped[str] = mapped_column(String(128), nullable=False, comment="定时任务名称")
-    job_id: Mapped[str | None] = mapped_column(String(64), nullable=True, comment="调度器 job id")
-    status: Mapped[str] = mapped_column(
-        String(32),
-        nullable=False,
-        comment="状态: started/finished/failed/skipped",
-    )
-    event_time: Mapped[datetime] = mapped_column(nullable=False, comment="事件发生时间")
-    biz_dt: Mapped[str | None] = mapped_column(String(8), nullable=True, comment="业务日 YYYYMMDD")
-    started_at: Mapped[datetime | None] = mapped_column(nullable=True, comment="任务开始时间")
-    finished_at: Mapped[datetime | None] = mapped_column(nullable=True, comment="任务结束时间")
-    duration_seconds: Mapped[float | None] = mapped_column(
-        Numeric(10, 2),
-        nullable=True,
-        comment="执行耗时(秒)",
-    )
-    error_message: Mapped[str | None] = mapped_column(Text, nullable=True, comment="失败原因")
-    detail: Mapped[str | None] = mapped_column(Text, nullable=True, comment="执行结果摘要 JSON")
-    create_time: Mapped[datetime] = mapped_column(
-        nullable=False,
-        server_default=func.now(),
-        comment="创建时间",
-    )
-    update_time: Mapped[datetime] = mapped_column(
-        nullable=False,
-        server_default=func.now(),
-        onupdate=func.now(),
-        comment="更新时间",
-    )

+ 0 - 4
supply_infra/db/repositories/__init__.py

@@ -40,9 +40,6 @@ from supply_infra.db.repositories.pipeline_lock_repo import PipelineLockReposito
 from supply_infra.db.repositories.pipeline_outbox_repo import PipelineOutboxRepository
 from supply_infra.db.repositories.pipeline_run_repo import PipelineRunRepository
 from supply_infra.db.repositories.pipeline_step_run_repo import PipelineStepRunRepository
-from supply_infra.db.repositories.scheduler_job_execution_repo import (
-    SchedulerJobExecutionRepository,
-)
 from supply_infra.db.repositories.video_discovery_repo import VideoDiscoveryRepository
 
 __all__ = [
@@ -68,6 +65,5 @@ __all__ = [
     "PipelineOutboxRepository",
     "PipelineRunRepository",
     "PipelineStepRunRepository",
-    "SchedulerJobExecutionRepository",
     "VideoDiscoveryRepository",
 ]

+ 0 - 50
supply_infra/db/repositories/scheduler_job_execution_repo.py

@@ -1,50 +0,0 @@
-from __future__ import annotations
-
-from datetime import datetime
-
-from sqlalchemy import select
-
-from supply_infra.db.models.scheduler_job_execution import SchedulerJobExecution
-from supply_infra.db.repositories.base import BaseRepository
-
-
-class SchedulerJobExecutionRepository(BaseRepository[SchedulerJobExecution]):
-    model = SchedulerJobExecution
-
-    def log_event(
-        self,
-        *,
-        run_id: str,
-        job_name: str,
-        status: str,
-        event_time: datetime,
-        job_id: str | None = None,
-        biz_dt: str | None = None,
-        started_at: datetime | None = None,
-        finished_at: datetime | None = None,
-        duration_seconds: float | None = None,
-        error_message: str | None = None,
-        detail: str | None = None,
-    ) -> SchedulerJobExecution:
-        entity = SchedulerJobExecution(
-            run_id=run_id,
-            job_name=job_name,
-            job_id=job_id,
-            status=status,
-            event_time=event_time,
-            biz_dt=biz_dt,
-            started_at=started_at,
-            finished_at=finished_at,
-            duration_seconds=duration_seconds,
-            error_message=error_message,
-            detail=detail,
-        )
-        return self.add(entity)
-
-    def list_recent(self, *, limit: int = 50) -> list[SchedulerJobExecution]:
-        stmt = (
-            select(SchedulerJobExecution)
-            .order_by(SchedulerJobExecution.event_time.desc(), SchedulerJobExecution.id.desc())
-            .limit(limit)
-        )
-        return list(self.session.scalars(stmt).all())

+ 0 - 4
supply_infra/pipeline/run_service.py

@@ -72,7 +72,6 @@ def submit_pipeline_run(
     trigger_type: str = "api",
     trigger_source: str | None = None,
     trigger_reason: str | None = None,
-    scheduled_for: datetime | None = None,
     settings: InfraSettings | None = None,
 ) -> RunSubmission:
     active = settings or get_infra_settings()
@@ -112,12 +111,10 @@ def submit_pipeline_run(
                     "biz_dt": resolved_biz_dt,
                     "trigger_type": trigger_type,
                     "trigger_source": trigger_source or trigger_type,
-                    "triggered_by": None,
                     "trigger_reason": trigger_reason,
                     "run_mode": "full",
                     "dry_run": False,
                     "status": initial_status,
-                    "scheduled_for": scheduled_for,
                     "deadline_at": deadline_for_trigger(
                         trigger_type,
                         settings=active,
@@ -340,7 +337,6 @@ def serialize_run(run: PipelineRun) -> dict[str, Any]:
         "dry_run": run.dry_run,
         "status": run.status,
         "current_step": run.current_step,
-        "scheduled_for": _iso(run.scheduled_for),
         "deadline_at": _iso(run.deadline_at),
         "started_at": _iso(run.started_at),
         "finished_at": _iso(run.finished_at),

+ 0 - 38
supply_infra/scheduler/app.py

@@ -84,44 +84,6 @@ def start_scheduler() -> BackgroundScheduler | None:
     return _scheduler
 
 
-def get_scheduler_status() -> dict:
-    """Return scheduler runtime status for health checks."""
-    settings = get_infra_settings()
-    if not settings.scheduler_enabled:
-        return {
-            "enabled": False,
-            "running": False,
-            "timezone": settings.scheduler_timezone,
-            "jobs": [],
-        }
-
-    if _scheduler is None or not _scheduler.running:
-        return {
-            "enabled": True,
-            "running": False,
-            "timezone": settings.scheduler_timezone,
-            "jobs": [],
-        }
-
-    jobs = []
-    for job in _scheduler.get_jobs():
-        next_run = getattr(job, "next_run_time", None)
-        jobs.append(
-            {
-                "id": job.id,
-                "name": job.name,
-                "next_run_time": next_run.isoformat() if next_run else None,
-            }
-        )
-
-    return {
-        "enabled": True,
-        "running": True,
-        "timezone": settings.scheduler_timezone,
-        "jobs": jobs,
-    }
-
-
 def stop_scheduler() -> None:
     """Shut down the background scheduler if running."""
     global _scheduler

+ 0 - 4
supply_infra/scheduler/constants.py

@@ -5,7 +5,3 @@ SUPPLY_PIPELINE_JOB_NAME = "供给数据流水线"
 
 # find_agent:当日全部 S/A 需求(有拓展点位),单线程串行找片;有效视频满 200 提前结束
 PIPELINE_FIND_AGENT_WORKERS = 1
-
-# 供给流水线每日触发时间(Asia/Shanghai)
-PIPELINE_CRON_HOUR = 14
-PIPELINE_CRON_MINUTE = 30

+ 0 - 120
supply_infra/scheduler/job_execution.py

@@ -1,120 +0,0 @@
-"""定时任务执行记录写入辅助。"""
-from __future__ import annotations
-
-import json
-import logging
-import uuid
-from typing import Any
-
-from supply_infra.db.repositories.scheduler_job_execution_repo import (
-    SchedulerJobExecutionRepository,
-)
-from supply_infra.db.session import get_session
-from supply_infra.pipeline.dates import china_now
-
-logger = logging.getLogger(__name__)
-
-STATUS_STARTED = "started"
-STATUS_FINISHED = "finished"
-STATUS_FAILED = "failed"
-STATUS_SKIPPED = "skipped"
-
-
-def _serialize_detail(payload: Any) -> str | None:
-    if payload is None:
-        return None
-    return json.dumps(payload, ensure_ascii=False, default=str)
-
-
-class JobExecutionRecorder:
-    """在任务开始/结束时各写一条执行记录。"""
-
-    def __init__(
-        self,
-        *,
-        job_name: str,
-        job_id: str | None = None,
-        biz_dt: str | None = None,
-    ) -> None:
-        self.run_id = str(uuid.uuid4())
-        self.job_name = job_name
-        self.job_id = job_id
-        self.biz_dt = biz_dt
-        self.started_at = china_now()
-
-    def record_started(self) -> None:
-        try:
-            with get_session() as session:
-                SchedulerJobExecutionRepository(session).log_event(
-                    run_id=self.run_id,
-                    job_name=self.job_name,
-                    job_id=self.job_id,
-                    status=STATUS_STARTED,
-                    event_time=self.started_at,
-                    biz_dt=self.biz_dt,
-                    started_at=self.started_at,
-                )
-        except Exception:
-            logger.exception("Failed to record job start: job_name=%s run_id=%s", self.job_name, self.run_id)
-
-    def record_finished(
-        self,
-        *,
-        success: bool,
-        result: dict[str, Any] | None = None,
-        error_message: str | None = None,
-    ) -> None:
-        finished_at = china_now()
-        duration_seconds = round((finished_at - self.started_at).total_seconds(), 2)
-        status = STATUS_FINISHED if success else STATUS_FAILED
-
-        try:
-            with get_session() as session:
-                SchedulerJobExecutionRepository(session).log_event(
-                    run_id=self.run_id,
-                    job_name=self.job_name,
-                    job_id=self.job_id,
-                    status=status,
-                    event_time=finished_at,
-                    biz_dt=self.biz_dt,
-                    started_at=self.started_at,
-                    finished_at=finished_at,
-                    duration_seconds=duration_seconds,
-                    error_message=error_message,
-                    detail=_serialize_detail(result),
-                )
-        except Exception:
-            logger.exception(
-                "Failed to record job finish: job_name=%s run_id=%s status=%s",
-                self.job_name,
-                self.run_id,
-                status,
-            )
-
-
-def record_skipped(
-    *,
-    job_name: str,
-    job_id: str | None = None,
-    biz_dt: str | None = None,
-    reason: str,
-) -> None:
-    """记录因并发等原因跳过的执行。"""
-    now = china_now()
-    run_id = str(uuid.uuid4())
-    try:
-        with get_session() as session:
-            SchedulerJobExecutionRepository(session).log_event(
-                run_id=run_id,
-                job_name=job_name,
-                job_id=job_id,
-                status=STATUS_SKIPPED,
-                event_time=now,
-                biz_dt=biz_dt,
-                started_at=now,
-                finished_at=now,
-                duration_seconds=0,
-                error_message=reason,
-            )
-    except Exception:
-        logger.exception("Failed to record skipped job: job_name=%s reason=%s", job_name, reason)

+ 1 - 5
supply_infra/scheduler/jobs/demand_pool/__init__.py

@@ -1,5 +1 @@
-"""策略需求池同步任务(ODPS → MySQL + 归属/热度/树权重/视频)。"""
-
-from supply_infra.scheduler.jobs.demand_pool.sync import sync_multi_demand_pool_odps_to_mysql
-
-__all__ = ["sync_multi_demand_pool_odps_to_mysql"]
+"""Demand pool sync job stages used by the durable pipeline."""

+ 0 - 15
supply_infra/scheduler/jobs/demand_pool/__main__.py

@@ -1,15 +0,0 @@
-"""策略需求池同步任务 CLI:``python -m supply_infra.scheduler.jobs.demand_pool``。"""
-from __future__ import annotations
-
-import sys
-
-from supply_infra.scheduler.cli_result import run_cli
-from supply_infra.scheduler.jobs.demand_pool.sync import sync_multi_demand_pool_odps_to_mysql
-
-if __name__ == "__main__":
-    run_cli(
-        lambda: sync_multi_demand_pool_odps_to_mysql(
-            partition_date=sys.argv[1] if len(sys.argv) > 1 else None,
-        ),
-        label="sync_multi_demand_pool_odps_to_mysql",
-    )

+ 10 - 95
supply_infra/scheduler/jobs/demand_pool/sync.py

@@ -1,19 +1,13 @@
 """
-定时任务:从 ODPS 同步策略需求天级表到 MySQL,并对新词做归属分类与热度统计。
-
-流程:
-1. 比对当天 ODPS / MySQL 行数,相同则跳过写入
-2. 有差异时拉取 ODPS,按 (strategy, demand_id) 只插入缺失、删除多余,并回填已有行 video_list
-3. video_list 每条最多保留前 10 个 video_id,video_count 与之保持一致
-4. 拉取近 7 日 rov_diff/vov_diff,按特征值匹配回填 real_rov_7d / real_vov_7d
-5. 查询当天全部 demand_name,按空格分词写入 set
-6. 过滤 demand_belong_category 中已存在的词
-7. 剩余词按 100 词一批调用 demand_belong_category_agent
-8. 遍历 demand_belong_category 全部词,按策略与 rov_diff/vov_diff 写入 demand_popularity_stats(avg/count)
-9. 基于三表计算整棵类目树节点加权平均分,写入 category_tree_weight
-10. 四维热度全局排名归一化打分,写入各维 score 与 total_score
-11. 增量同步 multi_demand_video_detail(全表 video_list → ODPS 昨天分区最终选题)
-12. 补充 demand_belong_pool_rel 匹配边,并回填词级 video_list(最多 10 个)
+Demand pool ODPS → MySQL sync stages used by the durable pipeline.
+
+Stages owned by this module:
+1. source row sync (`_sync_pool_rows`)
+2. demand word classification (`_classify_words`)
+3. real rov/vov enrichment (`enrich_real_rov_vov_7d`)
+4. popularity stats (`compute_popularity_stats`)
+
+Related pipeline stages live in sibling modules (belong_rel / tree_weight / videos).
 """
 from __future__ import annotations
 
@@ -21,7 +15,7 @@ import json
 import logging
 from datetime import datetime, timedelta
 from decimal import Decimal
-from typing import Any, Callable
+from typing import Any
 
 from agents.demand_belong_category_agent.run import main as classify_demand_words
 from supply_infra.db.repositories.demand_belong_category_repo import (
@@ -33,11 +27,6 @@ from supply_infra.db.repositories.demand_popularity_stats_repo import (
 from supply_infra.db.repositories.multi_demand_pool_di_repo import MultiDemandPoolDiRepository
 from supply_infra.db.session import get_session
 from supply_infra.odps.client import get_odps_client
-from supply_infra.scheduler.jobs.demand_pool.belong_rel import sync_demand_belong_pool_rel
-from supply_infra.scheduler.jobs.demand_pool.tree_weight import (
-    compute_category_tree_weight,
-)
-from supply_infra.scheduler.jobs.demand_pool.videos import sync_multi_demand_videos
 
 logger = logging.getLogger(__name__)
 
@@ -502,77 +491,3 @@ def _sync_pool_rows(partition_date: str) -> dict[str, Any]:
         "mysql_count": mysql_count,
         **_sync_diff(partition_date),
     }
-
-
-def _run_sync_stage(
-    stage_name: str,
-    action: Callable[[], dict[str, Any]],
-) -> tuple[dict[str, Any], str | None]:
-    """执行需求池同步子阶段,错误转为结构化结果并允许后续阶段继续。"""
-    try:
-        payload = action()
-    except Exception as exc:
-        logger.exception("Multi demand pool stage failed; continue: stage=%s", stage_name)
-        return {"success": False, "error": str(exc)}, str(exc)
-
-    if payload.get("success") is False:
-        error = str(payload.get("error") or f"{stage_name} returned success=False")
-        logger.error(
-            "Multi demand pool stage reported failure; continue: stage=%s error=%s",
-            stage_name,
-            error,
-        )
-        return payload, error
-    return payload, None
-
-
-def sync_multi_demand_pool_odps_to_mysql(partition_date: str | None = None) -> dict:
-    """
-    从 ODPS 增量同步策略需求天级数据到 MySQL,并对新词做归属分类与热度统计。
-
-    Args:
-        partition_date: 分区日期 (YYYYMMDD),默认当天
-    """
-    if partition_date is None:
-        partition_date = datetime.now().strftime("%Y%m%d")
-
-    logger.info(
-        "Starting multi demand pool ODPS → MySQL sync for partition: %s",
-        partition_date,
-    )
-
-    stage_definitions: list[tuple[str, Callable[[], dict[str, Any]]]] = [
-        ("source_sync", lambda: _sync_pool_rows(partition_date)),
-        ("classify", lambda: _classify_words(partition_date)),
-        ("belong_pool_rel", sync_demand_belong_pool_rel),
-        ("real_metrics", lambda: enrich_real_rov_vov_7d(partition_date)),
-        ("popularity", lambda: compute_popularity_stats(partition_date)),
-        ("tree_weight", lambda: compute_category_tree_weight(partition_date)),
-        ("videos", lambda: sync_multi_demand_videos(limit=None)),
-    ]
-    stage_results: dict[str, dict[str, Any]] = {}
-    errors: list[dict[str, str]] = []
-    for stage_name, action in stage_definitions:
-        payload, error = _run_sync_stage(stage_name, action)
-        stage_results[stage_name] = payload
-        if error:
-            errors.append({"stage": stage_name, "error": error})
-
-    sync_stats = stage_results["source_sync"]
-
-    result = {
-        "success": not errors,
-        "partition_date": partition_date,
-        **({} if errors and sync_stats.get("success") is False else sync_stats),
-        "source_sync": sync_stats,
-        "real_metrics": stage_results["real_metrics"],
-        "classify": stage_results["classify"],
-        "belong_pool_rel": stage_results["belong_pool_rel"],
-        "popularity": stage_results["popularity"],
-        "tree_weight": stage_results["tree_weight"],
-        "videos": stage_results["videos"],
-        "errors": errors,
-        "synced_at": datetime.now().isoformat(),
-    }
-    logger.info("Multi demand pool sync completed: %s", result)
-    return result

+ 0 - 66
supply_infra/scheduler/manual_jobs.py

@@ -1,66 +0,0 @@
-"""Deprecated compatibility facade for durable pipeline submissions."""
-from __future__ import annotations
-
-from dataclasses import dataclass
-from typing import Any
-
-from supply_infra.pipeline.run_service import get_pipeline_run, submit_pipeline_run
-from supply_infra.scheduler.constants import SUPPLY_PIPELINE_JOB_ID, SUPPLY_PIPELINE_JOB_NAME
-
-
-@dataclass(frozen=True)
-class ManualJobSpec:
-    job_id: str
-    name: str
-    description: str
-    accepts_biz_dt: bool = True
-    accepts_partition_date: bool = False
-
-
-MANUAL_JOBS = {
-    SUPPLY_PIPELINE_JOB_ID: ManualJobSpec(
-        job_id=SUPPLY_PIPELINE_JOB_ID,
-        name=SUPPLY_PIPELINE_JOB_NAME,
-        description="提交严格门禁的持久化供给流水线",
-    )
-}
-
-
-def list_manual_jobs() -> list[dict[str, Any]]:
-    return [
-        {
-            "id": item.job_id,
-            "name": item.name,
-            "description": item.description,
-            "accepts_biz_dt": item.accepts_biz_dt,
-            "accepts_partition_date": item.accepts_partition_date,
-        }
-        for item in MANUAL_JOBS.values()
-    ]
-
-
-def get_manual_job(job_id: str) -> ManualJobSpec | None:
-    return MANUAL_JOBS.get(job_id)
-
-
-def get_manual_run(run_id: str) -> dict[str, Any] | None:
-    return get_pipeline_run(run_id)
-
-
-def trigger_manual_job(
-    job_id: str,
-    *,
-    biz_dt: str | None = None,
-    partition_date: str | None = None,
-    wait: bool = False,
-) -> dict[str, Any]:
-    del partition_date
-    if job_id != SUPPLY_PIPELINE_JOB_ID:
-        raise KeyError(job_id)
-    if wait:
-        raise ValueError("wait=true is disabled; poll the persistent run_id")
-    return submit_pipeline_run(
-        biz_dt=biz_dt,
-        trigger_type="api",
-        trigger_source="legacy_manual_job",
-    ).to_dict()

+ 0 - 8
supply_infra/scheduler/step_runner.py

@@ -1,8 +0,0 @@
-"""Deprecated import compatibility for the durable pipeline step runner."""
-
-from supply_infra.pipeline.step_runner import (
-    StepExecutionRequest,
-    run_step_subprocess,
-)
-
-__all__ = ["StepExecutionRequest", "run_step_subprocess"]

+ 1 - 50
tests/supply_infra/scheduler/test_sync_multi_demand_pool.py

@@ -3,56 +3,7 @@ from __future__ import annotations
 
 from unittest.mock import MagicMock, patch
 
-from supply_infra.scheduler.jobs.demand_pool.sync import (
-    _classify_words,
-    sync_multi_demand_pool_odps_to_mysql,
-)
-
-
-@patch(
-    "supply_infra.scheduler.jobs.demand_pool.sync.sync_multi_demand_videos"
-)
-@patch(
-    "supply_infra.scheduler.jobs.demand_pool.sync.compute_category_tree_weight"
-)
-@patch(
-    "supply_infra.scheduler.jobs.demand_pool.sync.compute_popularity_stats"
-)
-@patch(
-    "supply_infra.scheduler.jobs.demand_pool.sync.enrich_real_rov_vov_7d"
-)
-@patch(
-    "supply_infra.scheduler.jobs.demand_pool.sync.sync_demand_belong_pool_rel"
-)
-@patch("supply_infra.scheduler.jobs.demand_pool.sync._classify_words")
-@patch("supply_infra.scheduler.jobs.demand_pool.sync._sync_pool_rows")
-def test_substep_failure_does_not_skip_later_substeps(
-    mock_source_sync,
-    mock_classify,
-    mock_rel,
-    mock_real_metrics,
-    mock_popularity,
-    mock_tree_weight,
-    mock_videos,
-) -> None:
-    mock_source_sync.return_value = {"inserted": 10}
-    mock_classify.side_effect = RuntimeError("agent unavailable")
-    mock_rel.return_value = {"inserted": 2}
-    mock_real_metrics.return_value = {"updated_rows": 3}
-    mock_popularity.return_value = {"upserted": 4}
-    mock_tree_weight.return_value = {"upserted": 5}
-    mock_videos.return_value = {"inserted": 6}
-
-    result = sync_multi_demand_pool_odps_to_mysql("20260721")
-
-    assert result["success"] is False
-    assert result["classify"] == {"success": False, "error": "agent unavailable"}
-    assert result["belong_pool_rel"] == {"inserted": 2}
-    assert result["real_metrics"] == {"updated_rows": 3}
-    assert result["popularity"] == {"upserted": 4}
-    assert result["tree_weight"] == {"upserted": 5}
-    assert result["videos"] == {"inserted": 6}
-    assert [item["stage"] for item in result["errors"]] == ["classify"]
+from supply_infra.scheduler.jobs.demand_pool.sync import _classify_words
 
 
 @patch(