Download src/knowledge_ingest/models.py from DataEyond/Agentic-Service-Data-Eyond-Catalog: direct link, hf CLI and curl.
- Browser
- Download file 5.68 kB
-
https://huggingface.co/spaces/DataEyond/Agentic-Service-Data-Eyond-Catalog/resolve/main/src/knowledge_ingest/models.py
- Command line
-
hf download hf://spaces/DataEyond/Agentic-Service-Data-Eyond-Catalog/src/knowledge_ingest/models.py
-
curl -L -o models.py https://huggingface.co/spaces/DataEyond/Agentic-Service-Data-Eyond-Catalog/resolve/main/src/knowledge_ingest/models.py
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 | |
| 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() | |