Преглед изворни кода

feat(pipeline): add cloud pipeline runners

SamLee пре 3 недеља
родитељ
комит
e5ac1a9722

+ 14 - 0
pipeline/__init__.py

@@ -0,0 +1,14 @@
+"""Formal creation pipeline orchestration package."""
+
+from pipeline.creation_pipeline import CreationPipelineResult, run_creation_pipeline
+from pipeline.dedupe import dedupe_candidate_items, item_dedupe_key
+from pipeline.models import PipelineJob, PipelineRun
+
+__all__ = [
+    "CreationPipelineResult",
+    "PipelineJob",
+    "PipelineRun",
+    "dedupe_candidate_items",
+    "item_dedupe_key",
+    "run_creation_pipeline",
+]

+ 51 - 0
pipeline/acquisition_runner.py

@@ -0,0 +1,51 @@
+"""Pipeline adapter for the acquisition stage."""
+from __future__ import annotations
+
+from dataclasses import dataclass
+from typing import Any
+from uuid import UUID
+
+from acquisition.repositories.base import AcquisitionRepository
+from acquisition.runner import RunBatchResult, run_batch
+from core.config import Settings
+from pipeline.models import PipelineJob, PipelineRun
+from pipeline.repository import PipelineRepository
+
+
+@dataclass(frozen=True)
+class AcquisitionStageResult:
+    pipeline_job: PipelineJob | None
+    acquisition: RunBatchResult
+
+
+def run_acquisition_stage(
+    *,
+    acquisition_repo: AcquisitionRepository,
+    batch_id: UUID,
+    settings: Settings,
+    pipeline_repo: PipelineRepository | None = None,
+    pipeline_run: PipelineRun | None = None,
+    **kwargs: Any,
+) -> AcquisitionStageResult:
+    job = None
+    if pipeline_repo and pipeline_run and pipeline_run.id:
+        job = pipeline_repo.save_pipeline_job(
+            run_id=pipeline_run.id,
+            stage="search",
+            target_id=batch_id,
+            status="running",
+            metadata={"batch_id": str(batch_id)},
+        )
+    try:
+        result = run_batch(acquisition_repo, batch_id=batch_id, settings=settings, **kwargs)
+        if pipeline_repo and job and job.id:
+            job = pipeline_repo.mark_job_status(
+                job.id,
+                status="done" if result.failed == 0 else "partial",
+                metadata=result.__dict__,
+            )
+        return AcquisitionStageResult(pipeline_job=job, acquisition=result)
+    except Exception as exc:
+        if pipeline_repo and job and job.id:
+            pipeline_repo.mark_job_status(job.id, status="failed", error_message=str(exc)[:300])
+        raise

+ 66 - 0
pipeline/creation_pipeline.py

@@ -0,0 +1,66 @@
+"""Formal query -> acquisition -> decode orchestration."""
+from __future__ import annotations
+
+from dataclasses import dataclass
+from uuid import UUID
+
+from acquisition.repositories.base import AcquisitionRepository
+from core.config import Settings
+from decode_content.service import DecodeService
+from pipeline.acquisition_runner import AcquisitionStageResult, run_acquisition_stage
+from pipeline.decode_runner import DecodeBatchResult, run_decode_stage
+from pipeline.models import PipelineRun, ResumePolicy
+from pipeline.repository import PipelineRepository
+
+
+@dataclass(frozen=True)
+class CreationPipelineResult:
+    pipeline_run: PipelineRun | None
+    acquisition: AcquisitionStageResult
+    decode: DecodeBatchResult | None = None
+
+
+def run_creation_pipeline(
+    *,
+    acquisition_repo: AcquisitionRepository,
+    batch_id: UUID,
+    settings: Settings,
+    pipeline_repo: PipelineRepository | None = None,
+    decode_service: DecodeService | None = None,
+    run_key: str | None = None,
+    resume_policy: ResumePolicy | None = None,
+    decode_limit: int = 100,
+) -> CreationPipelineResult:
+    resume_policy = resume_policy or ResumePolicy()
+    pipeline_run = (
+        pipeline_repo.create_pipeline_run(
+            run_key=run_key or f"creation-pipeline:{batch_id}",
+            batch_id=batch_id,
+            status="running",
+            metadata=resume_policy.model_dump(),
+        )
+        if pipeline_repo
+        else None
+    )
+    acquisition_result = run_acquisition_stage(
+        acquisition_repo=acquisition_repo,
+        batch_id=batch_id,
+        settings=settings,
+        pipeline_repo=pipeline_repo,
+        pipeline_run=pipeline_run,
+        resume=resume_policy.resume,
+        skip_done=resume_policy.skip_done,
+    )
+    decode_result = None
+    if decode_service is not None:
+        decode_result = run_decode_stage(
+            candidate_repo=acquisition_repo,
+            decode_service=decode_service,
+            run_id=acquisition_result.acquisition.run_id,
+            limit=decode_limit,
+        )
+    return CreationPipelineResult(
+        pipeline_run=pipeline_run,
+        acquisition=acquisition_result,
+        decode=decode_result,
+    )

