xueyiming пре 6 часа
родитељ
комит
0e75bfb5e6

+ 2 - 0
.env.example

@@ -62,6 +62,8 @@ SCHEDULER_ENABLED=false
 SCHEDULER_TIMEZONE=Asia/Shanghai
 SCHEDULER_CRON_HOUR=15
 SCHEDULER_CRON_MINUTE=0
+# AIGC 候选分发独立轮询间隔(分钟);与日批流水线解耦
+SCHEDULER_AIGC_PUBLISH_INTERVAL_MINUTES=10
 
 # find_agent P0 质量门禁(每次运行会保存完整配置快照)
 FIND_AGENT_RULE_VERSION=find-agent-gate-v1

+ 42 - 11
api/services/scheduler.py

@@ -8,19 +8,39 @@ 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 SUPPLY_PIPELINE_JOB_ID, SUPPLY_PIPELINE_JOB_NAME
+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": "提交严格门禁的 12 步持久化供给流水线",
+            "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": [],
+        },
     ]
 
 
@@ -29,13 +49,19 @@ def run_scheduler_job(
     *,
     biz_dt: str | None = None,
 ) -> dict[str, Any]:
-    if job_id != SUPPLY_PIPELINE_JOB_ID:
-        raise KeyError(job_id)
-    return submit_pipeline_run(
-        biz_dt=biz_dt,
-        trigger_type="api",
-        trigger_source="legacy_scheduler_api",
-    ).to_dict()
+    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:
@@ -62,7 +88,12 @@ def scheduler_status() -> dict[str, Any]:
                 "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,
     }

+ 6 - 0
supply_infra/config.py

@@ -108,6 +108,12 @@ class InfraSettings(BaseSettings):
         le=59,
         alias="SCHEDULER_CRON_MINUTE",
     )
+    scheduler_aigc_publish_interval_minutes: int = Field(
+        default=10,
+        ge=1,
+        le=1440,
+        alias="SCHEDULER_AIGC_PUBLISH_INTERVAL_MINUTES",
+    )
 
     # find_agent quality gates
     find_agent_rule_version: str = Field(

+ 0 - 1
supply_infra/pipeline/dag.py

@@ -26,7 +26,6 @@ _RAW_STEPS: tuple[tuple[str, int, bool], ...] = (
     ("demand_grade", 21600, False),
     ("demand_expand", 21600, False),
     ("video_discovery", 21600, False),
-    ("aigc_write_record", 1800, True),
 )
 
 

+ 0 - 15
supply_infra/pipeline/gates.py

@@ -31,19 +31,4 @@ def evaluate_step_gate(step_key: str, payload: dict[str, Any]) -> GateDecision:
                 f"{step_key} reported {failed} failed item(s)",
             )
 
-    if step_key == "aigc_write_record":
-        if payload.get("effect_recorded") is not True or not payload.get("payload_hash"):
-            return GateDecision(
-                False,
-                "aigc_record_incomplete",
-                "AIGC publish payload/hash was not durably recorded",
-            )
-        failed_batches = int(payload.get("failed_batch_count", 0) or 0)
-        if failed_batches > 0:
-            return GateDecision(
-                False,
-                "aigc_publish_failed",
-                f"AIGC publish reported {failed_batches} failed batch(es)",
-            )
-
     return GateDecision(True)

+ 0 - 77
supply_infra/pipeline/registry.py

@@ -1,19 +1,11 @@
 from __future__ import annotations
 
-import hashlib
-import json
 from collections.abc import Callable
-from pathlib import Path
 from typing import Any
 
-from supply_infra.config import get_infra_settings
-from supply_infra.db.repositories.pipeline_outbox_repo import PipelineOutboxRepository
-from supply_infra.db.session import get_session
 from supply_infra.pipeline.contracts import StepContext
 
 StepHandler = Callable[[StepContext], dict[str, Any]]
-_REPO_ROOT = Path(__file__).resolve().parents[2]
-_MAX_INLINE_EFFECT_BYTES = 512_000
 
 
 def _global_tree(context: StepContext) -> dict[str, Any]:
@@ -104,74 +96,6 @@ def _discover(context: StepContext) -> dict[str, Any]:
     )
 
 
