| """Job-related domain models — used by services, storage, and API.""" |
|
|
| from __future__ import annotations |
|
|
| import enum |
| import uuid |
| from datetime import datetime, timezone |
| from typing import Any, List, Optional |
|
|
| from pydantic import BaseModel, Field |
|
|
|
|
| class JobKind(str, enum.Enum): |
| DETECTION = "detection" |
| RECOGNITION = "recognition" |
| SEARCH = "search" |
| IMAGE_ANALYSIS = "image_analysis" |
| METADATA = "metadata" |
| FORENSICS = "forensics" |
| OCR = "ocr" |
| OBJECT_DETECTION = "object_detection" |
| SCENE_RECOGNITION = "scene_recognition" |
| NSFW_DETECTION = "nsfw_detection" |
| AI_IMAGE_DETECTION = "ai_image_detection" |
| EMBEDDING = "embedding" |
| FULL_PIPELINE = "full_pipeline" |
|
|
|
|
| class JobStatus(str, enum.Enum): |
| PENDING = "pending" |
| QUEUED = "queued" |
| RUNNING = "running" |
| COMPLETED = "completed" |
| FAILED = "failed" |
| CANCELLED = "cancelled" |
| TIMEOUT = "timeout" |
|
|
|
|
| class JobRequest(BaseModel): |
| """Inbound request to create a job.""" |
| kind: Optional[JobKind] = None |
| image_url: Optional[str] = None |
| image_base64: Optional[str] = None |
| providers: List[str] = Field( |
| default_factory=list, |
| description="Optional whitelist of provider names. Empty = use all enabled.", |
| ) |
| options: dict = Field(default_factory=dict) |
|
|
| def has_image_input(self) -> bool: |
| return bool(self.image_url or self.image_base64) |
|
|
|
|
| class Job(BaseModel): |
| """Persisted job record.""" |
| id: str = Field(default_factory=lambda: str(uuid.uuid4())) |
| kind: JobKind |
| status: JobStatus = JobStatus.PENDING |
| created_at: datetime = Field(default_factory=lambda: datetime.now(timezone.utc)) |
| started_at: Optional[datetime] = None |
| completed_at: Optional[datetime] = None |
| request: JobRequest |
| image_hash: Optional[str] = None |
| error: Optional[str] = None |
|
|
| def mark_running(self) -> None: |
| self.status = JobStatus.RUNNING |
| self.started_at = datetime.now(timezone.utc) |
|
|
| def mark_completed(self) -> None: |
| self.status = JobStatus.COMPLETED |
| self.completed_at = datetime.now(timezone.utc) |
|
|
| def mark_failed(self, error: str) -> None: |
| self.status = JobStatus.FAILED |
| self.completed_at = datetime.now(timezone.utc) |
| self.error = error |
|
|
| def mark_timeout(self) -> None: |
| self.status = JobStatus.TIMEOUT |
| self.completed_at = datetime.now(timezone.utc) |
| self.error = "Job exceeded timeout" |
|
|
|
|
| class JobResult(BaseModel): |
| """The output of a completed job — wraps a UnifiedFaceReport.""" |
| job_id: str |
| status: JobStatus |
| report: Optional[Any] = None |
| error: Optional[str] = None |
| elapsed_ms: float = 0.0 |
|
|