feat(web): show what the pipeline is doing while a job runs
"running" for two and a half minutes tells the user nothing. The job page now shows a spinner, the elapsed time, and the pipeline step Nextflow is actually on. - the API streams the Nextflow output into jobs.log as it arrives, instead of keeping it only when the run dies. Writes are throttled to one every 3s, or immediately when a new process starts, and are skipped for a job that has already finished, so a late line cannot overwrite the loader's result. - web/src/lib/progress.ts formats the elapsed time and picks the latest [PROCESS] line. No percentage: the pipeline cannot honestly estimate one. - the spinner animates only under prefers-reduced-motion: no-preference. Verified on a live run: the page showed "VEP (tiny)" for the duration, then the variant table replaced it on success. Tests: api 53, web 20; ruff, mypy, svelte-check clean.
This commit is contained in:
@@ -12,11 +12,15 @@ import json
|
||||
import logging
|
||||
import os
|
||||
import shutil
|
||||
import time
|
||||
import uuid
|
||||
from collections import deque
|
||||
from datetime import UTC, datetime
|
||||
from functools import cache
|
||||
from typing import Any
|
||||
|
||||
from sqlalchemy import update
|
||||
|
||||
from app.config import settings
|
||||
from app.db import SessionLocal
|
||||
from app.models import Job, JobStatus
|
||||
@@ -24,7 +28,8 @@ from app.models import Job, JobStatus
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
PIPELINE_ENTRYPOINT = "/pipeline/main.nf"
|
||||
LOG_TAIL_CHARS = 4000
|
||||
PROGRESS_LINES = 12 # rolling window of pipeline output kept on the job
|
||||
PROGRESS_INTERVAL_S = 3 # between writes, unless a new process starts
|
||||
PUBLISH_TIMEOUT_S = 30
|
||||
_watchers: set[asyncio.Task[None]] = set()
|
||||
|
||||
@@ -129,14 +134,36 @@ async def _run_local(job_id: uuid.UUID, vcf_uri: str, assembly: str) -> str:
|
||||
|
||||
|
||||
async def _watch(job_id: uuid.UUID, proc: Any) -> None:
|
||||
out, err = await proc.communicate()
|
||||
if proc.returncode == 0:
|
||||
"""Follow the pipeline: record progress while it runs, and the reason if it dies."""
|
||||
tail: deque[str] = deque(maxlen=PROGRESS_LINES)
|
||||
last_write = 0.0
|
||||
if proc.stdout is not None: # stderr is merged into stdout
|
||||
async for raw in proc.stdout:
|
||||
line = raw.decode(errors="replace").rstrip()
|
||||
if not line:
|
||||
continue
|
||||
tail.append(line)
|
||||
now = time.monotonic()
|
||||
# A new pipeline step is worth showing at once; otherwise throttle the writes.
|
||||
if line.startswith("[PROCESS") or now - last_write >= PROGRESS_INTERVAL_S:
|
||||
last_write = now
|
||||
await _record_progress(job_id, "\n".join(tail))
|
||||
returncode = await proc.wait()
|
||||
if returncode == 0:
|
||||
return # the loader marks the job succeeded
|
||||
tail = ((out or b"") + (err or b"")).decode(errors="replace")[-LOG_TAIL_CHARS:]
|
||||
async with SessionLocal() as session:
|
||||
job = await session.get(Job, job_id)
|
||||
if job is not None and job.status != JobStatus.succeeded:
|
||||
job.status = JobStatus.failed
|
||||
job.log = f"nextflow exited with code {proc.returncode}\n{tail}"
|
||||
job.log = f"nextflow exited with code {returncode}\n" + "\n".join(tail)
|
||||
job.finished_at = datetime.now(UTC)
|
||||
await session.commit()
|
||||
|
||||
|
||||
async def _record_progress(job_id: uuid.UUID, text: str) -> None:
|
||||
"""Put the latest output on the job, leaving a job that has already finished alone."""
|
||||
async with SessionLocal() as session:
|
||||
await session.execute(
|
||||
update(Job).where(Job.id == job_id, Job.status == JobStatus.running).values(log=text)
|
||||
)
|
||||
await session.commit()
|
||||
|
||||
@@ -33,14 +33,29 @@ async def test_local_without_nextflow_fails_fast_with_instructions(
|
||||
assert job["finished_at"] is not None
|
||||
|
||||
|
||||
class FakeProcess:
|
||||
def __init__(self, returncode: int, stderr: bytes) -> None:
|
||||
self.returncode = returncode
|
||||
self._stderr = stderr
|
||||
self.pid = 4242
|
||||
class FakeStdout:
|
||||
"""Stands in for asyncio's StreamReader: async-iterates the process output."""
|
||||
|
||||
async def communicate(self) -> tuple[bytes, bytes]:
|
||||
return b"", self._stderr
|
||||
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")
|
||||
@@ -149,3 +164,49 @@ async def test_the_pipeline_gets_its_own_database_url(
|
||||
await client.post(f"/api/samples/{sample_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)
|
||||
sample_id = await new_sample(client)
|
||||
|
||||
r = await client.post(f"/api/samples/{sample_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 Job, JobStatus, Sample
|
||||
|
||||
async with SessionLocal() as s:
|
||||
sample = Sample(name=f"s-{uuid.uuid4()}", vcf_uri="gs://b/x.vcf.gz", assembly="GRCh38")
|
||||
job = Job(sample=sample, status=JobStatus.succeeded)
|
||||
s.add_all([sample, 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
|
||||
|
||||
Reference in New Issue
Block a user