-def _aigc_write_record(context: StepContext) -> dict[str, Any]:
-    from supply_infra.scheduler.jobs.publish_videos_from_discovery import (
-        publish_videos_from_discovery,
-    )
-
-    settings = get_infra_settings()
-    payload = publish_videos_from_discovery(biz_dt=context.biz_dt)
-    publish_success = bool(payload.get("success", True))
-    canonical = json.dumps(
-        payload,
-        ensure_ascii=False,
-        sort_keys=True,
-        separators=(",", ":"),
-        default=str,
-    )
-    encoded = canonical.encode("utf-8")
-    payload_hash = hashlib.sha256(encoded).hexdigest()
-    idempotency_key = f"aigc_write_record:{context.biz_dt}:{payload_hash}"
-    payload_uri: str | None = None
-    stored_payload: dict[str, Any] | None = payload
-    if len(encoded) > _MAX_INLINE_EFFECT_BYTES:
-        log_dir = Path(settings.pipeline_log_dir)
-        if not log_dir.is_absolute():
-            log_dir = _REPO_ROOT / log_dir
-        effect_dir = log_dir / context.run_id / "effects"
-        effect_dir.mkdir(parents=True, exist_ok=True)
-        effect_path = effect_dir / f"{context.step_run_id}-{payload_hash}.json"
-        temporary_path = effect_path.with_suffix(".json.tmp")
-        temporary_path.write_bytes(encoded)
-        temporary_path.replace(effect_path)
-        payload_uri = str(effect_path)
-        stored_payload = {
-            "externalized": True,
-            "payload_bytes": len(encoded),
-            "payload_hash": payload_hash,
-        }
-    with get_session() as session:
-        record = PipelineOutboxRepository(session).record_effect(
-            run_id=context.run_id,
-            step_run_id=context.step_run_id,
-            effect_type="aigc_write_record",
-            idempotency_key=idempotency_key,
-            payload_hash=payload_hash,
-            payload=stored_payload,
-            payload_uri=payload_uri,
-        )
-        outbox_id = record.outbox_id
-
-    result: dict[str, Any] = {
-        "success": publish_success,
-        "effect_recorded": True,
-        "external_request_made": int(payload.get("candidate_count", 0) or 0) > 0,
-        "outbox_id": outbox_id,
-        "payload_hash": payload_hash,
-        "payload_uri": payload_uri,
-        "candidate_count": payload.get("candidate_count", 0),
-        "batch_count": payload.get("batch_count", 0),
-        "failed_batch_count": payload.get("failed_batch_count", 0),
-    }
-    if not publish_success:
-        result["error"] = str(
-            payload.get("error")
-            or f"AIGC publish failed ({result['failed_batch_count']} batch(es))"
-        )
-        result["error_code"] = "aigc_publish_failed"
-    return result
-
-
 STEP_REGISTRY: dict[str, StepHandler] = {
     "global_tree_sync": _global_tree,
     "demand_pool_source_sync": _demand_source,
@@ -184,7 +108,6 @@ STEP_REGISTRY: dict[str, StepHandler] = {
     "demand_grade": _grade,
     "demand_expand": _expand,
     "video_discovery": _discover,
-    "aigc_write_record": _aigc_write_record,
 }
 
 

+ 44 - 3
supply_infra/scheduler/app.py

@@ -5,10 +5,16 @@ 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 SUPPLY_PIPELINE_JOB_ID, SUPPLY_PIPELINE_JOB_NAME
+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
@@ -26,8 +32,28 @@ def submit_scheduled_pipeline() -> dict:
     ).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 that only submits durable pipeline batches."""
+    """Create a Scheduler for the daily pipeline and AIGC publish polling."""
     settings = get_infra_settings()
     tz = settings.scheduler_timezone
     scheduler = BackgroundScheduler(timezone=tz)
@@ -46,13 +72,28 @@ def create_scheduler() -> BackgroundScheduler:
         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",
+        "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
 

+ 4 - 0
supply_infra/scheduler/constants.py

@@ -3,5 +3,9 @@
 SUPPLY_PIPELINE_JOB_ID = "run_supply_pipeline"
 SUPPLY_PIPELINE_JOB_NAME = "供给数据流水线"
 
+# AIGC:独立轮询,发现未发布的 primary 候选后分发
+AIGC_PUBLISH_JOB_ID = "publish_videos_from_discovery"
+AIGC_PUBLISH_JOB_NAME = "AIGC候选分发"
+
 # find_agent:当日全部 S/A 需求(有拓展点位),单线程串行找视频
 PIPELINE_FIND_AGENT_WORKERS = 1

+ 2 - 2
supply_infra/scheduler/jobs/discover_videos_from_demands.py

@@ -1,5 +1,5 @@
 """
