Download src/api/v1/knowledge.py from DataEyond/Agentic-Service-Data-Eyond-Catalog: direct link, hf CLI and curl.
- Browser
- Download file 20.6 kB
-
https://huggingface.co/spaces/DataEyond/Agentic-Service-Data-Eyond-Catalog/resolve/main/src/api/v1/knowledge.py
- Command line
-
hf download hf://spaces/DataEyond/Agentic-Service-Data-Eyond-Catalog/src/api/v1/knowledge.py
-
curl -L -o knowledge.py https://huggingface.co/spaces/DataEyond/Agentic-Service-Data-Eyond-Catalog/resolve/main/src/api/v1/knowledge.py
20.6 kB
| """Knowledge ingestion endpoints (KNOWLEDGE_SURFACE_PLAN.md Β§4, N10/N10f). | |
| One click over all of a user's parseable documents: | |
| POST /api/v1/knowledge/ingest -> 202 {job_id} | |
| GET /api/v1/knowledge/jobs/{job_id} -> job state, per-document results | |
| GET /api/v1/knowledge/jobs -> recent jobs for a user | |
| The work is asynchronous by necessity β parse (~80 s) plus extract (~170 s) for | |
| a nine-page document is far past any HTTP timeout β so the route admits a batch | |
| and returns immediately. Polling rather than SSE (decided 2026-09-07): a job row | |
| survives a refresh and a dropped connection, a stream does not, and a 4-minute | |
| stream would hold a worker for its duration. | |
| `GET /jobs/{job_id}` returns **404** for an unknown id, following | |
| `GET /traceability` rather than the tri-state 200 of `GET /charts`. The | |
| difference is real: charts is called unconditionally on every turn where absence | |
| is normal, whereas a polled `job_id` was handed to the caller, so its absence is | |
| genuinely exceptional. | |
| β οΈ **No caller authentication β decided 2026-09-07, not inherited.** `user_id` is | |
| caller-supplied, exactly like every other live endpoint. The service-secret gate | |
| (DEV_PLAN Β§0.7 #37) cannot be armed here for the same reason it was unwired | |
| there: the browser SPA is the sole caller, we do not own that repo, and a secret | |
| shipped to a browser is not a secret. Arming it would 401 the feature. | |
| But this is a **write that spends LLM budget**, which is a different risk class | |
| from the read surface, so the posture is not inherited silently β it is bounded | |
| by four controls, and the gap they leave is stated rather than papered over: | |
| 1. **Rate limit** (below) β caps calls per client per window. | |
| 2. **Single-flight per scope** β a second click returns the running job. | |
| 3. **Pre-flight spend ceiling** β the free stages produce an exact estimate and | |
| the job stops rather than spending past it (N10c, lands with the pipeline). | |
| 4. **Content-addressed idempotency** β re-ingesting an unchanged document | |
| re-spends nothing (N10c). | |
| **What they do not do is identify the caller.** `user_id` is still whatever the | |
| request says, so one user can spend budget under another's id. That closes only | |
| with a real forwarded identity (DEV_PLAN #43). Recorded here so the next reader | |
| knows this is a bounded, accepted posture and not an oversight. | |
| `scope_id` is deliberately NOT accepted from the caller: it is derived | |
| server-side from `users.company`. Accepting a tenant key on an unauthenticated | |
| route is a cross-tenant write with extra steps. | |
| """ | |
| import uuid | |
| from datetime import datetime | |
| from typing import Literal | |
| from fastapi import APIRouter, HTTPException, Query, Request, status | |
| from pydantic import BaseModel, Field, field_validator | |
| from src.knowledge_domain.service import KnowledgeService | |
| from src.knowledge_extraction.queue.review_queue import build_queue | |
| from src.knowledge_ingest.knowledge_store import KnowledgeStore | |
| from src.knowledge_ingest.models import IngestJob | |
| from src.knowledge_ingest.selector import resolve_scope_id | |
| from src.knowledge_ingest.service import IngestService | |
| from src.middlewares.logging import get_logger, log_execution | |
| from src.middlewares.rate_limit import limiter | |
| # NOTE: this module must NOT use `from __future__ import annotations`. | |
| # With it, FastAPI cannot resolve a Pydantic body model's annotation through | |
| # the @log_execution wrapper and silently demotes it to a query parameter β | |
| # every POST then 422s with `loc: [query, payload]`. No router here carries it. | |
| logger = get_logger("knowledge_api") | |
| router = APIRouter(prefix="/api/v1/knowledge", tags=["Knowledge"]) | |
| # Warm, process-shared service β mirrors the chat and charts endpoints' | |
| # module-level instances. Holding one instance is also what makes single-flight | |
| # meaningful: the lock lives here. | |
| _service = IngestService() | |
| _store = KnowledgeStore() | |
| _knowledge = KnowledgeService(_store) | |
| def _require_uuid(value: str, field: str) -> str: | |
| """422 on a malformed id rather than a confusing downstream miss (F-22).""" | |
| try: | |
| uuid.UUID(value) | |
| except (ValueError, AttributeError, TypeError) as e: | |
| raise ValueError(f"{field} must be a UUID string") from e | |
| return value | |
| class IngestRequest(BaseModel): | |
| user_id: str = Field(..., description="The caller. Scope is derived from it.") | |
| document_ids: list[str] | None = Field( | |
| default=None, | |
| description=( | |
| "Optional. Omit to ingest every parseable document this user owns β " | |
| "the one-click case. Named ids are still filtered by owner and type." | |
| ), | |
| ) | |
| reparse: bool = Field( | |
| default=False, | |
| description=( | |
| "Re-parse even when the content hash already matches a stored " | |
| "artifact. Off by default: on, it silently re-spends on every call." | |
| ), | |
| ) | |
| review_mode: Literal["expert", "llm"] = Field( | |
| default="expert", | |
| description=( | |
| "Who rules on what this batch produces. 'expert' leaves every " | |
| "entry waiting for a human. 'llm' approves them on arrival with " | |
| "NO reviewer step and no second opinion β the entries were " | |
| "written by a model and are accepted as-is. Choosing it records " | |
| "`authorised_by` against the caller, so an entry approved this " | |
| "way can always say who allowed it." | |
| ), | |
| ) | |
| def _user_id_is_uuid(cls, v: str) -> str: | |
| return _require_uuid(v, "user_id") | |
| class IngestAccepted(BaseModel): | |
| job_id: str | |
| status: str | |
| stage: str | None = None | |
| n_requested: int | |
| reused: bool = Field( | |
| ..., | |
| description=( | |
| "True when a job for this scope was already in flight and is being " | |
| "returned instead of a second one starting." | |
| ), | |
| ) | |
| skipped: list[dict] = Field( | |
| default_factory=list, | |
| description="Documents the batch declined, each with a reason.", | |
| ) | |
| message: str | |
| # N11, control 1 of 4. Deliberately far below the chat route's 30/minute: this | |
| # starts GPU-less but LLM-spending batch work, and a legitimate caller clicks it | |
| # a few times a year, not a few times a minute. `slowapi` needs a Starlette | |
| # `Request` param named `request`, which is why the body is `payload`. | |
| # NOTE: `get_remote_address` buckets by client IP β if calls ever arrive through | |
| # a proxy, every caller shares one bucket. Revisit when identity is forwarded. | |
| async def ingest(request: Request, payload: IngestRequest) -> IngestAccepted: | |
| """Admit a batch and return immediately. Poll `GET /jobs/{job_id}`.""" | |
| job, reused = await _service.start( | |
| user_id=payload.user_id, | |
| document_ids=payload.document_ids, | |
| review_mode=payload.review_mode, | |
| ) | |
| if reused: | |
| message = "A job for this scope is already running; returning it." | |
| elif job.n_requested: | |
| message = f"Ingesting {job.n_requested} document(s)." | |
| else: | |
| message = job.error_message or "Nothing to ingest." | |
| return IngestAccepted( | |
| job_id=job.id, | |
| status=job.status, | |
| stage=job.stage, | |
| n_requested=job.n_requested, | |
| reused=reused, | |
| skipped=[s.model_dump() for s in job.skipped], | |
| message=message, | |
| ) | |
| async def get_job(job_id: str) -> IngestJob: | |
| """Job state and per-document outcomes. 404 when the id is unknown.""" | |
| job = await _service.get(job_id) | |
| if job is None: | |
| raise HTTPException( | |
| status_code=status.HTTP_404_NOT_FOUND, | |
| detail=f"No ingestion job with id {job_id}.", | |
| ) | |
| return job | |
| async def list_jobs( | |
| user_id: str = Query(..., description="Owner of the jobs"), | |
| limit: int = Query(20, ge=1, le=100), | |
| ) -> list[IngestJob]: | |
| """Lets a caller recover after a refresh without having stashed a job_id.""" | |
| return await _service.list_for_user(user_id, limit) | |
| # ββ Read + review surface (N10b, N12b) βββββββββββββββββββββββββββββββββββββββ | |
| # | |
| # Everything below resolves `scope_id` from `user_id` server-side and filters on | |
| # it. The tenant key is never accepted from the caller: taking it from the | |
| # request on an unauthenticated route would be a cross-tenant read with extra | |
| # steps. | |
| class KnowledgeDocumentSummary(BaseModel): | |
| knowledge_document_id: str | |
| doc_id: str | |
| document_id: str | None = None | |
| n_pages: int | |
| n_chunks: int | |
| schema_version: str | |
| parser: str | |
| content_hash: str | |
| created_at: datetime | None = None | |
| class EntrySummary(BaseModel): | |
| entity_id: str | |
| kind: str | |
| label: str | None = None | |
| mention_count: int | |
| extraction_status: str | None = None | |
| diff_status: str | None = None | |
| definition_conflict: bool | |
| payload: dict | |
| review: dict | None = Field( | |
| default=None, | |
| description="Current expert decision for this entity, if any.", | |
| ) | |
| class ReviewRequest(BaseModel): | |
| decision: Literal["approved", "edited", "rejected"] | |
| reviewer_id: str | |
| edited_payload: dict | None = Field( | |
| default=None, | |
| description="Required for `edited` β the expert's corrected entry.", | |
| ) | |
| note: str | None = None | |
| class ReviewAccepted(BaseModel): | |
| review_id: str | |
| entity_id: str | |
| decision: str | |
| content_hash: str | |
| message: str | |
| async def list_documents( | |
| user_id: str = Query(..., description="Scope is derived from this user"), | |
| limit: int = Query(50, ge=1, le=200), | |
| ) -> list[KnowledgeDocumentSummary]: | |
| """What can be re-extracted without re-parsing.""" | |
| scope_id = await resolve_scope_id(user_id) | |
| rows = await _store.list_documents(scope_id, limit) | |
| return [ | |
| KnowledgeDocumentSummary( | |
| knowledge_document_id=r.id, | |
| doc_id=r.doc_id, | |
| document_id=r.document_id, | |
| n_pages=r.n_pages or 0, | |
| n_chunks=len((r.artifact or {}).get("chunks") or []), | |
| schema_version=r.schema_version or "", | |
| parser=f"{r.parser_name}/{r.parser_backend or ''}".rstrip("/"), | |
| content_hash=r.content_hash, | |
| created_at=r.created_at, | |
| ) | |
| for r in rows | |
| ] | |
| async def reextract( | |
| request: Request, | |
| knowledge_document_id: str, | |
| user_id: str = Query(..., description="Scope is derived from this user"), | |
| ) -> IngestAccepted: | |
| """Re-run extraction over an artifact already in the database. | |
| The common case, and the reason it exists: prompt retuning must not re-pay | |
| the parser for a byte-identical artifact. Async like `/ingest` β extraction | |
| alone is ~170 s for nine pages β so poll `GET /jobs/{job_id}`. | |
| """ | |
| job, reused = await _service.start_reextract(user_id, knowledge_document_id) | |
| if reused: | |
| message = "A job for this scope is already running; returning it." | |
| elif job.status == "failed": | |
| message = job.error_message or "Nothing to extract." | |
| else: | |
| message = "Re-extracting the stored artifact." | |
| return IngestAccepted( | |
| job_id=job.id, | |
| status=job.status, | |
| stage=job.stage, | |
| n_requested=job.n_requested, | |
| reused=reused, | |
| skipped=[], | |
| message=message, | |
| ) | |
| async def list_entries( | |
| user_id: str = Query(..., description="Scope is derived from this user"), | |
| doc_id: str | None = Query(None, description="Filter to one document"), | |
| kind: Literal["glossary", "rule", "formula", "domain"] | None = Query(None), | |
| limit: int = Query(200, ge=1, le=1000), | |
| ) -> list[EntrySummary]: | |
| """The raw artifact read, annotated with any decision already recorded.""" | |
| scope_id = await resolve_scope_id(user_id) | |
| rows = await _store.list_entries( | |
| scope_id=scope_id, doc_id=doc_id, kind=kind, limit=limit | |
| ) | |
| reviews = await _store.latest_reviews(scope_id, [r.entity_id for r in rows]) | |
| return [ | |
| EntrySummary( | |
| entity_id=r.entity_id, | |
| kind=r.kind, | |
| label=r.label, | |
| mention_count=r.mention_count or 0, | |
| extraction_status=r.extraction_status, | |
| diff_status=r.diff_status, | |
| definition_conflict=bool(r.definition_conflict), | |
| payload=r.payload, | |
| review=_review_summary(reviews.get(r.entity_id)), | |
| ) | |
| for r in rows | |
| ] | |
| async def review_queue( | |
| user_id: str = Query(..., description="Scope is derived from this user"), | |
| doc_id: str | None = Query(None, description="Filter to one document"), | |
| ) -> list[dict]: | |
| """Rows an expert works through, in order. | |
| Derived at read time by the pipeline's own `build_queue`, never stored: | |
| persisting a derived ordering creates a second truth that drifts from the | |
| first. Reusing that function rather than re-sorting here is what keeps the | |
| endpoint and the CLI showing the same queue. | |
| Rows carry `row_kind` (`term` | `key_parameter`) and a consumer must switch | |
| on it: an unresolved key parameter has a surface and no term, no definition | |
| and no mention count. | |
| """ | |
| scope_id = await resolve_scope_id(user_id) | |
| glossary = await _store.list_entries( | |
| scope_id=scope_id, doc_id=doc_id, kind="glossary", limit=1000 | |
| ) | |
| domain_rows = await _store.list_entries( | |
| scope_id=scope_id, doc_id=doc_id, kind="domain", limit=1 | |
| ) | |
| domain = domain_rows[0].payload if domain_rows else None | |
| rows = build_queue([r.payload for r in glossary], domain) | |
| reviews = await _store.latest_reviews(scope_id) | |
| for row in rows: | |
| # `term_id` on a term row, absent on a key_parameter row. | |
| row["review"] = _review_summary(reviews.get(row.get("term_id") or "")) | |
| return rows | |
| async def record_review( | |
| entity_id: str, payload: ReviewRequest, user_id: str = Query(...) | |
| ) -> ReviewAccepted: | |
| """Approve, edit or reject one candidate. APPEND-ONLY. | |
| Keyed on `entity_id`, not the entry's row id: `entity_id` is deterministic | |
| and content-derived, so it survives re-extraction while row ids do not. That | |
| is what lets an expert's approvals outlive a prompt retune. | |
| `content_hash` is derived server-side from the artifact the entry came from, | |
| never taken from the request β a stale client would otherwise attribute a | |
| ruling to the wrong artifact version. | |
| β οΈ Approved and edited entries become the baseline a later extraction diffs | |
| against (`approved_glossary`), so an `edited_payload` written here can reach | |
| a future prompt. On an unauthenticated route that is a content-injection | |
| path, and it is the strongest argument for putting real auth in front of | |
| this surface (DEV_PLAN #43) before an expert relies on it. | |
| """ | |
| if payload.decision == "edited" and not payload.edited_payload: | |
| raise HTTPException( | |
| status_code=status.HTTP_422_UNPROCESSABLE_ENTITY, | |
| detail="decision 'edited' requires edited_payload.", | |
| ) | |
| scope_id = await resolve_scope_id(user_id) | |
| try: | |
| row = await _store.save_review( | |
| entity_id=entity_id, | |
| scope_id=scope_id, | |
| decision=payload.decision, | |
| reviewer_id=payload.reviewer_id, | |
| edited_payload=payload.edited_payload, | |
| note=payload.note, | |
| ) | |
| except LookupError as e: | |
| raise HTTPException( | |
| status_code=status.HTTP_404_NOT_FOUND, | |
| detail=f"No entry {entity_id} in this scope.", | |
| ) from e | |
| # The review IS the invalidation event. Without this the expert approves | |
| # something and the planner keeps reading the pre-approval card for up to an | |
| # hour, which makes the loop look broken in exactly the way that is hardest | |
| # to debug β the write succeeded and nothing changed. | |
| await _knowledge.invalidate(scope_id) | |
| return ReviewAccepted( | |
| review_id=row.id, | |
| entity_id=entity_id, | |
| decision=row.decision, | |
| content_hash=row.content_hash, | |
| message=f"Decision '{row.decision}' recorded.", | |
| ) | |
| def _review_summary(row) -> dict | None: | |
| """The decision a caller needs to see beside a row, without the payload.""" | |
| if row is None: | |
| return None | |
| return { | |
| "decision": row.decision, | |
| "reviewer_id": row.reviewer_id, | |
| "content_hash": row.content_hash, | |
| "note": row.note, | |
| "created_at": row.created_at.isoformat() if row.created_at else None, | |
| } | |
| # ββ The domain declaration (V4-9 / V4-10) ββββββββββββββββββββββββββββββββ | |
| # | |
| # The DECLARED half of a domain context: what it is, and what it deliberately is | |
| # not. Derived where it can be (the boundary is the complement of the covered | |
| # subdomains over a closed enum β checkable, unlike a model-written sentence), | |
| # left empty where only a human can answer. | |
| # | |
| # There is no new review machinery here on purpose. A declaration is stored as a | |
| # `kind='domain'` entry, so an expert approves or corrects it through the SAME | |
| # `POST /entries/{entity_id}/reviews` route as a glossary term. | |
| async def propose_domain(request: Request, user_id: str = Query(...)) -> dict: | |
| """Derive a draft declaration and store it awaiting review. | |
| **Calls no model.** Everything is derived from what the corpus demonstrably | |
| contains, so nothing in the draft can be wrong in a way a reviewer would | |
| have to detect rather than simply complete. Re-proposing replaces the | |
| standing draft; the expert's decisions key on `entity_id` and survive it. | |
| The result reaches no planner until it is approved. | |
| """ | |
| payload = await _knowledge.propose_declaration(user_id) | |
| return { | |
| "entity_id": payload["domain_id"], | |
| "declaration": payload, | |
| "message": ( | |
| "Draft stored. Approve or correct it via " | |
| f"POST /api/v1/knowledge/entries/{payload['domain_id']}/reviews" | |
| ), | |
| } | |
| async def get_domain( | |
| user_id: str = Query(...), | |
| rendered: bool = Query(False, description="Return the prompt text instead of the object"), | |
| ) -> dict: | |
| """The composed domain context, plus the standing declaration if there is one. | |
| `declaration_approved` is the field to read: until it is true the domain has | |
| no name and no boundary, and a planner cannot refuse on domain grounds. | |
| """ | |
| context = await _knowledge.get_context(user_id) | |
| declaration = await _knowledge.get_declaration(user_id) | |
| return { | |
| "scope_id": context.scope_id, | |
| "declaration": declaration, | |
| "declaration_approved": context.authority.coverage.declared, | |
| "n_measures": len(context.measures), | |
| "coverage": context.authority.coverage.model_dump(), | |
| "context": None if rendered else context.model_dump(mode="json"), | |
| "card": await _knowledge.render_context(user_id) if rendered else None, | |
| } | |