"""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." ), ) @field_validator("user_id") @classmethod 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 @router.post( "/ingest", response_model=IngestAccepted, status_code=status.HTTP_202_ACCEPTED, summary="Extract knowledge from a user's documents (async batch)", ) # 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. @limiter.limit("5/minute") @log_execution(logger) 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, ) @router.get( "/jobs/{job_id}", response_model=IngestJob, summary="Poll one ingestion job", ) @log_execution(logger) 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 @router.get( "/jobs", response_model=list[IngestJob], summary="Recent ingestion jobs for a user", ) @log_execution(logger) 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 @router.get( "/documents", response_model=list[KnowledgeDocumentSummary], summary="Parsed artifacts stored for this user's company", ) @log_execution(logger) 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 ] @router.post( "/documents/{knowledge_document_id}/extract", response_model=IngestAccepted, status_code=status.HTTP_202_ACCEPTED, summary="Re-extract a stored artifact (no re-parsing)", ) @limiter.limit("5/minute") @log_execution(logger) 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, ) @router.get( "/entries", response_model=list[EntrySummary], summary="Candidate knowledge entries", ) @log_execution(logger) 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 ] @router.get( "/queue", summary="The review queue — most-decision-worthy first", ) @log_execution(logger) 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 @router.post( "/entries/{entity_id}/reviews", response_model=ReviewAccepted, status_code=status.HTTP_201_CREATED, summary="Record an expert decision on one entry", ) @log_execution(logger) 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. @router.post( "/domain/propose", status_code=status.HTTP_201_CREATED, summary="Draft this company's domain declaration for review", ) @limiter.limit("10/minute") @log_execution(logger) 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" ), } @router.get( "/domain", summary="The domain context this company's agents would read", ) @log_execution(logger) 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, }