-从全部 S/A 级需求及其拓展点位触发 find_agent 视频发现;有效视频满 100 提前结束。
+从全部 S/A 级需求及其拓展点位触发 find_agent 视频发现;有效视频满 300 提前结束。
 
 任务层负责查库与组装上下文;Agent 负责搜索、画像与分池落库。
 """
@@ -28,7 +28,7 @@ from supply_infra.services.video_discovery_service import get_video_discovery_se
 logger = logging.getLogger(__name__)
 
 _DEFAULT_WORKERS = 1
-_DAILY_PASSED_VIDEO_LIMIT = 100
+_DAILY_PASSED_VIDEO_LIMIT = 300
 
 
 def _resolve_biz_dt(biz_dt: str | None) -> str:

+ 5 - 1
supply_infra/scheduler/jobs/publish_videos_from_discovery.py

@@ -87,14 +87,16 @@ def publish_videos_from_discovery(
     run_id: str | None = None,
     skip_published: bool = True,
     limit: int | None = None,
+    any_biz_dt: bool = False,
 ) -> dict[str, Any]:
     """
     从 video_discovery_candidate 读取视频,按轮询均匀分配到各 AIGC 计划对。
 
     仅处理 decision_bucket 为 primary 且 aweme_id 非空的候选,不区分品类。
