Files
rarelens/api/app/services/events.py
T
Kemal Yaylali 33e788122b 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.
2026-09-12 07:46:44 +01:00

170 lines
6.6 KiB
Python

"""Start the annotation pipeline for a job.
Serverless (CLOUDRUN_JOB set): execute a Cloud Run job that runs the Nextflow driver, which
submits the pipeline tasks to Google Batch. Nothing runs, or is billed, between annotations.
Kubernetes (PUBSUB_TOPIC set): publish to Pub/Sub; an Argo Events sensor submits the annotate-vcf
WorkflowTemplate, whose exit handler marks the job failed if the workflow does not succeed.
Local: run Nextflow as a child process and watch it. The docker-compose API image has neither
Nextflow nor Docker, so there the job fails at once with instructions to run it from the host.
"""
import asyncio
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
logger = logging.getLogger(__name__)
PIPELINE_ENTRYPOINT = "/pipeline/main.nf"
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()
class LaunchError(Exception):
"""The pipeline could not be started; the message is stored on the job for the user."""
async def launch(job_id: uuid.UUID, vcf_uri: str, assembly: str) -> str:
"""Start the pipeline and return a reference to the run (stored as jobs.workflow_ref)."""
if settings.cloudrun_job:
return await _run_cloud_run_job(job_id, vcf_uri, assembly)
if settings.pubsub_topic:
message = {"job_id": str(job_id), "vcf_uri": vcf_uri, "assembly": assembly}
return await _publish(settings.pubsub_topic, message)
return await _run_local(job_id, vcf_uri, assembly)
async def drain() -> None:
"""Wait for local pipeline watchers (tests, graceful shutdown)."""
await asyncio.gather(*_watchers)
@cache
def _jobs_client() -> Any:
from google.cloud import run_v2
return run_v2.JobsClient()
async def _run_cloud_run_job(job_id: uuid.UUID, vcf_uri: str, assembly: str) -> str:
if not settings.gcp_project:
raise LaunchError("CLOUDRUN_JOB is set but GCP_PROJECT is not")
name = (
f"projects/{settings.gcp_project}/locations/{settings.gcp_region}"
f"/jobs/{settings.cloudrun_job}"
)
args = [
"run", PIPELINE_ENTRYPOINT, "-profile", "gcp",
"--vcf", vcf_uri, "--job_id", str(job_id), "--assembly", assembly,
]
request = {"name": name, "overrides": {"container_overrides": [{"args": args}]}}
try:
operation = await asyncio.to_thread(_jobs_client().run_job, request=request)
except Exception as e:
logger.exception("starting Cloud Run job %s failed", name)
raise LaunchError(f"could not start {settings.cloudrun_job}: {e}") from e
name_or_job = getattr(getattr(operation, "metadata", None), "name", "")
execution = name_or_job or settings.cloudrun_job or "started"
return f"cloudrun:{execution.rsplit('/', 1)[-1]}"
@cache
def _publisher() -> Any:
from google.cloud import pubsub_v1
return pubsub_v1.PublisherClient()
async def _publish(topic_name: str, message: dict[str, str]) -> str:
if not settings.gcp_project:
raise LaunchError("PUBSUB_TOPIC is set but GCP_PROJECT is not")
publisher = _publisher()
topic = publisher.topic_path(settings.gcp_project, topic_name)
try:
future = publisher.publish(topic, json.dumps(message).encode())
message_id = await asyncio.to_thread(future.result, timeout=PUBLISH_TIMEOUT_S)
except Exception as e:
logger.exception("publishing to %s failed", topic)
raise LaunchError(f"could not publish to {topic}: {e}") from e
return f"pubsub:{message_id}"
async def _run_local(job_id: uuid.UUID, vcf_uri: str, assembly: str) -> str:
nextflow = shutil.which("nextflow")
if nextflow is None:
raise LaunchError(
"Nextflow is not installed where the API runs. Start the pipeline from the host with "
f"`make annotate JOB={job_id} VCF=<path to the VCF>`."
)
cmd = [
nextflow, "run", str(settings.pipeline_dir / "main.nf"),
"-profile", settings.nextflow_profile,
"--vcf", vcf_uri, "--job_id", str(job_id), "--assembly", assembly,
]
# The loader reads DATABASE_URL from its environment; keep it off the command line.
env = {**os.environ, "DATABASE_URL": settings.pipeline_database_url or settings.database_url}
try:
proc = await asyncio.create_subprocess_exec(
*cmd,
cwd=settings.pipeline_dir,
env=env,
stdout=asyncio.subprocess.PIPE,
stderr=asyncio.subprocess.STDOUT,
)
except OSError as e:
raise LaunchError(f"could not start Nextflow: {e}") from e
task = asyncio.create_task(_watch(job_id, proc))
_watchers.add(task)
task.add_done_callback(_watchers.discard)
return f"nextflow:pid-{proc.pid}"
async def _watch(job_id: uuid.UUID, proc: Any) -> None:
"""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
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 {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()