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
52 changes: 45 additions & 7 deletions aiter/utility/mp_tuner.py
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@
# Copyright (C) 2024-2026, Advanced Micro Devices, Inc. All rights reserved.
import math
import multiprocessing as mp
import os
import time
from multiprocessing import TimeoutError as MPTimeoutError

Expand All @@ -11,6 +12,7 @@
from aiter.test_common import checkAllclose

_TASK_START_TIMES = None
_TASK_PIDS = None


def _is_mapping_error(exc: BaseException) -> bool:
Expand All @@ -21,18 +23,37 @@ def _is_accelerator_error(exc: BaseException) -> bool:
return type(exc).__name__ == "AcceleratorError"


def _init_task_start_times(task_start_times):
global _TASK_START_TIMES
def _init_task_start_times(task_start_times, task_pids=None):
global _TASK_START_TIMES, _TASK_PIDS
_TASK_START_TIMES = task_start_times
_TASK_PIDS = task_pids


def _run_with_start_tracking(task_index, func, args):
if _TASK_START_TIMES is None:
raise RuntimeError("Task start-time storage is not initialized")
if _TASK_PIDS is not None:
_TASK_PIDS[task_index] = os.getpid()
_TASK_START_TIMES[task_index] = time.monotonic()
return func(*args)


def _task_worker_exited(task_pids, task_index):
"""True once the worker that started this task is gone. A worker killed by a
GPU memory fault never returns a result, so waiting on it only ends at the
task timeout (or never, without one)."""
pid = task_pids[task_index]
if pid == 0:
return False
try:
os.kill(pid, 0)
except ProcessLookupError:
return True
except PermissionError:
return False
return False


def _elapsed_since_task_start(task_start_times, task_index, now=None):
started_at = task_start_times[task_index]
if started_at == 0:
Expand All @@ -41,11 +62,13 @@ def _elapsed_since_task_start(task_start_times, task_index, now=None):
return current_time - started_at


def _reset_task_start_times(task_start_times, task_indices):
def _reset_task_start_times(task_start_times, task_indices, task_pids=None):
"""Mark tasks as queued again, so a resubmit is not judged against the
timestamp its previous attempt left behind."""
timestamp (or worker) its previous attempt left behind."""
for k in task_indices:
task_start_times[k] = 0
if task_pids is not None:
task_pids[k] = 0


def _merge_error_ratio(current, observed):
Expand Down Expand Up @@ -440,7 +463,7 @@ def mp_tuner(
def submit_tasks(pool, gpu_map, task_indices):
"""Submit tasks to the pool and return async results as a dict"""
task_indices = list(task_indices)
_reset_task_start_times(task_start_times, task_indices)
_reset_task_start_times(task_start_times, task_indices, task_pids)
return {
k: pool.apply_async(
_run_with_start_tracking,
Expand All @@ -462,10 +485,11 @@ def submit_tasks(pool, gpu_map, task_indices):

# Create initial pool and submit all tasks
task_start_times = mp.RawArray("d", len(task_group))
task_pids = mp.RawArray("i", len(task_group))
pool = mp.Pool(
processes=parallel_num,
initializer=_init_task_start_times,
initargs=(task_start_times,),
initargs=(task_start_times, task_pids),
)
pids = [pool.apply_async(get_pid) for i in range(start_idx, mp_num)]
gpu_map = {el.get(): i + start_idx for i, el in enumerate(pids)}
Expand Down Expand Up @@ -542,6 +566,20 @@ def add_dummy_result(k, results_list):
)

except MPTimeoutError:
if _task_worker_exited(task_pids, k):
print(
f"[!] Task {k} lost its worker (pid {task_pids[k]} exited) - likely a GPU fault",
flush=True,
)
failed_tasks.append((k, "worker exited"))
dummy_results = []
add_dummy_result(k, dummy_results)
result_dict[k] = (
dummy_results if shape_grouped else [dummy_results[0]]
)
completed_this_round.append((k, async_result))
pool_restart_needed = True
break
# Check if this specific task has exceeded its timeout (only if timeout is set)
if timeout is not None:
elapsed = _elapsed_since_task_start(task_start_times, k)
Expand Down Expand Up @@ -643,7 +681,7 @@ def add_dummy_result(k, results_list):
pool = mp.Pool(
processes=parallel_num,
initializer=_init_task_start_times,
initargs=(task_start_times,),
initargs=(task_start_times, task_pids),
)

# Recreate gpu_map for new processes (new PIDs)
Expand Down
45 changes: 45 additions & 0 deletions op_tests/tuning_tests/test_mp_tuner_logic.py
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,9 @@

import importlib
import multiprocessing as mp
import os
import subprocess
import sys
import time
import unittest
from multiprocessing import TimeoutError as MPTimeoutError
Expand Down Expand Up @@ -369,6 +372,48 @@ def test_reset_clears_only_the_given_slots(self):
reset_start_times(slots, [0, 2])
self.assertEqual(list(slots), [0, 22.0, 0])

def test_reset_clears_worker_pids(self):
tuner = importlib.import_module("aiter.utility.mp_tuner")
slots, pids = [11.0, 22.0], [101, 202]
tuner._reset_task_start_times(slots, [1], pids)
self.assertEqual(pids, [101, 0])


class TestTaskWorkerExit(unittest.TestCase):

def test_exited_worker_is_detected(self):
tuner = importlib.import_module("aiter.utility.mp_tuner")
worker_exited = getattr(tuner, "_task_worker_exited", None)
self.assertIsNotNone(
worker_exited,
"a task whose worker died (e.g. GPU fault) must not wait for the timeout",
)
proc = subprocess.Popen([sys.executable, "-c", "pass"])
proc.wait()
self.assertTrue(worker_exited([proc.pid], 0))
self.assertFalse(worker_exited([os.getpid()], 0))
self.assertFalse(worker_exited([0], 0), "a queued task has no worker yet")

def test_worker_pid_recorded_when_task_starts(self):
tuner = importlib.import_module("aiter.utility.mp_tuner")
ctx = mp.get_context("spawn")
start_times = ctx.RawArray("d", 1)
pids = ctx.RawArray("i", 1)
pool = ctx.Pool(
1,
initializer=tuner._init_task_start_times,
initargs=(start_times, pids),
)
try:
worker_pid = pool.apply_async(
tuner._run_with_start_tracking, (0, os.getpid, ())
).get(timeout=60)
self.assertEqual(pids[0], worker_pid)
self.assertGreater(start_times[0], 0)
finally:
pool.terminate()
pool.join()


class TestWorkerErrorRatio(unittest.TestCase):

Expand Down
Loading