From 15c392b36c970e85141a2a8e3d8239294a27fcb2 Mon Sep 17 00:00:00 2001 From: Bryon Lewis Date: Tue, 21 Jul 2026 13:24:06 -0400 Subject: [PATCH] group worker processes to make it easier to kill them --- server/dive_tasks/utils.py | 106 ++++++++++++++++++++----- server/tests/test_stream_subprocess.py | 86 +++++++++++++++++++- 2 files changed, 172 insertions(+), 20 deletions(-) diff --git a/server/dive_tasks/utils.py b/server/dive_tasks/utils.py index 0f25c77c3..118b9d074 100644 --- a/server/dive_tasks/utils.py +++ b/server/dive_tasks/utils.py @@ -8,6 +8,7 @@ from subprocess import Popen import tempfile import threading +import time from typing import List, Optional, Tuple from urllib import request from urllib.parse import urlencode, urljoin @@ -21,6 +22,10 @@ TIMEOUT_COUNT = 'timeout_count' TIMEOUT_LAST_CHECKED = 'last_checked' TIMEOUT_CHECK_INTERVAL = 30 +# How often the cancel monitor polls task/job status while a subprocess runs. +CANCEL_MONITOR_INTERVAL = 30 +# Grace period after SIGTERM before escalating to SIGKILL of the process group. +CANCEL_TERM_GRACE_SECONDS = 5 def make_directory(path: Path): @@ -152,6 +157,36 @@ def describe_exit(code: int) -> str: return f'exited with nonzero status code {code}' +def _kill_process_group( + process: Popen, *, grace_seconds: float = CANCEL_TERM_GRACE_SECONDS +) -> None: + """ + Terminate a subprocess and its descendants. + + Pipelines/training are launched with shell=True, so signaling only the shell + PID leaves KWIVER/training children alive and holding stdout open. Starting + the process in a new session lets us kill the whole group. + """ + if process.poll() is not None: + return + try: + pgid = os.getpgid(process.pid) + except ProcessLookupError: + return + + for sig in (signal.SIGTERM, signal.SIGKILL): + try: + os.killpg(pgid, sig) + except ProcessLookupError: + return + if sig == signal.SIGTERM: + deadline = time.monotonic() + grace_seconds + while time.monotonic() < deadline: + if process.poll() is not None: + return + time.sleep(0.1) + + def stream_subprocess( task: Task, context: dict, @@ -172,22 +207,36 @@ def stream_subprocess( assert 'args' in popen_kwargs, "popen_kwargs must contain key 'args'" stop_event = threading.Event() + cancel_requested = threading.Event() + # Copy so callers are not mutated; default to a new session so cancel can + # kill the full process tree (shell + VIAME children). + launch_kwargs = dict(popen_kwargs) + launch_kwargs.setdefault('start_new_session', True) def monitor_cancellation(): """Thread that periodically checks for cancellation.""" - while not stop_event.wait(30): # Check every 30 seconds + while True: + if stop_event.is_set(): + return manager.refreshStatus() if check_canceled(task, context, force=True) or manager.status == JobStatus.CANCELING: - manager.write('\nCancellation detected. Stopping subprocess...\n', forceFlush=True) - process.send_signal(signal.SIGTERM) - process.send_signal(signal.SIGKILL) - process.send_signal(signal.SIGINT) - return # Stop the thread + cancel_requested.set() + _kill_process_group(process) + # Unblock the main thread's readline if any child still holds + # the write end open after the kill attempt. + if process.stdout is not None: + try: + process.stdout.close() + except (ValueError, OSError): + pass + return + if stop_event.wait(CANCEL_MONITOR_INTERVAL): + return with tempfile.TemporaryFile() as stderr_file: - manager.write(f"Running command: {str(popen_kwargs['args'])}\n", forceFlush=True) + manager.write(f"Running command: {str(launch_kwargs['args'])}\n", forceFlush=True) process = Popen( - **popen_kwargs, + **launch_kwargs, stdout=subprocess.PIPE, stderr=stderr_file, ) @@ -199,23 +248,42 @@ def monitor_cancellation(): cancel_thread = threading.Thread(target=monitor_cancellation, daemon=True) cancel_thread.start() - # call readline until it returns empty bytes - for line in iter(process.stdout.readline, b''): - # Pipeline tools may emit Latin-1 / CP1252 (e.g. 0xa0 NBSP) in log lines. - line_str = line.decode('utf-8', errors='replace') - manager.write(line_str) - if keep_stdout: - stdout += line_str + try: + # call readline until it returns empty bytes + for line in iter(process.stdout.readline, b''): + # Pipeline tools may emit Latin-1 / CP1252 (e.g. 0xa0 NBSP) in log lines. + line_str = line.decode('utf-8', errors='replace') + manager.write(line_str) + if keep_stdout: + stdout += line_str + except (ValueError, OSError): + # stdout may be closed by the cancel monitor to unblock this loop. + pass stop_event.set() - cancel_thread.join() + cancel_thread.join(timeout=CANCEL_TERM_GRACE_SECONDS + 5) + + if cancel_requested.is_set(): + manager.write('\nCancellation detected. Stopping subprocess...\n', forceFlush=True) # flush logs manager._flush() - # Wait for exit up to 30 seconds after kill - code = process.wait(30) - if check_canceled(task, context): + try: + # Wait for exit up to 30 seconds after kill + code = process.wait(30) + except subprocess.TimeoutExpired: + _kill_process_group(process) + try: + code = process.wait(5) + except subprocess.TimeoutExpired: + if cancel_requested.is_set() or check_canceled(task, context): + manager.write('\nCanceled during subprocess run.\n') + manager.updateStatus(JobStatus.CANCELED) + raise CanceledError('Job was canceled') + raise RuntimeError('Subprocess did not exit after kill') + + if cancel_requested.is_set() or check_canceled(task, context): manager.write('\nCanceled during subprocess run.\n') manager.updateStatus(JobStatus.CANCELED) raise CanceledError('Job was canceled') diff --git a/server/tests/test_stream_subprocess.py b/server/tests/test_stream_subprocess.py index 37722f364..9a02d08c0 100644 --- a/server/tests/test_stream_subprocess.py +++ b/server/tests/test_stream_subprocess.py @@ -1,9 +1,12 @@ import signal +import time from unittest.mock import MagicMock import pytest +from girder_worker.utils import JobStatus -from dive_tasks.utils import describe_exit, stream_subprocess +from dive_tasks import utils +from dive_tasks.utils import CanceledError, describe_exit, stream_subprocess def run_command(args): @@ -43,3 +46,84 @@ def test_nonzero_exit_raises(): def test_successful_exit_returns_normally(): assert run_command(['bash', '-c', 'exit 0']) == "" + + +def test_cancel_kills_process_group_and_raises_canceled(monkeypatch): + """ + Cancel must kill shell descendants, not only the shell PID. + + Long VIAME jobs use shell=True. If only the shell is signaled, children + keep stdout open and stream_subprocess blocks forever in CANCELING. + """ + monkeypatch.setattr(utils, 'CANCEL_MONITOR_INTERVAL', 0.05) + monkeypatch.setattr(utils, 'CANCEL_TERM_GRACE_SECONDS', 0.5) + + # Background child keeps writing; shell waits on it. Without process-group + # kill this hangs on readline after the shell alone is killed. + script = r""" +python3 -c " +import sys, time +while True: + print('tick', flush=True) + time.sleep(0.05) +" & +wait +""" + task = MagicMock() + task.canceled = False + manager = MagicMock() + manager.status = None + + def flip_cancel(*_args, **_kwargs): + task.canceled = True + manager.status = JobStatus.CANCELING + + # First status refresh arms cancel so the monitor notices quickly. + manager.refreshStatus.side_effect = flip_cancel + + start = time.monotonic() + with pytest.raises(CanceledError, match='canceled'): + stream_subprocess( + task, + {}, + manager, + { + 'args': script, + 'shell': True, + 'executable': '/bin/bash', + }, + ) + elapsed = time.monotonic() - start + + manager.updateStatus.assert_called_with(JobStatus.CANCELED) + # Must finish promptly; a hang here is the production failure mode. + assert elapsed < 10 + + +def test_cancel_via_task_canceled_flag(monkeypatch): + monkeypatch.setattr(utils, 'CANCEL_MONITOR_INTERVAL', 0.05) + monkeypatch.setattr(utils, 'CANCEL_TERM_GRACE_SECONDS', 0.5) + + task = MagicMock() + task.canceled = False + manager = MagicMock() + manager.status = None + + def flip_cancel(*_args, **_kwargs): + task.canceled = True + + manager.refreshStatus.side_effect = flip_cancel + + with pytest.raises(CanceledError, match='canceled'): + stream_subprocess( + task, + {}, + manager, + { + 'args': 'while true; do echo tick; sleep 0.05; done', + 'shell': True, + 'executable': '/bin/bash', + }, + ) + + manager.updateStatus.assert_called_with(JobStatus.CANCELED)