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:p@host.docker.internal: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:p@host.docker.internal: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