"""Job models — the shape of one ingestion batch. Mirrors the `knowledge_jobs` table in KNOWLEDGE_PERSISTENCE_CONTRACT.md §2 field for field, deliberately: the table does not exist yet (Harry owns the dedorch migration, CLAUDE.md §2.2), so these models are what the in-memory store holds today and what the Postgres store will map one-to-one tomorrow. Keeping them aligned is what makes that swap a mapping exercise rather than a redesign. """ from __future__ import annotations import uuid from datetime import UTC, datetime from typing import Literal from pydantic import BaseModel, Field # Terminal states are succeeded / partial / failed / blocked. # # `partial` exists because one unreadable PDF must not fail a batch of twenty — # the caller needs to know which documents landed, not just that "something" # went wrong. # # `blocked` is NOT a failure: it means the pre-flight estimate exceeded the # spend ceiling and the job declined to continue. Folding it into `failed` # would hide the one outcome an operator can actually act on (raise the # ceiling, or ingest fewer documents at a time). JobStatus = Literal["queued", "running", "succeeded", "partial", "failed", "blocked"] # The stage machine, ordered by cost. Everything up to and including # `estimating` is free or near-free; `extracting` is where the LLM spend # happens, which is why the gate sits between them. JobStage = Literal["parsing", "filtering", "estimating", "extracting", "persisting"] DocumentStatus = Literal["pending", "succeeded", "failed", "skipped"] # r1 (2026-09-08 discussion). Chosen by the caller at ingest time. # "expert" - entries wait for a human to rule on them. # "llm" - entries go straight to approved, with NO reviewer step at all. # The second is not a second model call: the entries were produced by a model # already, so "review by LLM" means accepting them without human review. What # makes it safe to offer is that the decision is RECORDED against the person # who chose it (r2/r3), not that the pipeline checks itself. ReviewMode = Literal["expert", "llm"] def _now() -> datetime: return datetime.now(UTC) class DocumentResult(BaseModel): """One document's outcome inside a batch. Carries the ids of the rows it produced so the caller can go straight to the artifact without a second lookup. """ document_id: str filename: str | None = None status: DocumentStatus = "pending" knowledge_document_id: str | None = None run_id: str | None = None n_entries: int = 0 error: str | None = None # The pre-flight estimate, recorded whether or not the document was # refused. An operator cannot plan a batch from a number they only ever see # when it is already too late, and "why was this one refused when that one # ran" is not answerable without the shape of a run that succeeded. estimate: dict | None = None # Which ceiling refused this document (`clusters` / `prompt_tokens`), or # None if it was not refused. Distinguishes "declined before spending" from # every other failure, which is what lets a wholly-refused batch settle # `blocked` rather than `failed`. blocked_by: str | None = None class SkippedDocument(BaseModel): """A document the batch declined to take, and why. Reported rather than silently dropped: 6 of the 8 live `documents` rows are `xlsx`/`csv` and belong to the data path, so "we ignored most of your files" is information the caller needs, not noise. """ document_id: str filename: str | None = None reason: str class IngestJob(BaseModel): """One ingestion batch — the unit the user is waiting on.""" id: str = Field(default_factory=lambda: str(uuid.uuid4())) scope_id: str user_id: str status: JobStatus = "queued" stage: JobStage | None = None review_mode: ReviewMode = "expert" document_ids: list[str] = Field(default_factory=list) n_requested: int = 0 n_succeeded: int = 0 n_failed: int = 0 results: list[DocumentResult] = Field(default_factory=list) skipped: list[SkippedDocument] = Field(default_factory=list) estimated_cost: dict | None = None error_message: str | None = None created_at: datetime = Field(default_factory=_now) started_at: datetime | None = None completed_at: datetime | None = None @property def is_terminal(self) -> bool: return self.status in ("succeeded", "partial", "failed", "blocked") def settle(self) -> None: """Derive the terminal status from the per-document outcomes. Called once every document has been attempted. A job that produced nothing is a failure even though no exception escaped — "ran fine, extracted nothing" is not success. """ self.n_succeeded = sum(1 for r in self.results if r.status == "succeeded") self.n_failed = sum(1 for r in self.results if r.status == "failed") failures = [r for r in self.results if r.status == "failed"] if self.n_succeeded and self.n_failed: self.status = "partial" elif self.n_succeeded: self.status = "succeeded" elif failures and all(r.blocked_by for r in failures): # `blocked` was documented from the start and was unreachable until # 2026-09-11: a spend refusal settled as `failed`, which is the one # outcome an operator can act on (raise the ceiling, or split the # document) hidden inside the one they cannot. self.status = "blocked" else: self.status = "failed" self.stage = None self.completed_at = _now()