+ 68 - 0
pipeline/decode_runner.py

@@ -0,0 +1,68 @@
+"""Pipeline adapter for decoding creation candidate items."""
+from __future__ import annotations
+
+from dataclasses import dataclass
+from typing import Protocol
+from uuid import UUID
+
+from acquisition.domain import CandidateItem, MediaAsset
+from decode_content.readers.service import post_from_candidate_item
+from decode_content.service import DecodeService, DecodeWorkflowOutput
+from pipeline.dedupe import dedupe_candidate_items, should_decode_item
+
+
+class DecodeCandidateRepository(Protocol):
+    def list_creation_candidate_items(
+        self,
+        *,
+        run_id: UUID | None = None,
+        limit: int = 100,
+    ) -> list[CandidateItem]:
+        ...
+
+    def list_media_assets_for_item(self, item_id: UUID) -> list[MediaAsset]:
+        ...
+
+
+@dataclass(frozen=True)
+class DecodeBatchResult:
+    total: int
+    decoded: int
+    skipped: int
+    failed: int
+    outputs: list[DecodeWorkflowOutput]
+
+
+def run_decode_stage(
+    *,
+    candidate_repo: DecodeCandidateRepository,
+    decode_service: DecodeService,
+    run_id: UUID | None = None,
+    limit: int = 100,
+    decoded_item_ids: set[str] | None = None,
+) -> DecodeBatchResult:
+    items = dedupe_candidate_items(candidate_repo.list_creation_candidate_items(run_id=run_id, limit=limit))
+    outputs: list[DecodeWorkflowOutput] = []
+    decoded = skipped = failed = 0
+    for item in items:
+        if item.id is None:
+            skipped += 1
+            continue
+        decision = should_decode_item(item, decoded_item_ids=decoded_item_ids)
+        if not decision.keep:
+            skipped += 1
+            continue
+        try:
+            media = candidate_repo.list_media_assets_for_item(item.id)
+            post = post_from_candidate_item(item, media)
+            outputs.append(decode_service.decode_post(item_id=item.id, post=post))
+            decoded += 1
+        except Exception:
+            failed += 1
+    return DecodeBatchResult(
+        total=len(items),
+        decoded=decoded,
+        skipped=skipped,
+        failed=failed,
+        outputs=outputs,
+    )

+ 46 - 0
pipeline/dedupe.py

@@ -0,0 +1,46 @@
+"""Deduplication helpers for acquisition -> decode -> ingest pipelines."""
+from __future__ import annotations
+
+from dataclasses import dataclass
+from typing import Iterable
+
+from acquisition.domain import CandidateItem
+
+
+@dataclass(frozen=True)
+class DedupeDecision:
+    keep: bool
+    key: str
+    reason: str = ""
+
+
+def item_dedupe_key(item: CandidateItem) -> str:
+    if item.platform_item_id:
+        return f"{item.platform}:id:{item.platform_item_id}"
+    if item.canonical_url:
+        return f"url:{item.canonical_url.strip().lower()}"
+    return f"item:{item.id}"
+
+
+def dedupe_candidate_items(items: Iterable[CandidateItem]) -> list[CandidateItem]:
+    seen: set[str] = set()
+    out: list[CandidateItem] = []
+    for item in items:
+        key = item_dedupe_key(item)
+        if key in seen:
+            continue
+        seen.add(key)
+        out.append(item)
+    return out
+
+
+def should_decode_item(
+    item: CandidateItem,
+    *,
+    decoded_item_ids: set[str] | None = None,
+) -> DedupeDecision:
+    key = item_dedupe_key(item)
+    decoded = decoded_item_ids or set()
+    if str(item.id) in decoded or key in decoded:
+        return DedupeDecision(keep=False, key=key, reason="already_decoded")
+    return DedupeDecision(keep=True, key=key)

