ishaq101's picture
/fix parsing and term extract (#21)
f07443e
Raw History Blame Contribute Delete
5.68 kB
"""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()