+    any_biz_dt=True 时不按业务日过滤,用于定时轮询扫全部未发布候选。
     """
     settings = get_infra_settings()
-    resolved_biz_dt = _resolve_biz_dt(biz_dt)
+    resolved_biz_dt = None if any_biz_dt else _resolve_biz_dt(biz_dt)
     plan_pairs = list_unique_plan_pairs()
     if not plan_pairs:
         raise RuntimeError("AIGC 计划映射为空,无法分发")
@@ -111,6 +113,7 @@ def publish_videos_from_discovery(
             "success": True,
             "message": "没有待发布的候选视频",
             "biz_dt": resolved_biz_dt,
+            "any_biz_dt": any_biz_dt,
             "run_id": run_id,
             "plan_count": len(plan_pairs),
             "candidate_count": 0,
@@ -195,6 +198,7 @@ def publish_videos_from_discovery(
     return {
         "success": not failed_batches,
         "biz_dt": resolved_biz_dt,
+        "any_biz_dt": any_biz_dt,
         "run_id": run_id,
         "plan_count": len(plan_pairs),
         "candidate_count": len(candidates),

+ 12 - 31
tests/supply_infra/pipeline/test_dag_and_gates.py

@@ -4,48 +4,29 @@ from supply_infra.pipeline.dag import PIPELINE_STEPS
 from supply_infra.pipeline.gates import evaluate_step_gate
 
 
-def test_pipeline_has_twelve_strictly_ordered_steps() -> None:
-    assert len(PIPELINE_STEPS) == 12
+def test_pipeline_has_eleven_strictly_ordered_steps() -> None:
+    assert len(PIPELINE_STEPS) == 11
     assert PIPELINE_STEPS[0].key == "global_tree_sync"
-    assert PIPELINE_STEPS[-1].key == "aigc_write_record"
+    assert PIPELINE_STEPS[-1].key == "video_discovery"
+    assert "aigc_write_record" not in {step.key for step in PIPELINE_STEPS}
     for previous, current in zip(PIPELINE_STEPS, PIPELINE_STEPS[1:]):
         assert current.dependencies == (previous.key,)
         assert current.critical is True
 
 
-def test_aigc_gate_requires_durable_effect_record() -> None:
-    failed = evaluate_step_gate("aigc_write_record", {"success": True})
-    assert failed.passed is False
-    live_ok = evaluate_step_gate(
-        "aigc_write_record",
-        {
-            "success": True,
-            "effect_recorded": True,
-            "payload_hash": "abc",
-            "external_request_made": True,
-        },
-    )
-    assert live_ok.passed is True
-
-
-def test_aigc_gate_rejects_failed_batches() -> None:
+def test_explicit_step_failure_fails_closed() -> None:
     decision = evaluate_step_gate(
-        "aigc_write_record",
-        {
-            "success": True,
-            "effect_recorded": True,
-            "payload_hash": "abc",
-            "failed_batch_count": 2,
-        },
+        "global_tree_sync",
+        {"success": False, "error": "source missing"},
     )
     assert decision.passed is False
-    assert decision.error_code == "aigc_publish_failed"
+    assert decision.error_code == "step_reported_failure"
 
 
-def test_explicit_step_failure_fails_closed() -> None:
+def test_video_discovery_gate_rejects_business_failures() -> None:
     decision = evaluate_step_gate(
-        "global_tree_sync",
-        {"success": False, "error": "source missing"},
+        "video_discovery",
+        {"success": True, "failed": 2},
     )
     assert decision.passed is False
-    assert decision.error_code == "step_reported_failure"
+    assert decision.error_code == "business_failures"

+ 36 - 0
tests/supply_infra/scheduler/test_aigc_publish_scheduler.py

@@ -0,0 +1,36 @@
+from __future__ import annotations
+
+from unittest.mock import patch
+
+from supply_infra.scheduler.app import create_scheduler, submit_aigc_publish
+from supply_infra.scheduler.constants import (
+    AIGC_PUBLISH_JOB_ID,
+    SUPPLY_PIPELINE_JOB_ID,
+)
+
+
+def test_create_scheduler_registers_pipeline_and_aigc_jobs() -> None:
+    scheduler = create_scheduler()
+    job_ids = {job.id for job in scheduler.get_jobs()}
+    assert SUPPLY_PIPELINE_JOB_ID in job_ids
+    assert AIGC_PUBLISH_JOB_ID in job_ids
+    aigc_job = scheduler.get_job(AIGC_PUBLISH_JOB_ID)
+    assert aigc_job is not None
+    assert aigc_job.trigger.interval.total_seconds() == 600
+
+
+@patch(
+    "supply_infra.scheduler.jobs.publish_videos_from_discovery.publish_videos_from_discovery"
+)
+def test_submit_aigc_publish_scans_all_unpublished(mock_publish) -> None:
+    mock_publish.return_value = {
+        "success": True,
+        "candidate_count": 0,
+        "batch_count": 0,
+        "failed_batch_count": 0,
+    }
+
+    result = submit_aigc_publish()
+
+    assert result["success"] is True
+    mock_publish.assert_called_once_with(skip_published=True, any_biz_dt=True)

+ 3 - 3
tests/supply_infra/scheduler/test_discover_videos_from_demands.py

@@ -838,7 +838,7 @@ def test_evaluate_find_agent_run_fails_without_candidates(
 @patch(
     "supply_infra.scheduler.jobs.discover_videos_from_demands.list_find_demand_contexts"
 )
-def test_stops_discovery_after_100_passed_videos(
+def test_stops_discovery_after_300_passed_videos(
     mock_list_contexts,
     mock_filter_contexts,
     mock_count_passed,
@@ -863,14 +863,14 @@ def test_stops_discovery_after_100_passed_videos(
         contexts,
         {"total_loaded": 2, "skipped_already_done": 0},
     )
-    mock_count_passed.side_effect = [99, 100]
+    mock_count_passed.side_effect = [299, 300]
     mock_process.return_value = {"success": True, "skipped": False}
 
     result = discover_videos_from_demands("20260727", workers=1)
 
     assert mock_process.call_count == 1
     assert result["processed"] == 1
-    assert result["passed_videos"] == 100
+    assert result["passed_videos"] == 300
     assert result["stopped_by_passed_video_limit"] is True
 
 

+ 38 - 0
tests/supply_infra/scheduler/test_publish_videos_from_discovery.py

@@ -59,3 +59,41 @@ def test_publish_marks_candidates_when_bind_fails(
     )
     assert call_kwargs["crawler_plan_id"] == "crawler-1"
     assert call_kwargs["plan_label"].startswith("BIND_FAILED:")
+
+
+@patch(
+    "supply_infra.scheduler.jobs.publish_videos_from_discovery.get_video_discovery_service"
+)
+@patch(
+    "supply_infra.scheduler.jobs.publish_videos_from_discovery.list_unique_plan_pairs"
+)
+def test_publish_any_biz_dt_skips_date_filter(
+    mock_plan_pairs,
+    mock_get_service,
+) -> None:
+    mock_plan_pairs.return_value = [
+        type(
+            "Plan",
+            (),
+            {
+                "label": "plan-a",
+                "produce_plan_id": "produce-1",
+                "publish_plan_id": "publish-1",
+            },
+        )()
+    ]
+    mock_get_service.return_value = MagicMock(
+        list_publishable_candidates=MagicMock(return_value=[]),
+    )
+
+    result = publish_videos_from_discovery(any_biz_dt=True)
+
+    assert result["success"] is True
+    assert result["any_biz_dt"] is True
+    assert result["biz_dt"] is None
+    mock_get_service.return_value.list_publishable_candidates.assert_called_once_with(
+        run_id=None,
+        biz_dt=None,
+        skip_published=True,
+        limit=None,
+    )