+ 58 - 0
pipeline/models.py

@@ -0,0 +1,58 @@
+"""Pipeline run and resume models for formal orchestration."""
+from __future__ import annotations
+
+from datetime import datetime
+from typing import Any, Literal
+from uuid import UUID
+
+from pydantic import BaseModel, ConfigDict, Field
+
+
+PipelineStage = Literal["query", "search", "classify", "decode", "payload", "ingest"]
+PipelineStatus = Literal["pending", "running", "done", "failed", "skipped", "partial"]
+PipelineStep = PipelineStage
+RunStatus = PipelineStatus
+JobStatus = PipelineStatus
+
+
+class PipelineModel(BaseModel):
+    model_config = ConfigDict(extra="forbid")
+
+
+class ResumePolicy(PipelineModel):
+    resume: bool = True
+    skip_done: bool = True
+    retry_failed: bool = False
+    max_attempts: int = 3
+
+
+class ResumeCursor(PipelineModel):
+    stage: PipelineStage
+    target_id: UUID | None = None
+    last_successful_job_id: UUID | None = None
+    metadata: dict[str, Any] = Field(default_factory=dict)
+
+
+class PipelineRun(PipelineModel):
+    id: UUID | None = None
+    run_key: str | None = None
+    batch_id: UUID | None = None
+    status: PipelineStatus = "pending"
+    current_stage: PipelineStage | None = None
+    metadata: dict[str, Any] = Field(default_factory=dict)
+    error_message: str | None = None
+    started_at: datetime | None = None
+    finished_at: datetime | None = None
+
+
+class PipelineJob(PipelineModel):
+    id: UUID | None = None
+    run_id: UUID | None = None
+    stage: PipelineStage
+    status: PipelineStatus = "pending"
+    target_id: UUID | None = None
+    attempt_count: int = 0
+    metadata: dict[str, Any] = Field(default_factory=dict)
+    error_message: str | None = None
+    started_at: datetime | None = None
+    finished_at: datetime | None = None

+ 52 - 0
pipeline/repository.py

@@ -0,0 +1,52 @@
+"""Repository contracts for formal pipeline state."""
+from __future__ import annotations
+
+from typing import Any, Protocol
+from uuid import UUID
+
+from pipeline.models import PipelineJob, PipelineRun, PipelineStage, ResumeCursor
+
+
+class PipelineRepository(Protocol):
+    """Persistence boundary for orchestration and resumable jobs.
+
+    The first migration does not create dedicated pipeline tables yet, so this
+    interface captures the intended boundary before a persistence decision.
+    """
+
+    def create_pipeline_run(
+        self,
+        *,
+        run_key: str | None = None,
+        batch_id: UUID | None = None,
+        status: str = "pending",
+        metadata: dict[str, Any] | None = None,
+    ) -> PipelineRun:
+        ...
+
+    def get_pipeline_run(self, run_id: UUID) -> PipelineRun:
+        ...
+
+    def save_pipeline_job(
+        self,
+        *,
+        run_id: UUID,
+        stage: PipelineStage,
+        target_id: UUID | None = None,
+        status: str = "pending",
+        metadata: dict[str, Any] | None = None,
+    ) -> PipelineJob:
+        ...
+
+    def mark_job_status(
+        self,
+        job_id: UUID,
+        *,
+        status: str,
+        error_message: str | None = None,
+        metadata: dict[str, Any] | None = None,
+    ) -> PipelineJob:
+        ...
+
+    def get_resume_cursor(self, run_id: UUID) -> ResumeCursor | None:
+        ...