A table with filters made the user do the work. Rare disease triage is a different task:
which few variants could explain *this* patient's phenotype, and why. The app now answers
that, and lets a reviewer act on the answer.
Domain
- a case is a proband: a VCF plus the HPO terms observed in the patient (samples -> cases)
- HPO's gene-to-phenotype annotations are loaded as reference data (scripts/load-hpo.py)
- each candidate can be shortlisted or dismissed with a reason and a note
Ranking (app/services/triage.py, 21 tests)
- weighted sum of phenotype match, rarity, consequence severity and the model's score,
with every component shown next to the candidate
- rarity and consequence filter; phenotype only ranks, because a real diagnosis can sit in
a gene nobody has annotated yet and filtering on it would hide exactly that case
- ClinVar is deliberately not an input: it appears beside the result as independent
confirmation, so nothing ranks highly merely because ClinVar already said pathogenic
UI
- the funnel is the headline: variants called -> rare -> coding candidates -> phenotype-matched
- ranked candidates with evidence chips, not a grid of everything; filters are demoted
- a variant panel showing the score breakdown, the matched HPO terms, the raw VEP record and
links out to Ensembl/gnomAD/ClinVar, with the decision controls
- a printable case report: phenotype, funnel, shortlisted variants with reasons, provenance
API: /cases with phenotypes, /cases/{id}/candidates (funnel + ranked + weights),
/variants/{id}, /variants/{id}/decision, /cases/{id}/report, /phenotypes for the picker.
Scoring moved under the case and now answers 503 with the reason when no model registry is
reachable, instead of a 500.
Verified end to end on a simulated proband (scripts/make-demo-case.sh: real GIAB HG002
background + one real ClinVar 2-star pathogenic NF2 variant). 13 variants called -> 1 coding
candidate, and the planted variant ranks first at 0.80 on phenotype 1.00, rarity 1.00 and
consequence 1.00, with ClinVar agreeing afterwards.
Tests: api 75, ml 18, loader 16, web 27; ruff, mypy, svelte-check, terraform validate, both
kustomize overlays and the Nextflow stub run all clean.
213 lines
7.6 KiB
Python
213 lines
7.6 KiB
Python
import uuid
|
|
from typing import Any
|
|
|
|
import pytest
|
|
from httpx import AsyncClient
|
|
|
|
from app.config import settings
|
|
from app.services import events
|
|
|
|
|
|
async def new_case(client: AsyncClient) -> str:
|
|
r = await client.post(
|
|
"/api/cases",
|
|
json={"name": f"s-{uuid.uuid4()}", "vcf_uri": "gs://bucket/x.vcf.gz", "assembly": "GRCh37"},
|
|
)
|
|
assert r.status_code == 201, r.text
|
|
return r.json()["id"]
|
|
|
|
|
|
@pytest.mark.usefixtures("db")
|
|
async def test_local_without_nextflow_fails_fast_with_instructions(
|
|
client: AsyncClient, monkeypatch: pytest.MonkeyPatch
|
|
) -> None:
|
|
monkeypatch.setattr(settings, "pubsub_topic", None)
|
|
monkeypatch.setattr(events.shutil, "which", lambda _: None)
|
|
case_id = await new_case(client)
|
|
|
|
r = await client.post(f"/api/cases/{case_id}/annotate")
|
|
assert r.status_code == 202
|
|
job = r.json()
|
|
assert job["status"] == "failed"
|
|
assert "make annotate" in job["log"]
|
|
assert job["finished_at"] is not None
|
|
|
|
|
|
class FakeStdout:
|
|
"""Stands in for asyncio's StreamReader: async-iterates the process output."""
|
|
|
|
def __init__(self, lines: list[bytes]) -> None:
|
|
self._lines = lines
|
|
|
|
def __aiter__(self) -> "FakeStdout":
|
|
return self
|
|
|
|
async def __anext__(self) -> bytes:
|
|
if not self._lines:
|
|
raise StopAsyncIteration
|
|
return self._lines.pop(0)
|
|
|
|
|
|
class FakeProcess:
|
|
def __init__(self, returncode: int, output: bytes = b"", lines: list[bytes] | None = None) -> None:
|
|
self.returncode = returncode
|
|
self.pid = 4242
|
|
self.stdout = FakeStdout(lines if lines is not None else ([output] if output else []))
|
|
|
|
async def wait(self) -> int:
|
|
return self.returncode
|
|
|
|
|
|
@pytest.mark.usefixtures("db")
|
|
async def test_local_run_is_watched_and_a_crash_marks_the_job_failed(
|
|
client: AsyncClient, monkeypatch: pytest.MonkeyPatch
|
|
) -> None:
|
|
monkeypatch.setattr(settings, "pubsub_topic", None)
|
|
monkeypatch.setattr(events.shutil, "which", lambda _: "/usr/bin/nextflow")
|
|
launched: dict[str, Any] = {}
|
|
|
|
async def fake_exec(*cmd: str, **kw: Any) -> FakeProcess:
|
|
launched["cmd"], launched["env"] = cmd, kw["env"]
|
|
return FakeProcess(1, b"ERROR ~ VEP cache not found")
|
|
|
|
monkeypatch.setattr(events.asyncio, "create_subprocess_exec", fake_exec)
|
|
case_id = await new_case(client)
|
|
|
|
r = await client.post(f"/api/cases/{case_id}/annotate")
|
|
assert r.status_code == 202
|
|
await events.drain()
|
|
|
|
job = (await client.get(f"/api/jobs/{r.json()['id']}")).json()
|
|
assert job["status"] == "failed"
|
|
assert "VEP cache not found" in job["log"]
|
|
# The DB password travels in the environment, never on the command line.
|
|
assert settings.database_url not in " ".join(launched["cmd"])
|
|
assert launched["env"]["DATABASE_URL"] == settings.database_url
|
|
assert launched["cmd"][launched["cmd"].index("--assembly") + 1] == "GRCh37"
|
|
|
|
|
|
class FakeFuture:
|
|
def __init__(self, result: str | Exception) -> None:
|
|
self._result = result
|
|
|
|
def result(self, timeout: float | None = None) -> str:
|
|
if isinstance(self._result, Exception):
|
|
raise self._result
|
|
return self._result
|
|
|
|
|
|
class FakePublisher:
|
|
def __init__(self, result: str | Exception) -> None:
|
|
self.result = result
|
|
self.published: list[tuple[str, bytes]] = []
|
|
|
|
def topic_path(self, project: str, topic: str) -> str:
|
|
return f"projects/{project}/topics/{topic}"
|
|
|
|
def publish(self, topic: str, data: bytes) -> FakeFuture:
|
|
self.published.append((topic, data))
|
|
return FakeFuture(self.result)
|
|
|
|
|
|
@pytest.mark.usefixtures("db")
|
|
async def test_pubsub_publishes_to_the_full_topic_path(
|
|
client: AsyncClient, monkeypatch: pytest.MonkeyPatch
|
|
) -> None:
|
|
publisher = FakePublisher("msg-123")
|
|
monkeypatch.setattr(settings, "pubsub_topic", "vcf-uploaded")
|
|
monkeypatch.setattr(settings, "gcp_project", "my-proj")
|
|
monkeypatch.setattr(events, "_publisher", lambda: publisher)
|
|
case_id = await new_case(client)
|
|
|
|
r = await client.post(f"/api/cases/{case_id}/annotate")
|
|
job = r.json()
|
|
assert (job["status"], job["workflow_ref"]) == ("running", "pubsub:msg-123")
|
|
[(topic, data)] = publisher.published
|
|
assert topic == "projects/my-proj/topics/vcf-uploaded"
|
|
assert b'"assembly": "GRCh37"' in data and job["id"].encode() in data
|
|
|
|
|
|
@pytest.mark.usefixtures("db")
|
|
async def test_pubsub_failure_marks_the_job_failed(
|
|
client: AsyncClient, monkeypatch: pytest.MonkeyPatch
|
|
) -> None:
|
|
monkeypatch.setattr(settings, "pubsub_topic", "vcf-uploaded")
|
|
monkeypatch.setattr(settings, "gcp_project", "my-proj")
|
|
monkeypatch.setattr(events, "_publisher", lambda: FakePublisher(RuntimeError("403 denied")))
|
|
case_id = await new_case(client)
|
|
|
|
job = (await client.post(f"/api/cases/{case_id}/annotate")).json()
|
|
assert job["status"] == "failed"
|
|
assert "403 denied" in job["log"]
|
|
|
|
|
|
@pytest.mark.usefixtures("db")
|
|
async def test_the_pipeline_gets_its_own_database_url(
|
|
client: AsyncClient, monkeypatch: pytest.MonkeyPatch
|
|
) -> None:
|
|
"""The loader runs in a container, where the API's own localhost URL would point at itself."""
|
|
monkeypatch.setattr(settings, "pubsub_topic", None)
|
|
monkeypatch.setattr(settings, "cloudrun_job", None)
|
|
monkeypatch.setattr(
|
|
settings, "pipeline_database_url", "postgresql+asyncpg://u:[email protected]:5432/db"
|
|
)
|
|
monkeypatch.setattr(events.shutil, "which", lambda _: "/usr/bin/nextflow")
|
|
launched: dict[str, Any] = {}
|
|
|
|
async def fake_exec(*cmd: str, **kw: Any) -> FakeProcess:
|
|
launched["env"] = kw["env"]
|
|
return FakeProcess(0, b"")
|
|
|
|
monkeypatch.setattr(events.asyncio, "create_subprocess_exec", fake_exec)
|
|
case_id = await new_case(client)
|
|
|
|
await client.post(f"/api/cases/{case_id}/annotate")
|
|
await events.drain()
|
|
assert launched["env"]["DATABASE_URL"] == "postgresql+asyncpg://u:[email protected]:5432/db"
|
|
|
|
|
|
@pytest.mark.usefixtures("db")
|
|
async def test_progress_is_recorded_while_the_pipeline_runs(
|
|
client: AsyncClient, monkeypatch: pytest.MonkeyPatch
|
|
) -> None:
|
|
"""The UI needs something better than "running" for two minutes."""
|
|
monkeypatch.setattr(settings, "pubsub_topic", None)
|
|
monkeypatch.setattr(settings, "cloudrun_job", None)
|
|
monkeypatch.setattr(events.shutil, "which", lambda _: "/usr/bin/nextflow")
|
|
|
|
async def fake_exec(*cmd: str, **kw: Any) -> FakeProcess:
|
|
return FakeProcess(0, lines=[
|
|
b"[PIPELINE] main.nf | profile=docker\n",
|
|
b"[PROCESS f2/5c13bd] NORMALISE (tiny)\n",
|
|
b"[PROCESS cc/d62082] VEP (tiny)\n",
|
|
])
|
|
|
|
monkeypatch.setattr(events.asyncio, "create_subprocess_exec", fake_exec)
|
|
case_id = await new_case(client)
|
|
|
|
r = await client.post(f"/api/cases/{case_id}/annotate")
|
|
await events.drain()
|
|
|
|
job = (await client.get(f"/api/jobs/{r.json()['id']}")).json()
|
|
assert "VEP (tiny)" in job["log"]
|
|
assert job["status"] == "running" # the loader is what marks it succeeded
|
|
|
|
|
|
@pytest.mark.usefixtures("db")
|
|
async def test_progress_never_overwrites_a_finished_job(client: AsyncClient) -> None:
|
|
from app.db import SessionLocal
|
|
from app.models import Case, Job, JobStatus
|
|
|
|
async with SessionLocal() as s:
|
|
case = Case(name=f"c-{uuid.uuid4()}", vcf_uri="gs://bucket/x.vcf.gz", assembly="GRCh38")
|
|
job = Job(case=case, status=JobStatus.succeeded)
|
|
s.add_all([case, job])
|
|
await s.commit()
|
|
job_id = job.id
|
|
|
|
await events._record_progress(job_id, "[PROCESS 1/a] VEP (late)")
|
|
|
|
job_json = (await client.get(f"/api/jobs/{job_id}")).json()
|
|
assert job_json["status"] == "succeeded"
|
|
assert job_json["log"] is None
|