| """FastAPI backend: serves extraction results + page images, accepts |
| confirm/edit actions (docs/prd.md section 8.4), and -- when the frontend has been |
| built -- serves the SPA itself so the whole app lives on one origin. |
| |
| Every JSON endpoint lives under `/api`. That prefix is load-bearing rather |
| than cosmetic: the SPA's client-side routes `/dashboard`, `/export`, and |
| `/cases/:barcode` are character-for-character the same paths as three of the |
| API's own routes, so serving both from one origin at the root would make a |
| deep link to `/dashboard` return JSON instead of the app. |
| |
| Endpoints that spend money or change reviewer state require a per-user login |
| (see `authorize`, #128). Reads stay open: the corpus is public, de-identified |
| TCGA data, and the SPA has to render before anyone can authenticate. |
| """ |
|
|
| from __future__ import annotations |
|
|
| import json |
| import logging |
| import os |
| import re |
| import sqlite3 |
| import uuid |
| from contextlib import asynccontextmanager |
| from datetime import datetime, timezone |
| from pathlib import Path |
| from typing import Any, Optional |
|
|
| from fastapi import APIRouter, Depends, FastAPI, HTTPException |
| from fastapi.middleware.cors import CORSMiddleware |
| from fastapi.responses import FileResponse, RedirectResponse |
| from fastapi.security import HTTPAuthorizationCredentials, HTTPBearer |
| from pydantic import BaseModel, Field, ValidationError |
| from sqlalchemy import text |
|
|
| from endopath import auth, colpali_retrieval, fields, precomputed_retrieval, taxonomy, textsource |
| from endopath.auth import Identity |
| from endopath.dictionary import DataDictionary |
| from endopath.schema import CaseStatus, EndometrialChecklist, EvidenceSpan |
| from endopath.storage import ( |
| DEFAULT_DB_PATH, |
| DEFAULT_PROJECT_ID, |
| connect, |
| count_users, |
| create_project, |
| get_case, |
| get_user, |
| list_cases, |
| list_projects, |
| upsert_case, |
| upsert_user, |
| ) |
|
|
| logger = logging.getLogger(__name__) |
|
|
| IMAGE_DIR = Path("data") / "images" |
|
|
| FRONTEND_DIST_ENV = "ENDOPATH_FRONTEND_DIST" |
|
|
| |
| |
| |
| SPACE_ENV = "SPACE_ID" |
|
|
| _bearer_scheme = HTTPBearer(auto_error=False) |
|
|
|
|
| def authorize( |
| credentials: Optional[HTTPAuthorizationCredentials] = Depends(_bearer_scheme), |
| ) -> Optional[Identity]: |
| """Resolve the authenticated account for a write endpoint, or None when auth is disabled (#128). |
| |
| Expects `Authorization: Bearer <session-token>`, the token endopath.auth issues at login. With no |
| session secret configured the gate is open and this returns None, which keeps local dev and the test |
| suite free of ceremony; that default is only safe because a Space without a secret refuses to start |
| (see the lifespan below). With a secret set, a missing, malformed, or expired token is a 401, and a |
| valid one resolves to the Identity the endpoints read the confirmer and role from. |
| """ |
| secret = auth.session_secret() |
| if not secret: |
| return None |
| token = credentials.credentials if credentials is not None else "" |
| identity = auth.decode_token(token, secret) |
| if identity is None: |
| raise HTTPException( |
| status_code=401, |
| detail="missing or invalid session token", |
| headers={"WWW-Authenticate": "Bearer"}, |
| ) |
| return identity |
|
|
|
|
| def _attribution( |
| identity: Optional[Identity], body_confirmer: Optional[str], body_role: Optional[str] |
| ) -> tuple[Optional[str], Optional[str]]: |
| """The confirmer name and role recorded on a decision. When authenticated, the account's, ignoring any |
| body-supplied values, so attribution is non-repudiable (PRD section 4.4). When auth is disabled (local |
| dev, the test suite), the body-supplied session identity, preserving the open-path behavior.""" |
| if identity is not None: |
| return identity.confirmer, identity.role |
| return body_confirmer, body_role |
|
|
|
|
| def _require_licence(identity: Optional[Identity], licence_required: bool) -> None: |
| """A licence_required decision may only be confirmed by an account holding a pathology licence (PRD |
| section 4.4). Enforced only when authenticated; the open dev path does not gate on licence.""" |
| if identity is not None and licence_required and not identity.holds_licence: |
| raise HTTPException( |
| status_code=403, |
| detail="confirming this requires a pathology licence; the signed-in account holds none", |
| ) |
|
|
|
|
| def _seed_users() -> None: |
| """Provision the accounts named in ENDOPATH_SEED_USERS into the store at startup, so a fresh |
| deployment has logins. Upserts, so a redeploy rotates a seeded password in place. A malformed seed |
| raises, failing the boot loudly rather than starting an auth-enabled Space with no way to log in.""" |
| seeds = auth.parse_seed_users(os.environ.get(auth.SEED_USERS_ENV, "")) |
| if not seeds: |
| return |
| conn = connect(_db_path) |
| try: |
| for seed in seeds: |
| upsert_user( |
| conn, |
| username=seed.username, |
| password_hash=auth.hash_password(seed.password), |
| role=seed.role, |
| holds_licence=seed.holds_licence, |
| display_name=seed.display_name, |
| ) |
| total = count_users(conn) |
| finally: |
| conn.close() |
| logger.info("seeded %d reviewer account(s); %d total", len(seeds), total) |
|
|
|
|
| @asynccontextmanager |
| async def lifespan(_app: FastAPI): |
| if auth.session_secret(): |
| _seed_users() |
| logger.info( |
| "%s is set; write endpoints require a reviewer login", auth.SESSION_SECRET_ENV |
| ) |
| elif os.environ.get(SPACE_ENV): |
| |
| |
| raise RuntimeError( |
| f"{auth.SESSION_SECRET_ENV} is not set but {SPACE_ENV} is: refusing to expose " |
| "unauthenticated extraction and review endpoints on a public URL. Set " |
| f"{auth.SESSION_SECRET_ENV} and seed accounts via {auth.SEED_USERS_ENV} as Space Secrets, " |
| "or make the Space private." |
| ) |
| else: |
| logger.warning( |
| "%s is not set: extraction and review endpoints are unauthenticated. " |
| "Fine locally, not for a deployment reachable from the internet.", |
| auth.SESSION_SECRET_ENV, |
| ) |
|
|
| |
| |
| from endopath import llm, provenance |
|
|
| llm.configure_recorder(provenance.make_recorder(_db_path)) |
|
|
| |
| |
| |
| from endopath import worker |
|
|
| if worker.start(_db_path): |
| logger.info("Intake follow-up worker started against %s", _db_path) |
| try: |
| yield |
| finally: |
| worker.stop() |
|
|
|
|
| app = FastAPI(title="Provinans API", lifespan=lifespan) |
| api = APIRouter(prefix="/api") |
|
|
| |
| |
| |
| authenticated = [Depends(authorize)] |
|
|
| |
| |
| |
| app.add_middleware( |
| CORSMiddleware, |
| allow_origins=["http://localhost:5173", "http://127.0.0.1:5173"], |
| allow_methods=["*"], |
| allow_headers=["*"], |
| ) |
|
|
| _db_path = DEFAULT_DB_PATH |
| _frontend_dist = Path(os.environ.get(FRONTEND_DIST_ENV, "frontend/dist")) |
|
|
|
|
| def configure_db_path(path: Path | str) -> None: |
| """Test/CLI hook to point the API at a different SQLite file.""" |
| global _db_path |
| _db_path = Path(path) |
|
|
|
|
| def configure_frontend_dist(path: Path | str) -> None: |
| """Test/CLI hook to point the SPA fallback at a different build directory.""" |
| global _frontend_dist |
| _frontend_dist = Path(path) |
|
|
|
|
| def get_db(): |
| conn = connect(_db_path) |
| try: |
| yield conn |
| finally: |
| conn.close() |
|
|
|
|
| @api.get("/trace") |
| def model_trace(limit: int = 50, conn: sqlite3.Connection = Depends(get_db)) -> dict: |
| """The live model-call trace (#221): the most recent LLM calls with the requested and the resolved |
| model, token usage, latency, and the Anthropic request id, plus the running aggregate. Read only, |
| metadata only. A viewer watches this to confirm the app calls the model it claims to.""" |
| from endopath import provenance |
|
|
| return {"calls": provenance.recent(conn, limit), "summary": provenance.summary(conn)} |
|
|
|
|
| @api.get("/coverage") |
| def coverage(project: Optional[str] = None, conn: sqlite3.Connection = Depends(get_db)) -> dict: |
| """The Field coverage view (#237): a project's target dictionary, DB-backed, with the record type that |
| fulfils each field and its fill config. The live fill is computed in the browser from /api/dashboard and |
| /api/longitudinal/cohort. Seeds the default project's coverage on first read.""" |
| from endopath import coverage as coverage_mod |
| from endopath import coverage_data |
|
|
| coverage_mod.seed_coverage(conn) |
| return coverage_mod.build_coverage(conn, project or coverage_data.PROJECT_ID) |
|
|
|
|
| class ReportSource(BaseModel): |
| """Where the reports come from. `prepared_corpus` walks the pre-fetched TCGA |
| UCEC corpus; `folder` ingests a local directory of report PDFs (#25).""" |
|
|
| kind: str = "prepared_corpus" |
| folder_path: Optional[str] = None |
|
|
|
|
| class IngestionRunRequest(BaseModel): |
| """One Intake submission: a report source, a target dictionary, and the text |
| source that reads each report. All optional, so a bare POST still runs the |
| prepared corpus against the built-in CAP checklist.""" |
|
|
| source: ReportSource = Field(default_factory=ReportSource) |
| dictionary: Optional[DataDictionary] = None |
| text_source: Optional[str] = None |
| limit: Optional[int] = None |
|
|
|
|
| def _default_text_source_for(source_kind: str) -> str: |
| """The prepared corpus already carries recognition text, so it reads through |
| corpus OCR by default; a folder source has none, so it reads the rendered |
| pages with the vision model, the MVP path. Either is overridable per |
| submission.""" |
| return textsource.CORPUS_OCR if source_kind == "prepared_corpus" else textsource.DEFAULT_SOURCE_ID |
|
|
|
|
| @api.post("/ingestion/run", dependencies=authenticated) |
| def trigger_ingestion( |
| limit: Optional[int] = None, |
| body: Optional[IngestionRunRequest] = None, |
| conn: sqlite3.Connection = Depends(get_db), |
| ) -> dict: |
| """Intake screen action (docs/prd.md section 8.1): point the tool at a report |
| source and a target dictionary, ingest, and enqueue the follow-up jobs for |
| that pair (report induction #92, dictionary compilation #91). |
| |
| The report's text is read through a pluggable text source (docs/prd.md |
| section 7, "Two points plug out"): corpus OCR, the PDF text layer (#45), or a |
| vision model read of the rendered pages. This call site names no provider; |
| it selects a source by id, so a local on-premise model swaps in unchanged. |
| |
| The corpus pipeline is idempotent, so re-running it (re-clicking Intake) is |
| safe and cheap once the corpus is already local. |
| """ |
| from pathlib import Path as _Path |
|
|
| import pandas as pd |
|
|
| from endopath import dictionary as dictionary_mod |
| from endopath import ingestion, jobs |
|
|
| req = body or IngestionRunRequest() |
| effective_limit = req.limit if req.limit is not None else limit |
| source = req.source |
| data_dictionary = dictionary_mod.register( |
| conn, req.dictionary or dictionary_mod.builtin_cap_dictionary() |
| ) |
|
|
| source_id = req.text_source or _default_text_source_for(source.kind) |
| try: |
| text_source = textsource.select(source_id) |
| except ValueError as exc: |
| raise HTTPException(status_code=422, detail=str(exc)) from exc |
|
|
| if source.kind == "folder": |
| if not source.folder_path: |
| raise HTTPException(status_code=422, detail="folder source requires folder_path") |
| folder = _Path(source.folder_path) |
| try: |
| results_df = ingestion.run_folder_pipeline(folder, limit=effective_limit) |
| except (NotADirectoryError, FileNotFoundError) as exc: |
| raise HTTPException(status_code=422, detail=str(exc)) from exc |
| |
| |
| empty_manifest = pd.DataFrame({"case_barcode": [], "text": []}) |
| seeded = ingestion.seed_cases_into_store( |
| conn, results_df, manifest=empty_manifest, text_source=text_source, folder=folder |
| ) |
| elif source.kind == "prepared_corpus": |
| results_df = ingestion.run_pipeline(limit=effective_limit) |
| seeded = ingestion.seed_cases_into_store(conn, results_df, text_source=text_source) |
| else: |
| raise HTTPException(status_code=422, detail=f"unknown report source kind: {source.kind!r}") |
|
|
| case_count = len(results_df) |
| failures = int(results_df.error.notna().sum()) if case_count else 0 |
| enqueued = jobs.enqueue_intake_followups( |
| conn, |
| dictionary_id=data_dictionary.id, |
| source_kind=source.kind, |
| case_count=case_count, |
| variable_count=len(data_dictionary.variables), |
| ) |
| return { |
| "ingested": case_count, |
| "failures": failures, |
| "seeded": seeded, |
| "source": source.model_dump(), |
| "text_source": {"id": text_source.id, "label": text_source.label}, |
| "dictionary": { |
| "id": data_dictionary.id, |
| "title": data_dictionary.title, |
| "variable_count": len(data_dictionary.variables), |
| "licence_required_variables": data_dictionary.licence_required_variables, |
| }, |
| "enqueued": [job.model_dump() for job in enqueued], |
| } |
|
|
|
|
| @api.get("/ingestion/sources") |
| def ingestion_text_sources() -> list[dict]: |
| """The selectable text sources for the Intake screen (id + label). An open |
| read: the SPA renders the picker before anyone authenticates.""" |
| return textsource.available_sources() |
|
|
|
|
| @api.get("/ingestion/dictionary") |
| def builtin_dictionary() -> dict: |
| """The built-in CAP endometrium checklist as a dictionary, the default |
| target an Intake submission points at. An open read, so the SPA can render |
| the dictionary preview before a reviewer has a token.""" |
| from endopath import dictionary as dictionary_mod |
|
|
| return dictionary_mod.builtin_cap_dictionary().model_dump() |
|
|
|
|
| @api.get("/ingestion/jobs") |
| def ingestion_jobs(conn: sqlite3.Connection = Depends(get_db)) -> list[dict]: |
| """The follow-up jobs Intake submissions have enqueued (report induction |
| #92, dictionary compilation #91), read back from the durable queue. An open read.""" |
| from endopath import jobs |
|
|
| return [job.model_dump() for job in jobs.pending(conn)] |
|
|
|
|
| @api.get("/fields") |
| def checklist_fields() -> list[dict]: |
| """The checklist definition, so the SPA renders a schema it is handed |
| rather than a hand-mirrored copy (#37). Order is meaningful: it is the |
| order Case Review lays out evidence cards. |
| |
| `enum_values` is served alongside the type so the edit path can offer a |
| dropdown of the permitted values instead of a free-text box. |
| """ |
| return [ |
| { |
| "name": spec.name, |
| "label": spec.label, |
| "value_type": spec.value_type, |
| "enum_values": spec.enum_values, |
| } |
| for spec in fields.FIELDS |
| ] |
|
|
|
|
| @api.get("/concepts") |
| def concepts() -> dict: |
| """The canonical concept list (#91) with its registry evidence, dependency |
| relationships, and decisions joined on (#35), for the registry browser |
| (#94, docs/prd.md section 8.1). The concept spine is |
| `mapping.concept_list()`, served rather than hand-mirrored, so the browser |
| never diverges from the exported column spec. An open read: the SPA renders |
| the browser before anyone authenticates. |
| """ |
| from endopath import concept_registry |
|
|
| return concept_registry.concept_registry() |
|
|
|
|
| @api.get("/mapping") |
| def mapping(dictionary: Optional[str] = None, conn: sqlite3.Connection = Depends(get_db)) -> dict: |
| """The crosswalk mapping-review surface (#95, #109, docs/prd.md sections 8.1, |
| 3; design-guideline section 9.1): the drafted edges from a source dictionary's |
| fields to the canonical concepts, colored by category, with the unresolved and |
| low-confidence edges marked as the reviewer queue and each queued edge carrying |
| its maieutic question and two candidate resolutions. |
| |
| `dictionary` selects a submitted dictionary's drafted crosswalk (#108); omitted, |
| it serves the built-in CAP `crosswalk.reference_crosswalk()`, the committed |
| reference instance. Joined to the confirmations recorded so far. An open read: |
| the SPA renders the diagram before anyone authenticates. |
| """ |
| from endopath import mapping_review |
|
|
| try: |
| return mapping_review.mapping_review(conn, dictionary) |
| except mapping_review.MappingReviewError as exc: |
| |
| raise HTTPException(status_code=404, detail=str(exc)) from exc |
|
|
|
|
| @api.get("/dictionaries") |
| def dictionaries(conn: sqlite3.Connection = Depends(get_db)) -> list[dict]: |
| """The dictionaries a reviewer can open on the mapping screen: the built-in CAP |
| checklist always, then any submitted dictionary, each flagged with whether its |
| crosswalk has been drafted yet (#109). An open read.""" |
| from endopath import dictionary as dictionary_mod |
| from endopath import storage as storage_mod |
|
|
| out = [ |
| { |
| "id": dictionary_mod.BUILTIN_DICTIONARY_ID, |
| "title": dictionary_mod.BUILTIN_DICTIONARY_TITLE, |
| "builtin": True, |
| "ready": True, |
| } |
| ] |
| for row in storage_mod.list_dictionaries(conn): |
| if row["id"] == dictionary_mod.BUILTIN_DICTIONARY_ID: |
| continue |
| ready = storage_mod.get_submission_artifact(conn, row["id"], "crosswalk_draft") is not None |
| out.append({"id": row["id"], "title": row["title"], "builtin": False, "ready": ready}) |
| return out |
|
|
|
|
| class LoginRequest(BaseModel): |
| """A reviewer's credentials at login (#128).""" |
|
|
| username: str |
| password: str |
|
|
|
|
| @api.post("/login") |
| def login(body: LoginRequest, conn: sqlite3.Connection = Depends(get_db)) -> dict: |
| """Exchange a username and password for a signed session token (#128). The token carries the account's |
| role and licence, which the write endpoints read the confirmer and role from, so a confirmation's |
| attribution is the authenticated account's rather than a self-declared field. Only meaningful when auth |
| is enabled; with no session secret the write endpoints are open and no login is needed.""" |
| secret = auth.session_secret() |
| if not secret: |
| raise HTTPException(status_code=409, detail="authentication is not enabled on this deployment") |
| user = get_user(conn, body.username) |
| if user is None or not auth.verify_password(body.password, user["password_hash"]): |
| |
| |
| raise HTTPException(status_code=401, detail="invalid username or password") |
| identity = Identity( |
| username=user["username"], |
| role=user["role"], |
| holds_licence=user["holds_licence"], |
| display_name=user["display_name"], |
| ) |
| return { |
| "token": auth.issue_token(identity, secret), |
| "username": identity.username, |
| "display_name": identity.display_name, |
| "role": identity.role, |
| "holds_licence": identity.holds_licence, |
| } |
|
|
|
|
| |
| |
| |
|
|
| class LongitudinalRecord(BaseModel): |
| """A single longitudinal event (follow-up, treatment, recurrence, vital status).""" |
|
|
| record_id: str |
| participant_id: str |
| case_barcode: Optional[str] = None |
| record_type: str |
| record_index: int |
| days_to_followup: Optional[int] = None |
| days_to_treatment: Optional[int] = None |
| days_to_recurrence: Optional[int] = None |
| biopsy_result: Optional[str] = None |
| followup_treatment_success: Optional[str] = None |
| recurrence_type: Optional[str] = None |
| therapy_type: Optional[str] = None |
| structured_data: dict[str, Any] = {} |
|
|
|
|
| class LongitudinalJoin(BaseModel): |
| """A confirmed association between a case and a longitudinal record.""" |
|
|
| join_id: str |
| case_barcode: str |
| record_id: str |
| disposition: str |
| confirmed_by: Optional[str] = None |
| role: Optional[str] = None |
| confirmed_at: Optional[str] = None |
| note: Optional[str] = None |
| question: Optional[str] = None |
|
|
|
|
| class JoinRequest(BaseModel): |
| """Request to confirm/reject a longitudinal record with a case.""" |
|
|
| record_id: str |
| disposition: str |
| note: Optional[str] = None |
| question: Optional[str] = None |
|
|
|
|
| class LinkageRequest(BaseModel): |
| """Request to affirm a case's whole GDC longitudinal record in one action (#180).""" |
|
|
| note: Optional[str] = None |
|
|
|
|
| |
|
|
| class MappingComponentInput(BaseModel): |
| """One component concept a split resolves a source field to (#103 follow-on): |
| the canonical concept and how the field carries it. Never 'unresolved'.""" |
|
|
| concept_id: str |
| relation: str |
|
|
|
|
| class MappingConfirmRequest(BaseModel): |
| """One reviewer's decision on a crosswalk edge (#95, #103): the edge, the |
| chosen concept and relation (concept null for an out-of-scope resolution), |
| the per-variable licence-required flag, and the decision-ledger context. |
| `disposition` is `confirmed`, `flagged`, or `deferred`; `decided_via` names |
| how the answer was reached; `question` and `options` record what the decision |
| weighed. `relation` is optional for a flag or a deferral, which resolves to none. |
| |
| `components`, when present, splits the field across several canonical concepts |
| (#103 follow-on): `concept_id` and `relation` become the derived headline, and |
| each component is a primitive the field also resolves to. A split is always a |
| `confirmed` resolution. |
| |
| `confirmer` and `role` are optional (#128): when auth is enabled the server |
| reads them from the authenticated account and ignores the body, so a |
| confirmation's attribution is non-repudiable; the open dev path still reads |
| them from the body-supplied session identity.""" |
|
|
| mapping_id: str |
| concept_id: Optional[str] = None |
| relation: str = "unresolved" |
| confirmer: Optional[str] = None |
| role: Optional[str] = None |
| note: Optional[str] = None |
| licence_required: bool = False |
| disposition: str = "confirmed" |
| decided_via: Optional[str] = None |
| question: Optional[str] = None |
| options: Optional[list] = None |
| components: Optional[list[MappingComponentInput]] = None |
|
|
|
|
| @api.post("/mapping/confirm", dependencies=authenticated) |
| def confirm_mapping( |
| body: MappingConfirmRequest, |
| conn: sqlite3.Connection = Depends(get_db), |
| identity: Optional[Identity] = Depends(authorize), |
| ) -> dict: |
| """Record one reviewer's decision on a crosswalk edge (#95, #103): the |
| disposition, the resolved concept and relation, the confirmer and their |
| role, the licence-required flag, the time, and the ledger context. Any edge |
| is decidable, not only the queued ones. Resolving the last queued edge, with |
| nothing flagged or deferred, confirms the whole mapping, at which point |
| `GET /api/mapping` exposes the confirmed field set (#96) reads. |
| |
| When auth is enabled (#128) the confirmer and role are the authenticated |
| account's, not the body's, and a licence_required edge is refused unless the |
| account holds a pathology licence. On the open dev path the body-supplied |
| session identity is used, as before. |
| """ |
| from endopath import mapping_review |
|
|
| _require_licence(identity, body.licence_required) |
| confirmer, role = _attribution(identity, body.confirmer, body.role) |
|
|
| try: |
| if body.components: |
| mapping_review.confirm_split( |
| conn, |
| mapping_id=body.mapping_id, |
| headline_concept_id=body.concept_id, |
| headline_relation=body.relation, |
| components=[c.model_dump() for c in body.components], |
| confirmer=confirmer, |
| role=role, |
| note=body.note, |
| licence_required=body.licence_required, |
| decided_via=body.decided_via, |
| question=body.question, |
| options=body.options, |
| ) |
| else: |
| mapping_review.confirm_edge( |
| conn, |
| mapping_id=body.mapping_id, |
| concept_id=body.concept_id, |
| relation=body.relation, |
| confirmer=confirmer, |
| role=role, |
| note=body.note, |
| licence_required=body.licence_required, |
| disposition=body.disposition, |
| decided_via=body.decided_via, |
| question=body.question, |
| options=body.options, |
| ) |
| except mapping_review.MappingReviewError as exc: |
| raise HTTPException(status_code=422, detail=str(exc)) from exc |
|
|
| return mapping_review.mapping_review(conn) |
|
|
|
|
| class MappingChatTurn(BaseModel): |
| """One prior turn of an edge chat: the speaker (`user` or `assistant`) and |
| the text, so a follow-up question carries its context.""" |
|
|
| role: str |
| text: str |
|
|
|
|
| class MappingChatRequest(BaseModel): |
| """A reviewer's question about one crosswalk edge (#103), with the prior |
| turns of that edge's conversation.""" |
|
|
| mapping_id: str |
| question: str |
| history: list[MappingChatTurn] = Field(default_factory=list) |
|
|
|
|
| @api.post("/mapping/edge/chat", dependencies=authenticated) |
| def mapping_edge_chat( |
| body: MappingChatRequest, conn: sqlite3.Connection = Depends(get_db) |
| ) -> dict: |
| """Reply to a reviewer's message about one crosswalk edge (#103): a concise |
| answer and, when the message points at a better resolution, one concrete remap |
| (a re-target or a split) the reviewer can apply. The exchange resolves into a |
| filed decision when the reviewer applies the proposal through |
| `POST /api/mapping/confirm` with `decided_via='chat'`. |
| |
| Gated by the shared token: the call spends the Anthropic key. The chat model |
| is named explicitly in `mapping_review` (claude-sonnet-5), behind a |
| provider-agnostic seam a fake substitutes for in the fast test suite. |
| """ |
| from endopath import mapping_review |
|
|
| try: |
| return mapping_review.explain_edge( |
| conn, |
| mapping_id=body.mapping_id, |
| question=body.question, |
| history=[turn.model_dump() for turn in body.history], |
| ) |
| except mapping_review.MappingReviewError as exc: |
| raise HTTPException(status_code=422, detail=str(exc)) from exc |
| except RuntimeError as exc: |
| raise HTTPException(status_code=503, detail=str(exc)) from exc |
|
|
|
|
| @api.get("/dashboard") |
| def dashboard( |
| project: Optional[str] = None, |
| conn: sqlite3.Connection = Depends(get_db), |
| ) -> dict: |
| """The cohort dashboard (#100, docs/prd.md sections 8.1, 4.3, 4.5, 4.7): the completion rate and |
| worst-field-first summary, then per field the confirmed-value and coded-absence-reason distributions |
| and the absence-versus-grade association, and per case the absent variables that would require ordering |
| slides. `project` scopes the whole payload to one project's cases (#149); omitted, the whole cohort is |
| summarized. An open read: the SPA renders the dashboard before anyone authenticates. |
| """ |
| from endopath.dashboard import cohort_dashboard |
|
|
| return cohort_dashboard(conn, project) |
|
|
|
|
| @api.get("/export") |
| def export( |
| staging_edition: Optional[str] = None, |
| project: Optional[str] = None, |
| conn: sqlite3.Connection = Depends(get_db), |
| ) -> dict: |
| """The cohort export (#101, docs/prd.md sections 3, 8.1): the confirmed case records projected onto |
| the canonical concept list as the column spec, one row per case, each value column carrying a parallel |
| coded absent-reason column, alongside the reproduction bundle and the three rendered download files. |
| |
| One projection feeds the on-screen preview and the CSV, XLSX, and JSON downloads, so the files and the |
| preview never disagree. `staging_edition` names the edition the researcher chose at export (#40); the |
| FIGO stage is re-projected under it through `taxonomy.derive_stage`, rendering a determined code, a |
| constrained candidate set, or an indeterminate reason distinctly. Omitted, each case projects under its |
| own recorded edition. `project` scopes the export to one project's cases (#149); omitted, every case in |
| the store is projected. An open read: the SPA renders the preview before anyone authenticates. |
| """ |
| from endopath import export as export_mod |
| from endopath import projection |
|
|
| if staging_edition is not None and not projection.is_offered_edition(staging_edition): |
| raise HTTPException(status_code=422, detail=f"unknown staging edition: {staging_edition}") |
| try: |
| result = export_mod.build_export( |
| conn, |
| staging_edition=staging_edition, |
| project=project, |
| statuses=export_mod.EXPORTABLE_STATUSES, |
| ) |
| except taxonomy.RegistryError as exc: |
| raise HTTPException(status_code=422, detail=str(exc)) from exc |
| return { |
| "columns": result["columns"], |
| "rows": result["rows"], |
| "provenance": result["provenance"], |
| "bundle": result["bundle"], |
| "files": export_mod.render_files(result), |
| } |
|
|
|
|
| @api.get("/staging_editions") |
| def staging_editions() -> list[dict]: |
| """The FIGO editions the Export target-schema selector offers (#40). An open read.""" |
| from endopath import projection |
|
|
| return projection.STAGING_EDITIONS |
|
|
|
|
| @api.get("/cases/{case_barcode}/projection") |
| def case_projection( |
| case_barcode: str, |
| edition: Optional[str] = None, |
| conn: sqlite3.Connection = Depends(get_db), |
| ) -> dict: |
| """Project one case's FIGO stage under a target edition (#40): a determined code, a constrained |
| candidate set the answer narrows to, or an indeterminate reason, each with its confidence and the |
| derivation that produced it. Omitted, the case's own recorded edition. An open read.""" |
| from endopath import projection |
|
|
| result = get_case(conn, case_barcode) |
| if result is None: |
| raise HTTPException(status_code=404, detail=f"case {case_barcode} not found") |
| case, _ = result |
| if edition is not None and not projection.is_offered_edition(edition): |
| raise HTTPException(status_code=422, detail=f"unknown staging edition: {edition}") |
| target = edition or case.staging_edition.value |
| try: |
| verdict = projection.project_stage(case.checklist, target) |
| except taxonomy.RegistryError as exc: |
| raise HTTPException(status_code=422, detail=str(exc)) from exc |
| return {"case_barcode": case_barcode, "edition": target, "figo_stage": verdict} |
|
|
|
|
| @api.get("/cases") |
| def worklist( |
| status: Optional[str] = None, |
| field_name: Optional[str] = None, |
| field_status: Optional[str] = None, |
| max_confidence: Optional[float] = None, |
| project: Optional[str] = None, |
| conn: sqlite3.Connection = Depends(get_db), |
| ) -> list[dict]: |
| return list_cases( |
| conn, |
| status=status, |
| field_name=field_name, |
| field_status=field_status, |
| max_confidence=max_confidence, |
| project=project, |
| ) |
|
|
|
|
| def _slugify_project_id(name: str, existing: set[str]) -> str: |
| """Derive a URL-safe project_id from a project name: lowercased, each run of non-alphanumeric |
| characters collapsed to a single hyphen, leading and trailing hyphens stripped. A numeric suffix is |
| appended when the base slug is already taken, so two projects sharing a name get distinct ids. Falls |
| back to 'project' when the name carries no alphanumeric character.""" |
| base = re.sub(r"[^a-z0-9]+", "-", name.lower()).strip("-") or "project" |
| if base not in existing: |
| return base |
| suffix = 2 |
| while f"{base}-{suffix}" in existing: |
| suffix += 1 |
| return f"{base}-{suffix}" |
|
|
|
|
| class CreateProjectRequest(BaseModel): |
| """A new project's human-readable name. The project_id slug is derived server-side from it.""" |
|
|
| name: str |
|
|
|
|
| @api.get("/projects") |
| def projects_list(conn: sqlite3.Connection = Depends(get_db)) -> dict: |
| """Every project with its case count, the default project first. An open read: the SPA renders |
| the project picker before anyone authenticates.""" |
| return {"projects": list_projects(conn), "default": DEFAULT_PROJECT_ID} |
|
|
|
|
| @api.post("/projects") |
| def projects_create( |
| body: CreateProjectRequest, |
| conn: sqlite3.Connection = Depends(get_db), |
| identity: Optional[Identity] = Depends(authorize), |
| ) -> dict: |
| """Create a project: derive a slug project_id from the name and insert it. Admin-gated when |
| auth is enabled: an authenticated non-admin account is refused with a 403; the open dev path (no |
| identity, local dev and the test suite) allows it. Returns the created project with a zero case count.""" |
| if identity is not None and identity.role != "admin": |
| raise HTTPException(status_code=403, detail="admin role required") |
| name = body.name.strip() |
| if not name: |
| raise HTTPException(status_code=422, detail="project name is required") |
| existing = {project["project_id"] for project in list_projects(conn)} |
| project_id = _slugify_project_id(name, existing) |
| created_by = identity.confirmer if identity is not None else None |
| try: |
| return create_project(conn, project_id=project_id, name=name, created_by=created_by) |
| except ValueError as exc: |
| raise HTTPException(status_code=409, detail=str(exc)) from exc |
|
|
|
|
| @api.post("/cases/extract_all", dependencies=authenticated) |
| def extract_all(conn: sqlite3.Connection = Depends(get_db)) -> dict: |
| """Enqueue a cohort-extraction job (#115): the worker drains every queued case through the shared |
| pipeline in the background, so a whole cohort is extracted without one /process call per case. Gated |
| by the shared token, since the job spends the Anthropic key. Returns the enqueued job.""" |
| from endopath import jobs |
|
|
| queued = list_cases(conn, status=CaseStatus.QUEUED.value) |
| job = jobs.enqueue_cohort_extraction(conn, case_count=len(queued)) |
| return job.model_dump() |
|
|
|
|
| def _page_count(patient_filename: Optional[str]) -> int: |
| if patient_filename is None: |
| return 0 |
| page_dir = IMAGE_DIR / patient_filename |
| if not page_dir.is_dir(): |
| return 0 |
| return len(list(page_dir.glob("page_*.png"))) |
|
|
|
|
| @api.get("/cases/{case_barcode}") |
| def case_detail(case_barcode: str, conn: sqlite3.Connection = Depends(get_db)) -> dict: |
| result = get_case(conn, case_barcode) |
| if result is None: |
| raise HTTPException(status_code=404, detail="case not found") |
| case, patient_filename = result |
| return { |
| "case_barcode": case.case_barcode, |
| "patient_filename": patient_filename, |
| "report_date": case.report_date.isoformat() if case.report_date else None, |
| "staging_edition": case.staging_edition.value, |
| "status": case.status.value, |
| "page_count": _page_count(patient_filename), |
| "fields": { |
| name: field.model_dump() for name, field in case.checklist.all_fields().items() |
| }, |
| } |
|
|
|
|
| @api.get("/cases/{case_barcode}/review") |
| def case_review_surface(case_barcode: str, conn: sqlite3.Connection = Depends(get_db)) -> dict: |
| """The case-review confirmation surface (#98, docs/prd.md sections 8.1, 8.3, 4.6; design-guideline |
| section 9.3): the case's fields grouped into sections as record cards, each carrying its B1 validation |
| block (#97) that promotes a field from needs-review to confirmed-eligible, and the decision ledger. |
| |
| An open read: the SPA renders the surface before anyone authenticates. The value chat, the confirm, |
| and the edit actions on it are gated. |
| """ |
| from endopath import case_review as case_review_mod |
|
|
| result = get_case(conn, case_barcode) |
| if result is None: |
| raise HTTPException(status_code=404, detail="case not found") |
| case, patient_filename = result |
| return case_review_mod.case_review( |
| conn, case, patient_filename=patient_filename, page_count=_page_count(patient_filename) |
| ) |
|
|
|
|
| @api.post("/cases/{case_barcode}/process", dependencies=authenticated) |
| def process_case_endpoint( |
| case_barcode: str, |
| use_visual: bool = False, |
| conn: sqlite3.Connection = Depends(get_db), |
| ) -> dict: |
| """Run first-pass inference on a queued case (docs/prd.md section 8.2: |
| Queued -> Processing -> Ready for Review). The LLM text-extraction pass is |
| the fast core; use_visual additionally runs the ColPali consensus/resolution |
| layers, which are minutes-per-page on CPU (see pipeline.py) and default off. |
| """ |
| from endopath import pipeline |
|
|
| result = get_case(conn, case_barcode) |
| if result is None: |
| raise HTTPException(status_code=404, detail="case not found") |
| case, _ = result |
| if case.status != CaseStatus.QUEUED: |
| raise HTTPException( |
| status_code=409, |
| detail=f"case is {case.status.value}, not queued; only queued cases can be processed", |
| ) |
|
|
| try: |
| pipeline.process_case_in_store(conn, case_barcode, use_visual=use_visual) |
| except pipeline.NoReportTextError: |
| raise HTTPException(status_code=422, detail="no OCR report text available for this case") |
| except RuntimeError as exc: |
| raise HTTPException(status_code=503, detail=str(exc)) |
|
|
| return case_detail(case_barcode, conn) |
|
|
|
|
| def _maybe_advance_case(conn: sqlite3.Connection, case_barcode: str) -> None: |
| """Section 7.2: confirming the last open field should move a case |
| toward Confirmed, not require a separate manual action. |
| """ |
| result = get_case(conn, case_barcode) |
| if result is None: |
| return |
| case, patient_filename = result |
| changed = False |
| |
| |
| |
| |
| if case.status == CaseStatus.READY_FOR_REVIEW: |
| case.transition_to(CaseStatus.IN_REVIEW) |
| changed = True |
| if case.status == CaseStatus.IN_REVIEW and case.is_ready_for_export(): |
| case.transition_to(CaseStatus.CONFIRMED) |
| changed = True |
| if changed: |
| upsert_case(conn, case, patient_filename) |
|
|
|
|
| class FieldDecisionRequest(BaseModel): |
| """The reviewer-session context recorded with a confirm on the case-review surface (#98, #103): the |
| confirmer and their role drawn from the session, an optional note, how the answer was reached, and the |
| question and options the decision weighed. All optional, so a bare confirm still works; a row is |
| written to the decision ledger only when the confirmer and role are present.""" |
|
|
| confirmer: Optional[str] = None |
| role: Optional[str] = None |
| note: Optional[str] = None |
| decided_via: Optional[str] = None |
| question: Optional[str] = None |
| options: Optional[list] = None |
|
|
|
|
| def _record_field_decision( |
| conn: sqlite3.Connection, |
| *, |
| case_barcode: str, |
| field_name: str, |
| field: Any, |
| default_decided_via: str, |
| ctx: Optional[FieldDecisionRequest], |
| ) -> None: |
| """Write a decision-ledger row for a confirm or edit, but only when the session identity is present. |
| A bare confirm with no identity (the backward-compatible path) performs the state change without a |
| ledger row; the confirmation surface always sends the session identity, so it always files one.""" |
| if ctx is None or not (ctx.confirmer and ctx.confirmer.strip() and ctx.role and ctx.role.strip()): |
| return |
| from endopath import case_review as case_review_mod |
|
|
| case_review_mod.record_decision( |
| conn, |
| case_barcode=case_barcode, |
| field=field, |
| field_name=field_name, |
| confirmer=ctx.confirmer, |
| role=ctx.role, |
| note=ctx.note, |
| disposition="confirmed", |
| decided_via=ctx.decided_via or default_decided_via, |
| question=ctx.question, |
| options=ctx.options, |
| ) |
|
|
|
|
| def _with_identity( |
| ctx: Optional[FieldDecisionRequest], identity: Optional[Identity] |
| ) -> Optional[FieldDecisionRequest]: |
| """Force a field decision's confirmer and role to the authenticated account's, so an authenticated |
| confirm files a ledger row under the real account even when the client sent no identity (#128). On the |
| open dev path (no identity) the client's context passes through unchanged.""" |
| if identity is None: |
| return ctx |
| base = ctx if ctx is not None else FieldDecisionRequest() |
| return base.model_copy(update={"confirmer": identity.confirmer, "role": identity.role}) |
|
|
|
|
| @api.post("/cases/{case_barcode}/fields/{field_name}/confirm", dependencies=authenticated) |
| def confirm_field( |
| case_barcode: str, |
| field_name: str, |
| body: Optional[FieldDecisionRequest] = None, |
| conn: sqlite3.Connection = Depends(get_db), |
| identity: Optional[Identity] = Depends(authorize), |
| ) -> dict: |
| result = get_case(conn, case_barcode) |
| if result is None: |
| raise HTTPException(status_code=404, detail="case not found") |
| case, patient_filename = result |
| if field_name not in case.checklist.all_fields(): |
| raise HTTPException(status_code=404, detail="field not found") |
|
|
| field = getattr(case.checklist, field_name) |
| field.confirm() |
| setattr(case.checklist, field_name, field) |
|
|
| upsert_case(conn, case, patient_filename) |
| _record_field_decision( |
| conn, |
| case_barcode=case_barcode, |
| field_name=field_name, |
| field=field, |
| default_decided_via="accept", |
| ctx=_with_identity(body, identity), |
| ) |
| _maybe_advance_case(conn, case_barcode) |
|
|
| return case_detail(case_barcode, conn) |
|
|
|
|
| class EditFieldRequest(BaseModel): |
| value: Any = None |
| evidence: Optional[EvidenceSpan] = None |
| |
| decision: Optional[FieldDecisionRequest] = None |
|
|
|
|
| @api.post("/cases/{case_barcode}/fields/{field_name}/edit", dependencies=authenticated) |
| def edit_field( |
| case_barcode: str, |
| field_name: str, |
| body: EditFieldRequest, |
| conn: sqlite3.Connection = Depends(get_db), |
| identity: Optional[Identity] = Depends(authorize), |
| ) -> dict: |
| result = get_case(conn, case_barcode) |
| if result is None: |
| raise HTTPException(status_code=404, detail="case not found") |
| case, patient_filename = result |
| if field_name not in case.checklist.all_fields(): |
| raise HTTPException(status_code=404, detail="field not found") |
|
|
| field = getattr(case.checklist, field_name) |
| field.edit(body.value, body.evidence) |
| setattr(case.checklist, field_name, field) |
|
|
| |
| |
| |
| |
| |
| |
| |
| try: |
| EndometrialChecklist.model_validate(case.checklist.model_dump()) |
| except ValidationError as exc: |
| raise HTTPException(status_code=422, detail=f"invalid value for {field_name}: {exc}") from exc |
|
|
| |
| |
| |
| |
| |
| case.checklist.apply_applicability_rules(case.staging_edition) |
|
|
| upsert_case(conn, case, patient_filename) |
| _record_field_decision( |
| conn, |
| case_barcode=case_barcode, |
| field_name=field_name, |
| field=field, |
| default_decided_via="edit", |
| ctx=_with_identity(body.decision, identity), |
| ) |
| _maybe_advance_case(conn, case_barcode) |
|
|
| return case_detail(case_barcode, conn) |
|
|
|
|
| @api.post("/cases/{case_barcode}/fields/{field_name}/escalate", dependencies=authenticated) |
| def escalate_field( |
| case_barcode: str, |
| field_name: str, |
| body: Optional[FieldDecisionRequest] = None, |
| conn: sqlite3.Connection = Depends(get_db), |
| identity: Optional[Identity] = Depends(authorize), |
| ) -> dict: |
| """Route a drafted field to the licensed reviewer's queue (#119, #161). Records an 'escalated' |
| disposition in the decision ledger without confirming the value: the field stays in review, and the |
| escalation queue reads these rows. Resolving it is confirming or editing the field, which upserts the |
| ledger row to a settled disposition and drops it from the queue.""" |
| result = get_case(conn, case_barcode) |
| if result is None: |
| raise HTTPException(status_code=404, detail="case not found") |
| case, _ = result |
| if field_name not in case.checklist.all_fields(): |
| raise HTTPException(status_code=404, detail="field not found") |
|
|
| field = getattr(case.checklist, field_name) |
| ctx = _with_identity(body, identity) |
| if ctx is not None and ctx.confirmer and ctx.role: |
| from endopath import case_review as case_review_mod |
|
|
| case_review_mod.record_decision( |
| conn, |
| case_barcode=case_barcode, |
| field=field, |
| field_name=field_name, |
| confirmer=ctx.confirmer, |
| role=ctx.role, |
| note=ctx.note, |
| disposition="escalated", |
| decided_via="escalate", |
| question=ctx.question, |
| options=ctx.options, |
| ) |
| return case_detail(case_barcode, conn) |
|
|
|
|
| @api.get("/escalations") |
| def list_escalations( |
| conn: sqlite3.Connection = Depends(get_db), |
| project: Optional[str] = None, |
| ) -> dict: |
| """The escalation queue (#161): fields a reviewer routed to the licensed account, across cases. Reads |
| the decision ledger for the escalated disposition; a later confirm or edit upserts the row to a settled |
| disposition, so it leaves the queue. Scoped to a project when one is given.""" |
| if project: |
| rows = conn.execute( |
| text( |
| "SELECT fc.case_barcode, fc.field_name, fc.value_json, fc.confirmed_by, fc.role, " |
| "fc.note, fc.confirmed_at FROM field_confirmations fc " |
| "JOIN cases c ON c.case_barcode = fc.case_barcode " |
| "WHERE fc.disposition = 'escalated' AND c.project_id = :project " |
| "ORDER BY fc.confirmed_at DESC" |
| ), |
| {"project": project}, |
| ) |
| else: |
| rows = conn.execute( |
| text( |
| "SELECT case_barcode, field_name, value_json, confirmed_by, role, note, confirmed_at " |
| "FROM field_confirmations WHERE disposition = 'escalated' ORDER BY confirmed_at DESC" |
| ) |
| ) |
| escalations = [] |
| for r in rows.mappings(): |
| value: Any = None |
| if r["value_json"]: |
| try: |
| value = json.loads(r["value_json"]) |
| except (ValueError, TypeError): |
| value = r["value_json"] |
| escalations.append( |
| { |
| "case_barcode": r["case_barcode"], |
| "field_name": r["field_name"], |
| "value": value, |
| "escalated_by": r["confirmed_by"], |
| "role": r["role"], |
| "note": r["note"], |
| "escalated_at": r["confirmed_at"], |
| } |
| ) |
| return {"escalations": escalations} |
|
|
|
|
| @api.post("/demo/reset") |
| def demo_reset( |
| conn: sqlite3.Connection = Depends(get_db), |
| case_barcode: Optional[str] = None, |
| field_name: str = "molecular_classification", |
| seed_escalation: bool = False, |
| clear_longitudinal: bool = False, |
| ) -> dict: |
| """Clear one beat of the integrated demo's slice so a run starts fresh. Called three ways, and none |
| must step on each other in the full walkthrough: |
| |
| - `clear_longitudinal` with a case: drop that case's GDC linkage joins (#180), so the annotator's |
| patient-record beat starts from an unconfirmed panel and its Confirm-linkage click has something to |
| do on every run. It touches only `longitudinal_joins`, leaving the case's field decisions and the |
| mapping resolution untouched, and returns without running the field or mapping reset below. |
| - No case: clear the data manager's molecular mapping resolution (and its split components). The |
| mapping and the case-field beats live in different tables, so this is the mapping-only reset. |
| - A case: clear that case's decision on the one escalated field (the annotator's escalation and the |
| pathologist's confirmation), and nothing else, so the case's other confirmed values and the data |
| manager's earlier mapping resolution both survive. |
| |
| With `seed_escalation`, it then re-routes that field to the escalation queue as the annotator, so the |
| pathologist story finds an escalation to resolve on its own, without the annotator story running |
| first. Demo-only and idempotent; it touches only the rows the demo writes, on a public demo corpus, |
| so it stays open like the reads.""" |
| if clear_longitudinal: |
| if case_barcode: |
| conn.execute( |
| text("DELETE FROM longitudinal_joins WHERE case_barcode = :cb"), |
| {"cb": case_barcode}, |
| ) |
| conn.commit() |
| return {"reset": True, "case_barcode": case_barcode, "cleared": "longitudinal_joins"} |
|
|
| if case_barcode: |
| conn.execute( |
| text( |
| "DELETE FROM field_confirmations WHERE case_barcode = :cb AND field_name = :fn" |
| ), |
| {"cb": case_barcode, "fn": field_name}, |
| ) |
| else: |
| conn.execute( |
| text( |
| "DELETE FROM mapping_confirmation_components WHERE mapping_id IN " |
| "(SELECT mapping_id FROM mapping_confirmations WHERE LOWER(source_field) LIKE '%molecular%')" |
| ) |
| ) |
| conn.execute( |
| text("DELETE FROM mapping_confirmations WHERE LOWER(source_field) LIKE '%molecular%'") |
| ) |
| conn.commit() |
|
|
| if case_barcode and seed_escalation: |
| from endopath import case_review as case_review_mod |
|
|
| seed_notes = { |
| "lymphovascular_space_invasion": "Report reads 'Identified' but not focal vs substantial; the FIGO stage turns on it.", |
| "molecular_classification": "Molecular class needs POLE, MMR, and p53 the report never assays.", |
| } |
| result = get_case(conn, case_barcode) |
| if result is not None: |
| case, _ = result |
| if field_name in case.checklist.all_fields(): |
| case_review_mod.record_decision( |
| conn, |
| case_barcode=case_barcode, |
| field=getattr(case.checklist, field_name), |
| field_name=field_name, |
| confirmer="Annotator", |
| role="reviewer", |
| note=seed_notes.get(field_name, "Needs a pathologist's read."), |
| disposition="escalated", |
| decided_via="escalate", |
| question=None, |
| options=None, |
| ) |
|
|
| return {"reset": True, "case_barcode": case_barcode, "field_name": field_name} |
|
|
|
|
| class FieldChatTurn(BaseModel): |
| """One prior turn of a value chat: the speaker (`user` or `assistant`) and the text, so a follow-up |
| question carries its context.""" |
|
|
| role: str |
| text: str |
|
|
|
|
| class FieldChatRequest(BaseModel): |
| """A reviewer's question about one checklist value (#98), with the prior turns of that value's |
| conversation.""" |
|
|
| question: str |
| history: list[FieldChatTurn] = Field(default_factory=list) |
|
|
|
|
| @api.post("/cases/{case_barcode}/fields/{field_name}/chat", dependencies=authenticated) |
| def field_chat( |
| case_barcode: str, |
| field_name: str, |
| body: FieldChatRequest, |
| conn: sqlite3.Connection = Depends(get_db), |
| ) -> dict: |
| """Answer a reviewer's question about one extracted value (#98): the engine's reasoning, the evidence |
| it read, and, when the question points at a correction, the value to file. The exchange resolves into |
| a filed decision when the reviewer accepts the suggestion through the edit endpoint with |
| `decided_via='chat'`. |
| |
| Gated by the shared token: the call spends the Anthropic key. The chat model is named explicitly in |
| `case_review` (claude-sonnet-5), behind a provider-agnostic seam a fake substitutes for in the fast |
| test suite. |
| """ |
| from endopath import case_review as case_review_mod |
|
|
| result = get_case(conn, case_barcode) |
| if result is None: |
| raise HTTPException(status_code=404, detail="case not found") |
| case, _ = result |
| if field_name not in case.checklist.all_fields(): |
| raise HTTPException(status_code=404, detail="field not found") |
|
|
| try: |
| return case_review_mod.explain_field( |
| conn, |
| case=case, |
| field_name=field_name, |
| question=body.question, |
| history=[turn.model_dump() for turn in body.history], |
| ) |
| except case_review_mod.CaseReviewError as exc: |
| raise HTTPException(status_code=422, detail=str(exc)) from exc |
| except RuntimeError as exc: |
| raise HTTPException(status_code=503, detail=str(exc)) from exc |
|
|
|
|
| |
| |
| |
| |
| |
| |
| |
| |
| |
| _page_embedding_cache: dict[str, object] = {} |
|
|
|
|
| def _get_case_page_embeddings(patient_filename: str): |
| if patient_filename not in _page_embedding_cache: |
| image_paths = colpali_retrieval.case_page_images(patient_filename, image_dir=IMAGE_DIR) |
| if not image_paths: |
| raise HTTPException(status_code=404, detail="no page images ingested for this case") |
| _page_embedding_cache[patient_filename] = colpali_retrieval.embed_pages(image_paths) |
| return _page_embedding_cache[patient_filename] |
|
|
|
|
| |
| |
| |
| _EMBEDDINGS_DIR = Path(os.environ.get("ENDOPATH_EMBEDDINGS_DIR", str(precomputed_retrieval.EMBEDDINGS_DIR))) |
| _query_embeddings_cache: Optional[dict] = None |
| _query_embeddings_loaded = False |
|
|
|
|
| def _get_query_embeddings(): |
| """The nine precomputed query vectors, loaded once. None when they aren't |
| present (then visual retrieval falls back to embedding live).""" |
| global _query_embeddings_cache, _query_embeddings_loaded |
| if not _query_embeddings_loaded: |
| path = _EMBEDDINGS_DIR / precomputed_retrieval.QUERY_EMBEDDINGS_NAME |
| _query_embeddings_cache = ( |
| precomputed_retrieval.load_query_embeddings(_EMBEDDINGS_DIR) if path.exists() else None |
| ) |
| _query_embeddings_loaded = True |
| return _query_embeddings_cache |
|
|
|
|
| def _visual_rank(patient_filename: str, field_name: str) -> tuple[int, float]: |
| """(1-indexed page, score) for a field's visual retrieval. Prefers |
| precomputed vectors scored on CPU (deployed path, no torch/ColPali); falls |
| back to embedding the case's pages live when they aren't precomputed.""" |
| queries = _get_query_embeddings() |
| if ( |
| queries is not None |
| and field_name in queries |
| and precomputed_retrieval.has_case_embeddings(patient_filename, _EMBEDDINGS_DIR) |
| ): |
| pages = precomputed_retrieval.load_case_page_embeddings(patient_filename, _EMBEDDINGS_DIR) |
| page_number, score, _ = precomputed_retrieval.rank_pages_with_margin(queries[field_name], pages) |
| return page_number, float(score) |
|
|
| if not colpali_retrieval.is_available(): |
| raise HTTPException( |
| status_code=503, |
| detail="no precomputed embeddings for this case and colpali_engine/torch not installed", |
| ) |
| page_embeddings = _get_case_page_embeddings(patient_filename) |
| page_number, score = colpali_retrieval.rank_pages( |
| colpali_retrieval.FIELD_QUERIES[field_name], page_embeddings |
| ) |
| return page_number, float(score) |
|
|
|
|
| @api.get("/cases/{case_barcode}/fields/{field_name}/visual_evidence", dependencies=authenticated) |
| def visual_evidence(case_barcode: str, field_name: str, conn: sqlite3.Connection = Depends(get_db)) -> dict: |
| """Visual retrieval for one checklist field, returning an EvidenceSpan the |
| same shape llm_extraction.py produces. Scores precomputed page vectors |
| against precomputed query vectors on CPU (issue #21, no model loaded), and |
| falls back to embedding live with ColPali when a case isn't precomputed. |
| Stateless. Issue #18: raw scores aren't comparable across different fields. |
| """ |
| if field_name not in colpali_retrieval.FIELD_QUERIES: |
| raise HTTPException(status_code=404, detail="field not found") |
|
|
| result = get_case(conn, case_barcode) |
| if result is None: |
| raise HTTPException(status_code=404, detail="case not found") |
| _, patient_filename = result |
| if patient_filename is None: |
| raise HTTPException(status_code=404, detail="no page images ingested for this case") |
|
|
| page_number, score = _visual_rank(patient_filename, field_name) |
| evidence = EvidenceSpan(quote="", page_number=page_number, source="visual") |
| return {"evidence": evidence.model_dump(), "score": score} |
|
|
|
|
| @api.get("/cases/{case_barcode}/pages/{page_number}") |
| def get_page_image(case_barcode: str, page_number: int, conn: sqlite3.Connection = Depends(get_db)): |
| result = get_case(conn, case_barcode) |
| if result is None: |
| raise HTTPException(status_code=404, detail="case not found") |
| _, patient_filename = result |
| if patient_filename is None: |
| raise HTTPException(status_code=404, detail="no page images ingested for this case") |
|
|
| image_path = IMAGE_DIR / patient_filename / f"page_{page_number:03d}.png" |
| if not image_path.exists(): |
| raise HTTPException(status_code=404, detail="page image not found") |
| return FileResponse(image_path) |
|
|
|
|
| |
| |
| |
|
|
| @api.get("/longitudinal/records/{case_barcode}") |
| def longitudinal_records(case_barcode: str, conn = Depends(get_db)) -> dict: |
| """List all longitudinal records for a participant (via case barcode).""" |
| |
| result = conn.execute( |
| text("SELECT case_barcode FROM cases WHERE case_barcode = :barcode"), |
| {"barcode": case_barcode} |
| ).fetchone() |
|
|
| if not result: |
| raise HTTPException(status_code=404, detail="case not found") |
|
|
| |
| participant_id = case_barcode[:12] if len(case_barcode) >= 12 else case_barcode |
|
|
| |
| rows = conn.execute( |
| text(""" |
| SELECT record_id, participant_id, case_barcode, record_type, record_index, |
| days_to_followup, days_to_treatment, days_to_recurrence, |
| biopsy_result, followup_treatment_success, recurrence_type, therapy_type, |
| structured_data |
| FROM longitudinal_records |
| WHERE participant_id = :pid |
| ORDER BY record_type, record_index |
| """), |
| {"pid": participant_id}, |
| ).mappings().fetchall() |
|
|
| records = [] |
| for row in rows: |
| record_dict = dict(row) |
| if record_dict.get("structured_data"): |
| try: |
| record_dict["structured_data"] = json.loads(record_dict["structured_data"]) |
| except (json.JSONDecodeError, TypeError): |
| record_dict["structured_data"] = {} |
| records.append(record_dict) |
|
|
| return {"case_barcode": case_barcode, "participant_id": participant_id, "records": records} |
|
|
|
|
| @api.get("/longitudinal/joins/{case_barcode}") |
| def longitudinal_joins(case_barcode: str, conn = Depends(get_db)) -> dict: |
| """List all confirmed joins for a case.""" |
| rows = conn.execute( |
| text(""" |
| SELECT lj.join_id, lj.case_barcode, lj.record_id, lj.disposition, |
| lj.confirmed_by, lj.role, lj.confirmed_at, lj.note, lj.question, |
| lr.record_type, lr.days_to_followup, lr.therapy_type, lr.recurrence_type |
| FROM longitudinal_joins lj |
| JOIN longitudinal_records lr ON lj.record_id = lr.record_id |
| WHERE lj.case_barcode = :barcode |
| ORDER BY lj.confirmed_at DESC |
| """), |
| {"barcode": case_barcode}, |
| ).mappings().fetchall() |
|
|
| joins = [dict(row) for row in rows] |
| return {"case_barcode": case_barcode, "joins": joins} |
|
|
|
|
| @api.post("/longitudinal/join/{case_barcode}", dependencies=authenticated) |
| def confirm_join( |
| case_barcode: str, |
| request: JoinRequest, |
| identity: Optional[Identity] = Depends(authorize), |
| conn = Depends(get_db), |
| ) -> dict: |
| """Confirm or reject a join between a case and a longitudinal record.""" |
| join_id = str(uuid.uuid4()) |
| confirmed_at = datetime.now(timezone.utc).isoformat() |
|
|
| |
| result = conn.execute( |
| text("SELECT record_id FROM longitudinal_records WHERE record_id = :rid"), |
| {"rid": request.record_id} |
| ).fetchone() |
|
|
| if not result: |
| raise HTTPException(status_code=404, detail="record not found") |
|
|
| |
| existing = conn.execute( |
| text("SELECT join_id FROM longitudinal_joins WHERE case_barcode = :cb AND record_id = :rid"), |
| {"cb": case_barcode, "rid": request.record_id} |
| ).fetchone() |
|
|
| if existing: |
| |
| conn.execute( |
| text(""" |
| UPDATE longitudinal_joins |
| SET disposition = :disp, confirmed_by = :cb, role = :role, confirmed_at = :cat, note = :note, question = :q |
| WHERE join_id = :jid |
| """), |
| { |
| "disp": request.disposition, |
| "cb": identity.confirmer if identity else None, |
| "role": identity.role if identity else None, |
| "cat": confirmed_at, |
| "note": request.note, |
| "q": request.question, |
| "jid": existing[0], |
| } |
| ) |
| join_id = existing[0] |
| else: |
| |
| conn.execute( |
| text(""" |
| INSERT INTO longitudinal_joins ( |
| join_id, case_barcode, record_id, disposition, confirmed_by, role, confirmed_at, note, question |
| ) VALUES (:jid, :cb, :rid, :disp, :cb2, :role, :cat, :note, :q) |
| """), |
| { |
| "jid": join_id, |
| "cb": case_barcode, |
| "rid": request.record_id, |
| "disp": request.disposition, |
| "cb2": identity.confirmer if identity else None, |
| "role": identity.role if identity else None, |
| "cat": confirmed_at, |
| "note": request.note, |
| "q": request.question, |
| } |
| ) |
|
|
| conn.commit() |
|
|
| return { |
| "join_id": join_id, |
| "case_barcode": case_barcode, |
| "record_id": request.record_id, |
| "disposition": request.disposition, |
| "confirmed_by": identity.confirmer if identity else None, |
| "confirmed_at": confirmed_at, |
| } |
|
|
|
|
| @api.post("/longitudinal/confirm-linkage/{case_barcode}", dependencies=authenticated) |
| def confirm_linkage( |
| case_barcode: str, |
| request: LinkageRequest = LinkageRequest(), |
| identity: Optional[Identity] = Depends(authorize), |
| conn = Depends(get_db), |
| ) -> dict: |
| """Affirm in one action that every GDC longitudinal record joined to this case belongs to it (#180). |
| |
| Upserts a `confirmed` join for each of the participant's records, recording the reviewer's identity and |
| role, and leaves a record the reviewer has already flagged `rejected` untouched so a confirm-all does not |
| silently un-flag it. Idempotent: re-running re-stamps the confirmations without duplicating rows.""" |
| result = conn.execute( |
| text("SELECT case_barcode FROM cases WHERE case_barcode = :b"), {"b": case_barcode} |
| ).fetchone() |
| if not result: |
| raise HTTPException(status_code=404, detail="case not found") |
|
|
| participant_id = case_barcode[:12] if len(case_barcode) >= 12 else case_barcode |
| record_ids = [ |
| row[0] |
| for row in conn.execute( |
| text("SELECT record_id FROM longitudinal_records WHERE participant_id = :p"), |
| {"p": participant_id}, |
| ).fetchall() |
| ] |
| existing = { |
| row[0]: row[1] |
| for row in conn.execute( |
| text("SELECT record_id, disposition FROM longitudinal_joins WHERE case_barcode = :b"), |
| {"b": case_barcode}, |
| ).fetchall() |
| } |
| confirmed_at = datetime.now(timezone.utc).isoformat() |
| confirmer = identity.confirmer if identity else None |
| role = identity.role if identity else None |
|
|
| confirmed = 0 |
| for record_id in record_ids: |
| if existing.get(record_id) == "rejected": |
| continue |
| if record_id in existing: |
| conn.execute( |
| text( |
| "UPDATE longitudinal_joins SET disposition='confirmed', confirmed_by=:c, " |
| "role=:r, confirmed_at=:t, note=:n WHERE case_barcode=:b AND record_id=:rid" |
| ), |
| {"c": confirmer, "r": role, "t": confirmed_at, "n": request.note, "b": case_barcode, "rid": record_id}, |
| ) |
| else: |
| conn.execute( |
| text( |
| "INSERT INTO longitudinal_joins (join_id, case_barcode, record_id, disposition, " |
| "confirmed_by, role, confirmed_at, note, question) " |
| "VALUES (:j, :b, :rid, 'confirmed', :c, :r, :t, :n, NULL)" |
| ), |
| {"j": str(uuid.uuid4()), "b": case_barcode, "rid": record_id, "c": confirmer, "r": role, "t": confirmed_at, "n": request.note}, |
| ) |
| confirmed += 1 |
| conn.commit() |
|
|
| rejected = sum(1 for disposition in existing.values() if disposition == "rejected") |
| return { |
| "case_barcode": case_barcode, |
| "confirmed": confirmed, |
| "rejected": rejected, |
| "total_records": len(record_ids), |
| "confirmed_by": confirmer, |
| "role": role, |
| "confirmed_at": confirmed_at, |
| } |
|
|
|
|
| @api.get("/longitudinal/summary/{case_barcode}") |
| def longitudinal_summary(case_barcode: str, conn = Depends(get_db)) -> dict: |
| """Summarize longitudinal data availability for a case.""" |
| |
| result = conn.execute( |
| text("SELECT case_barcode FROM cases WHERE case_barcode = :barcode"), |
| {"barcode": case_barcode} |
| ).fetchone() |
|
|
| if not result: |
| raise HTTPException(status_code=404, detail="case not found") |
|
|
| |
| participant_id = case_barcode[:12] if len(case_barcode) >= 12 else case_barcode |
|
|
| |
| rows = conn.execute( |
| text(""" |
| SELECT record_type, COUNT(*) as count |
| FROM longitudinal_records |
| WHERE participant_id = :pid |
| GROUP BY record_type |
| """), |
| {"pid": participant_id} |
| ).fetchall() |
|
|
| records_by_type = {row[0]: row[1] for row in rows} |
|
|
| |
| rows = conn.execute( |
| text(""" |
| SELECT disposition, COUNT(*) as count |
| FROM longitudinal_joins |
| WHERE case_barcode = :barcode |
| GROUP BY disposition |
| """), |
| {"barcode": case_barcode} |
| ).fetchall() |
|
|
| joins_by_disposition = {row[0]: row[1] for row in rows} |
|
|
| return { |
| "case_barcode": case_barcode, |
| "participant_id": participant_id, |
| "records_available": records_by_type, |
| "joins_confirmed": joins_by_disposition, |
| } |
|
|
|
|
| @api.get("/longitudinal/cohort") |
| def longitudinal_cohort(project: Optional[str] = None, conn = Depends(get_db)) -> dict: |
| """Cohort-level longitudinal outcomes for the dashboard module (#181): recurrence, survival, and |
| adjuvant-treatment counts, plus median follow-up and survival times, over the cases the GDC join |
| reached. Projects each case through the same `project_outcomes` the export uses, so the dashboard and |
| the export never disagree. `project` scopes to one project's cases; omitted, every case is aggregated. |
| An open read.""" |
| from endopath import export as export_mod |
| from endopath import longitudinal |
|
|
| by_participant = export_mod._load_longitudinal(conn) |
| joins_by_case = export_mod._load_longitudinal_joins(conn) |
| if project is not None: |
| case_rows = conn.execute( |
| text("SELECT case_barcode FROM cases WHERE project_id = :p ORDER BY case_barcode"), |
| {"p": project}, |
| ).fetchall() |
| else: |
| case_rows = conn.execute( |
| text("SELECT case_barcode FROM cases ORDER BY case_barcode") |
| ).fetchall() |
|
|
| outcomes_per_case = [] |
| for (barcode,) in case_rows: |
| participant_id = barcode[:12] if len(barcode) >= 12 else barcode |
| records = by_participant.get(participant_id, []) |
| if not records: |
| continue |
| |
| |
| rejected = export_mod._rejected_record_ids(joins_by_case.get(barcode, [])) |
| if rejected: |
| records = [r for r in records if r.record_id not in rejected] |
| outcomes_per_case.append(longitudinal.project_outcomes(records)) |
|
|
| return longitudinal.aggregate_cohort(outcomes_per_case) |
|
|
|
|
| app.include_router(api) |
|
|
|
|
| def _static_file_within_dist(dist: Path, relative_path: str) -> Optional[Path]: |
| """Resolve `relative_path` against the build directory, or None if it isn't |
| a real file inside it. A path that names a directory holding an index.html |
| resolves to that index.html, so a bundled subfolder page (the case study at |
| /case-study/) serves without falling through to the SPA and rendering blank. |
| |
| The containment check is not decorative: `spa_path` is whatever the client |
| put in the URL, and a percent-encoded `..` survives the routing layer, so |
| resolving without checking would hand out arbitrary files from the host. |
| """ |
| if not relative_path: |
| return None |
| dist_root = dist.resolve() |
| candidate = (dist_root / relative_path).resolve() |
| if dist_root not in candidate.parents: |
| return None |
| if candidate.is_file(): |
| return candidate |
| index = candidate / "index.html" |
| if candidate.is_dir() and index.is_file(): |
| return index |
| return None |
|
|
|
|
| @app.get("/{spa_path:path}", include_in_schema=False) |
| def serve_frontend(spa_path: str): |
| """Serve the built SPA: real files (hashed JS/CSS, favicon) as themselves, |
| every other path as index.html so client-side routes survive a hard refresh |
| or a deep link. |
| |
| Registered last, so it only sees paths no API route claimed. Unmatched |
| `/api/...` paths are excluded explicitly -- an API client that typos a URL |
| should get a JSON 404, not a 200 and a page of HTML. |
| """ |
| if spa_path == "api" or spa_path.startswith("api/"): |
| raise HTTPException(status_code=404, detail="not found") |
|
|
| index = _frontend_dist / "index.html" |
| if not index.is_file(): |
| raise HTTPException( |
| status_code=404, |
| detail=f"frontend build not found at {_frontend_dist}; run `npm run build` in frontend/", |
| ) |
|
|
| static_file = _static_file_within_dist(_frontend_dist, spa_path) |
| if static_file is None: |
| return FileResponse(index) |
| |
| |
| |
| |
| if ( |
| static_file.name == "index.html" |
| and static_file != index |
| and not spa_path.endswith("/") |
| and not spa_path.endswith("index.html") |
| ): |
| return RedirectResponse(url=f"/{spa_path}/", status_code=308) |
| return FileResponse(static_file) |
|
|