Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions flask-server/.env.example
Original file line number Diff line number Diff line change
Expand Up @@ -63,6 +63,11 @@ CANCERVERSE_LOWRES_PATH=/home/visitor/cancerverse_lowres
# GPU_WORKER_MAX_CPU_LOAD=0.75
# GPU_WORKER_COOLDOWN_SECONDS=600
# GPU_WORKER_SSH_MULTIPLEX=true
# Jobs allowed to run at once, one per worker (capped by the number of workers;
# 1 = strictly one at a time). A job waits up to MAX_WAIT for a busy worker
# before it may use this host's GPU.
# GPU_WORKER_PARALLEL=1
# GPU_WORKER_MAX_WAIT_SECONDS=300
# Keep BLAS/OpenMP from multiplying CPU threads inside each model subprocess.
# Set model-specific values only after measuring the host; these are safe
# process-wide defaults for a shared web server.
Expand Down
10 changes: 5 additions & 5 deletions flask-server/api/api_blueprint.py
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
from flask import Blueprint, send_file, make_response, request, jsonify, Response, current_app, stream_with_context
from werkzeug.utils import secure_filename
from services.session_manager import SessionManager, generate_uuid
from services.auto_segmentor import run_auto_segmentation, cancel_session, cancel_all_inference
from services.auto_segmentor import run_auto_segmentation, cancel_session, cancel_all_inference, max_parallel_jobs
from services.mesh_generation import (
bake_case_meshes,
generate_mesh_manifest,
Expand Down Expand Up @@ -2166,12 +2166,12 @@ def download_segmentation_zip(id):
# services.auto_segmentor serializes the actual model execution, but an
# unbounded number of accepted requests would still create one blocked Python
# thread and retain one uploaded volume per request. Keep the cap configurable
# so it can be sized to the host's RAM and GPU; the default includes the job
# currently running plus a small waiting queue.
# so it can be sized to the host's RAM and GPU; the default is the jobs that can
# run at once plus a small waiting queue (4 with one job at a time, as always).
try:
_INFERENCE_MAX_PENDING = max(1, int(os.getenv("INFERENCE_MAX_PENDING", "4")))
_INFERENCE_MAX_PENDING = max(1, int(os.getenv("INFERENCE_MAX_PENDING", str(3 + max_parallel_jobs()))))
except (TypeError, ValueError):
_INFERENCE_MAX_PENDING = 4
_INFERENCE_MAX_PENDING = 3 + max_parallel_jobs()
_INFERENCE_PENDING_SLOTS = threading.BoundedSemaphore(_INFERENCE_MAX_PENDING)


Expand Down
52 changes: 50 additions & 2 deletions flask-server/deploy/GPU_WORKERS.md
Original file line number Diff line number Diff line change
Expand Up @@ -68,15 +68,57 @@ download over the 100 Gbit/s LAN, typically 1-2 s. Transfers abort after 120 s
without progress (hard cap 15 min each). Model run time is the same as on
bdmap1 (same GB10 hardware).

Jobs still run one at a time (the same global lock as local runs), so
throughput is unchanged; the gain is that bdmap1's GPU stays free.
Jobs run one at a time unless `GPU_WORKER_PARALLEL` is raised (see Parallel
jobs); the main gain is that bdmap1's GPU stays free.

If a worker runs its own warm predictor on the same port as `EPAI_WARM_URL` /
`LESIONSEG_WARM_URL`, the forwarded URL (127.0.0.1) makes jobs on that worker
use it automatically; otherwise they use the cold path, as bdmap1 does when its
warm predictor is not running. Starting warm predictors on the workers makes
jobs faster than today.

## Assumptions

- **Each worker has one GPU** (GB10), and jobs are placed per worker, not per
GPU. Model commands carry the web host's `CUDA_VISIBLE_DEVICES` choice
(default 0) to the worker unchanged. A multi-GPU worker would need per-GPU
placement first.
- **One gunicorn process.** Worker reservations live in that process's memory,
which matches the documented deployment (one process, eight threads).
- **No strict first-come-first-served order** between waiting jobs, as with the
single lock before.

## Parallel jobs

By default one model job runs at a time, exactly as before. Set
`GPU_WORKER_PARALLEL=N` (in `.env`, with the workers enabled) to let up to N jobs
run at once, one per worker. It is capped at the number of workers in
`GPU_WORKER_HOSTS`, and ignored (1) while workers are disabled. Recommended: the
number of dependable workers (for example 2 with bdmap2 and bdmap4), leaving
the others, such as bdmap3, as spares that take over when one is down.

How it stays safe:
- **No double booking.** Workers are health-checked without being reserved, and
the winner is reserved atomically, so two jobs never get the same worker and
one job never sees "no worker free" just because another is checking.
- **Waiting beats overflowing.** If every usable worker is busy with our own
jobs, a job waits (up to `GPU_WORKER_MAX_WAIT_SECONDS`, 300; cancel works while
waiting) instead of running on bdmap1. bdmap1 is used only when no worker is
up at all, or after that wait. Jobs beyond N wait in the normal queue, and
the queue size (`INFERENCE_MAX_PENDING`) defaults to `3 + N`.
- **bdmap1's own GPU still runs at most one model at a time**, however many
jobs are in flight: fallback jobs queue on a lock of their own.
- **A session never runs against itself.** A repeat request for a session whose
job is still queued or running (a retry, a double click) waits for it, as it
did behind the old single lock, and does not use up a job slot while waiting.
- **The older single-host ePAI mode (`EPAI_REMOTE_ENABLED`) keeps jobs strictly
serial**, because it sends every ePAI job to one fixed GPU outside the pool.
- Placement is still by health and load, so a hot, throttled or busy worker is
skipped, and a failing one cools down.

To roll back to one job at a time, unset `GPU_WORKER_PARALLEL` (or set it to 1)
and reload.

## One-time setup (per worker)

1. Passwordless SSH from bdmap1 as `visitor`:
Expand Down Expand Up @@ -107,6 +149,12 @@ Restart gunicorn with the usual deploy procedure. To roll back, set

## Maintenance

- **Reload or restart gunicorn only when no job is running.** Worker reservations
and the web host GPU lock live in the process; a new process starts with none.
(The health check still keeps it off a worker that is busy, but a job in
flight is lost when its process exits anyway.) Check that no model process
runs on bdmap1 and no worker has a session folder or GPU process first.

- After installing or updating a model, conda env or weights on bdmap1, run
`flask-server/scripts/sync_gpu_workers.sh` (all hosts from `.env`).
`flask-server/scripts` itself is synced automatically with every job.
Expand Down
149 changes: 120 additions & 29 deletions flask-server/services/auto_segmentor.py
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
import contextlib
import os
import uuid
import signal
Expand All @@ -14,8 +15,89 @@
# Load environment variables
load_dotenv()

# Only one model inference runs at a time to avoid GPU OOM
_gpu_lock = threading.Lock()
def _is_truthy(value) -> bool:
return gpu_workers._truthy(value)


def max_parallel_jobs() -> int:
"""How many model jobs may run at once.

1 (strictly one at a time, as before) unless GPU workers are enabled and
GPU_WORKER_PARALLEL asks for more; never more than there are workers, since
each worker runs one job at a time, and 1 while the older EPAI_REMOTE_ENABLED
mode is on. Read once at import.
"""
if not gpu_workers.enabled():
return 1
if _is_truthy(os.getenv("EPAI_REMOTE_ENABLED", "false")):
# The older explicit single-host ePAI mode sends every ePAI job to one
# fixed GPU outside the worker pool; only serial jobs are safe there.
return 1
try:
wanted = int(os.getenv("GPU_WORKER_PARALLEL", "1"))
except ValueError:
wanted = 1
return max(1, min(wanted, len(gpu_workers.hosts())))


# A job holds one of these slots for its whole run, so concurrent requests queue
# instead of OOM-ing a GPU. One slot = the old single global lock.
_job_slots = threading.BoundedSemaphore(max_parallel_jobs())

# bdmap1's own GPU runs at most one model at a time, however many jobs are in
# flight: it is the machine that serves the website (shared CPU/GPU memory), so
# fallback jobs queue here rather than pile onto it.
_local_gpu_lock = threading.Lock()


_session_locks = {} # session id -> [lock, waiters+holder count]


@contextlib.contextmanager
def _session_exclusive(session_id):
"""A session never runs against itself.

A repeat request (a retry, a double click) for a session whose job is still
queued or running used to serialize behind it on the single global lock.
With several jobs in flight it would overlap it, two runs writing the same
workspace, outputs, process entry and job state. It waits here instead, and
it does so before taking a job slot, so a waiting duplicate costs no capacity.
"""
if not session_id:
yield
return
with _session_procs_lock:
entry = _session_locks.setdefault(session_id, [threading.Lock(), 0])
entry[1] += 1
try:
with entry[0]:
yield
finally:
with _session_procs_lock:
entry[1] -= 1
if entry[1] == 0:
_session_locks.pop(session_id, None)


def _current_session_cancelled() -> bool:
sid = getattr(_thread_session, "sid", None)
return bool(sid) and sid in _cancelled_sessions


@contextlib.contextmanager
def _local_gpu_slot(cancelled):
"""Hold the web host's GPU for one local model run; cancellable while waiting."""
while not _local_gpu_lock.acquire(timeout=1.0):
if cancelled():
raise RuntimeError("Inference cancelled")
try:
# Also covers a cancel that landed while waiting but just as the lock
# came free (no process existed to kill), and a lock that was free.
if cancelled():
raise RuntimeError("Inference cancelled")
yield
finally:
_local_gpu_lock.release()

# ── Per-session subprocess tracking (for user-initiated cancel) ──
# Each session's worker thread binds its session id thread-locally; _tracked_run
Expand Down Expand Up @@ -118,17 +200,23 @@ def _cancelled():
if _session_procs.get(sid) is registered[-1]:
_session_procs.pop(sid, None)

proc = subprocess.Popen(cmd, **kwargs)
if sid:
with _session_procs_lock:
_session_procs[sid] = proc
try:
stdout, stderr = proc.communicate()
finally:
is_model_cmd = bool(remote_dir) and bool(kwargs.get("shell")) and isinstance(cmd, str)
if is_model_cmd and gpu_workers.enabled():
local_slot = _local_gpu_slot(_current_session_cancelled)
else:
local_slot = contextlib.nullcontext()
with local_slot:
proc = subprocess.Popen(cmd, **kwargs)
if sid:
with _session_procs_lock:
if _session_procs.get(sid) is proc:
_session_procs.pop(sid, None)
_session_procs[sid] = proc
try:
stdout, stderr = proc.communicate()
finally:
if sid:
with _session_procs_lock:
if _session_procs.get(sid) is proc:
_session_procs.pop(sid, None)
if retry_of and proc.returncode == 0:
# Local succeeded where the worker failed: the worker is at fault.
gpu_workers.mark_failed(retry_of)
Expand Down Expand Up @@ -290,23 +378,24 @@ def _resolve_conda_activate_path():
def run_auto_segmentation(input_path, session_dir, model, session_id=None, on_start=None):
"""Run one model; see _run_auto_segmentation. Always clears per-job state."""
token = object()
try:
return _run_auto_segmentation(input_path, session_dir, model, session_id, on_start, token)
finally:
_thread_session.remote_session_dir = None
if session_id:
with _session_procs_lock:
# Runs after _gpu_lock is released, so a newer run of the same
# session may already be active; only clear our own state.
if _active_sessions.get(session_id) is token:
del _active_sessions[session_id]
_cancelled_sessions.discard(session_id)
with _session_exclusive(session_id):
try:
return _run_auto_segmentation(input_path, session_dir, model, session_id, on_start, token)
finally:
_thread_session.remote_session_dir = None
if session_id:
with _session_procs_lock:
# Only clear our own state.
if _active_sessions.get(session_id) is token:
del _active_sessions[session_id]
_cancelled_sessions.discard(session_id)


def _run_auto_segmentation(input_path, session_dir, model, session_id=None, on_start=None, token=None):
"""
Dispatch to the appropriate model inference function.
Serialized via _gpu_lock so concurrent requests queue instead of OOM-ing.
Limited by _job_slots (one at a time unless GPU workers allow more), so
concurrent requests queue instead of OOM-ing.
Returns the output directory path on success, raises on failure.

session_id: binds this worker thread so _tracked_run/cancel_session can
Expand All @@ -315,7 +404,7 @@ def _run_auto_segmentation(input_path, session_dir, model, session_id=None, on_s
run (used when the user cancelled while the job was still queued) and
makes this function return None.
"""
with _gpu_lock:
with _job_slots:
if session_id:
# Before on_start: a cancel landing after its status check is
# still recorded. A flag left from an earlier run is dropped.
Expand Down Expand Up @@ -1199,10 +1288,6 @@ def _run_lesionsegmenter_inference(input_path: str, session_dir: str, conda_path
return output_ct_dir


def _is_truthy(value: str) -> bool:
return str(value or "").strip().lower() in {"1", "true", "yes", "y", "on"}


def _run_checked_process(cmd: list[str], error_prefix: str):
process = _tracked_run(cmd, text=True, capture_output=True)
if process.returncode != 0:
Expand Down Expand Up @@ -1517,7 +1602,13 @@ def _run_media_agentic_inference(
print(f"[INFO] Running MedIA-Agentic {model_type} inference...")
print(f"[INFO] Command: {' '.join(cmd)}")

result = subprocess.run(cmd, capture_output=True, text=True)
# This model runs on the web host's own GPU (it is not routed to a worker),
# so with several jobs in flight it queues on the same lock as every other
# local model run. Only taken when workers are enabled (the disabled path is
# unchanged: one job at a time already).
gpu_slot = _local_gpu_slot(_current_session_cancelled) if gpu_workers.enabled() else contextlib.nullcontext()
with gpu_slot:
result = subprocess.run(cmd, capture_output=True, text=True)

if result.returncode != 0:
print(f"[ERROR] MedIA-Agentic inference failed: {result.stderr}")
Expand Down
Loading
Loading