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