Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
37 commits
Select commit Hold shift + click to select a range
ec752ba
Run multi-rank CPU unit tests in CI via LOCAL_SIZE
delock Sep 1, 2026
f2b7d0b
Disable dist-env reuse for the full multi-rank CPU run
delock Sep 1, 2026
f2793e7
Split CPU unit tests into halves with per-half timeouts
delock Sep 2, 2026
4d4019d
Merge branch 'master' into ci/cpu-multi-rank-local-size
delock Sep 15, 2026
89ca163
Isolate shm-allreduce segments per test pool and clean them up
delock Sep 16, 2026
453f60f
Pre-download HF fixtures and share one cache across all test phases
delock Sep 16, 2026
6799643
Download HF fixtures with the hf CLI
delock Sep 16, 2026
2654cac
Disable Xet transfers for the anonymous HF fixture download
delock Sep 16, 2026
d429f62
Merge remote-tracking branch 'origin/master' into pr-8381-new
delock Sep 17, 2026
132298c
Merge remote-tracking branch 'origin/master' into pr-8381-new
delock Sep 17, 2026
dfb0c50
Scope the bf16 version floors to NCCL transports in the test check
delock Sep 17, 2026
0194033
Merge remote-tracking branch 'origin/master' into pr-8381-new
delock Sep 18, 2026
3fba8dc
Stop pinning DDP device_ids in the autocast baseline on CPU
delock Sep 18, 2026
3b74fba
Merge remote-tracking branch 'origin/master' into pr-8381-new
delock Sep 18, 2026
33e7dcc
Skip the dynamic offload-state tests on the cpu accelerator
delock Sep 18, 2026
940d404
Merge remote-tracking branch 'origin/master' into ci/cpu-multi-rank-l…
delock Sep 21, 2026
61abf25
Fix pipe tests failures on CPU
jinyouzhi Sep 20, 2026
68fdd74
Skip fp16 ZeroPP tests on accelerators without fp16 support
jinyouzhi Sep 20, 2026
2adfc0d
Skip ZERO++ Quantized in CPU
jinyouzhi Sep 20, 2026
08c966f
Merge remote-tracking branch 'origin/master' into pr-8381-new
delock Sep 22, 2026
09e3cba
Skip fp16 universal-checkpoint fixtures without fp16 support
delock Sep 22, 2026
0cb5cdc
Merge remote-tracking branch 'origin/master' into pr-8381-new
delock Sep 24, 2026
2b23590
Skip fp16 coalesce tests on accelerators without fp16 support
delock Sep 24, 2026
e4ec54f
Loosen the late-iteration gradient tolerance in the checkpointing test
delock Sep 24, 2026
2bdbbdf
Allow fp16 ulp noise in autocast loss parity below ZeRO-3
delock Sep 24, 2026
8fa6e22
Pre-build the shm comm op so rank workers never race the JIT lock
delock Sep 24, 2026
18ee3ec
Skip fp16 no_sync tests on accelerators without fp16 support
delock Sep 25, 2026
ecc198e
Allow fp16 ulp noise in autocast loss parity at ZeRO-3
delock Sep 26, 2026
5770829
Move the CPU workflow to torch 2.14.0 for this branch
delock Sep 26, 2026
a84e1b7
Run only the multi-rank tests in this workflow
delock Sep 26, 2026
5914ddc
Fix Ulysses SP registration check and seqlen exchange on gloo
delock Sep 27, 2026
86eeb05
Adapt the ulysses_sp_hf tests for non-CUDA accelerators
delock Sep 27, 2026
c81a577
Raise the non-v1 half timeout to 180m
delock Sep 27, 2026
a0680d6
De-hardcode the disable-in-eval test device
delock Sep 27, 2026
1e6872a
Skip fp16 MoE checkpoint tests on accelerators without fp16 support
delock Sep 27, 2026
e6c3d95
Skip the remaining fp16-config tests without fp16 support
delock Sep 27, 2026
30a96e4
Raise the non-v1 half timeout to 210m
delock Sep 28, 2026
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: 49 additions & 3 deletions .github/workflows/cpu-torch-latest.yml
Original file line number Diff line number Diff line change
Expand Up @@ -81,7 +81,20 @@ jobs:
runs-on: ubuntu-24.04

env:
DEFAULT_TORCH_PRESET: '2.10.0-cpu'
# The runner is single-socket, so CPU_Accelerator.device_count() reports 1
# NUMA node and the per-device gate in tests/unit/common.py skips every
# multi-rank test. CPU ranks are plain processes over gloo, so advertise 4
# local devices to let world_size<=4 tests run. The test harness re-sets
# LOCAL_SIZE per worker, so this value only affects the launch gate.
LOCAL_SIZE: '4'

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1 Badge Keep CUDA-only distributed tests out of the CPU run

In the cpu-torch-latest job, advertising four devices admits every distributed test with world_size <= 4, not only CPU-safe tests. For example, tests/unit/ulysses_alst/test_ulysses_sp_hf.py:240-263 defines an unguarded two-rank test that creates tensors on cuda:<rank>; because CPU_Accelerator.is_available() returns true, the harness does not skip it, and the all-unit pytest invocation at workflow line 283 will fail on the CPU-only PyTorch installation. Scope this override to an explicitly CPU-compatible subset or add CPU capability checks before enabling the previously skipped tests.

Useful? React with 👍 / 👎.

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P1 Badge Add the required sign-off trailer

This is a non-merge commit, but its commit message has no Signed-off-by trailer. Add the author sign-off so the commit satisfies the repository's commit and CI requirements.

AGENTS.md reference: AGENTS.md:L8-L8

Useful? React with 👍 / 👎.

# Multi-rank tests churn mp pools for hours; reused pools eventually hang in
# cleanup and stall workers until the 6h job limit. Fresh pools per test are
# slower but let the suite finish (knob documented in tests/unit/common.py).
DS_DISABLE_REUSE_DIST_ENV: '1'
# Only the multi-rank tests this LOCAL_SIZE gate admits run here; single-rank
# tests are deselected because the plain cpu-torch-latest run guards them.
DS_MULTIRANK_ONLY: '1'
DEFAULT_TORCH_PRESET: '2.14.0-cpu'
DEFAULT_TRANSFORMERS_SOURCE: 'git'
# Manual PyPI fallback only; scheduled and default manual runs use Git.
DEFAULT_TRANSFORMERS_VERSION: '4.51.3'
Expand Down Expand Up @@ -158,6 +171,11 @@ jobs:
torchvision_install_version='0.25.0'
torch_test_version='2.10'
;;
'2.14.0-cpu')
torch_install_version='2.14.0'
torchvision_install_version='0.29.0'
torch_test_version='2.14'
;;
*)
echo "Unsupported torch_preset: $selected_preset" >&2
exit 1
Expand Down Expand Up @@ -270,9 +288,37 @@ jobs:
run: |
pip list

- name: Pre-build JIT ops
# Rank workers JIT-load ops concurrently on first use; torch's build
# lock has no stale handling, so a leftover baton from a killed build
# wedges every later load. Build the hot ops once up front so tests
# only hit the disk cache.
run: |
python -c "
from deepspeed.comm.torch import build_shm_op
build_shm_op()
import torch
from deepspeed.ops.adam import DeepSpeedCPUAdam
DeepSpeedCPUAdam([torch.nn.Parameter(torch.randn(8))])
"

- name: Pre-download HF test fixtures
# Concurrent xdist workers downloading the same model race on the
# transformers cache lock and fail with PermissionError; populate the
# shared cache once so tests only hit it read-only.
run: |
HF_HUB_DISABLE_XET=1 HF_HOME=/tmp/hf_home hf download bert-base-uncased

- name: Unit tests
run: |
unset TORCH_CUDA_ARCH_LIST # only jit compile for current arch
cd tests
HF_HOME=/tmp/hf_home/ pytest $PYTEST_OPTS --forked -n 4 unit/ --torch_ver="$TORCH_TEST_VERSION"
HF_HOME=/tmp/hf_home/ pytest $PYTEST_OPTS --forked -m 'sequential' unit/ --torch_ver="$TORCH_TEST_VERSION"
# The suite is split so each half gets fresh xdist workers: multi-rank pool
# teardown eventually wedges a worker, and one 3400-test process never
# reaches its summary inside the 6h job limit. maxfail is raised so every
# failure is listed, and timeout caps each half.
overall=0
HF_HOME=/tmp/hf_home timeout 210m pytest $PYTEST_OPTS --maxfail=100000 --forked -n 4 unit/ --ignore=unit/v1 --torch_ver="$TORCH_TEST_VERSION" || overall=$?
HF_HOME=/tmp/hf_home timeout 150m pytest $PYTEST_OPTS --maxfail=100000 --forked -n 4 unit/v1 --torch_ver="$TORCH_TEST_VERSION" || overall=$?
HF_HOME=/tmp/hf_home/ pytest $PYTEST_OPTS --forked -m 'sequential' unit/ --torch_ver="$TORCH_TEST_VERSION" || overall=$?
exit $overall
23 changes: 15 additions & 8 deletions deepspeed/runtime/sequence_parallel/ulysses_sp.py
Original file line number Diff line number Diff line change
Expand Up @@ -444,19 +444,24 @@ def register_with_transformers(
mpu.initialize_sequence_parallel(sequence_parallel_size=sequence_parallel_size)

from transformers import PreTrainedModel
if hasattr(model_name_or_path, "config") or isinstance(model_name_or_path, PreTrainedModel):
model_was_loaded = hasattr(model_name_or_path, "config") or isinstance(model_name_or_path, PreTrainedModel)
if model_was_loaded:
# we already have the model (or a PEFT wrapper with config attribute)
hf_model_config = model_name_or_path.config
else:
# if we don't have the model yet at this stage
hf_model_config = AutoConfig.from_pretrained(model_name_or_path)

model_attn_implementation = getattr(hf_model_config, "_attn_implementation", None)
if model_attn_implementation is not None and model_attn_implementation != core_attn_implementation:
raise ValueError(
f"core_attn_implementation='{core_attn_implementation}' does not match "
f"model config attn_implementation='{model_attn_implementation}'. "
"Set both to the same value so sequence-parallel wrapper can intercept the active attention path.")
# Only a loaded model's config carries the attn implementation resolved at load time;
# a bare AutoConfig still holds the unresolved 'eager' default, so there is nothing
# meaningful to compare for a string model path.
if model_was_loaded:
model_attn_implementation = getattr(hf_model_config, "_attn_implementation", None)
if model_attn_implementation is not None and model_attn_implementation != core_attn_implementation:
raise ValueError(
f"core_attn_implementation='{core_attn_implementation}' does not match "
f"model config attn_implementation='{model_attn_implementation}'. "
"Set both to the same value so sequence-parallel wrapper can intercept the active attention path.")

# eager always materializes a 4D attention_mask (O(n²) memory) and cannot fall back
# to is_causal=True like sdpa — so it's incompatible with SP which discards masks.
Expand Down Expand Up @@ -666,7 +671,9 @@ def refill(self):
"Ensure your data collator includes position_ids in its output.")

# we have batches of variable seqlen so in order to do all_gather on batches - we need to know the exact length of each tensor on each rank
seqlen = torch.tensor(batch["input_ids"].shape[1], dtype=torch.int64, device=self.device)
# gloo validates gather shapes strictly, so send a 1-element tensor to match the
# receive list; a 0-dim scalar only passes on backends that move raw bytes.
seqlen = torch.full((1, ), batch["input_ids"].shape[1], dtype=torch.int64, device=self.device)
seqlens = [torch.zeros(1, dtype=torch.int64, device=self.device) for _ in range(self.sp_world_size)]
dist.all_gather(seqlens, seqlen, group=self.sp_group)
seqlens = [x[0].item() for x in seqlens]
Expand Down
28 changes: 28 additions & 0 deletions tests/conftest.py
Original file line number Diff line number Diff line change
Expand Up @@ -61,6 +61,34 @@ def check_environment(pytestconfig):

# Override of pytest "runtest" for DistributedTest class
# This hook is run before the default pytest_runtest_call
# The multi-rank-only CI workflow (DS_MULTIRANK_ONLY=1) spends its budget only on
# the multi-rank tests its LOCAL_SIZE gate admits; single-rank tests are deselected
# because the plain cpu-torch-latest run already guards them. Per-test world_size
# marks are honored with the same precedence as the launcher (mark, then class
# attribute); a fixture-injected world_size falls back to the class attribute.
def pytest_collection_modifyitems(config, items):
if os.environ.get('DS_MULTIRANK_ONLY', '0') != '1':
return

def is_multirank(item):
cls = getattr(item, 'cls', None)
if not getattr(cls, 'is_dist_test', False):
return False
mark = item.get_closest_marker('world_size')
if mark is not None:
world_size = mark.args[0]
else:
world_size = getattr(cls, 'world_size', None)
if isinstance(world_size, (list, tuple)):
return any(ws > 1 for ws in world_size)
return isinstance(world_size, int) and world_size > 1

deselected = [item for item in items if not is_multirank(item)]
if deselected:
config.hook.pytest_deselected(items=deselected)
items[:] = [item for item in items if is_multirank(item)]


@pytest.hookimpl(tryfirst=True)
def pytest_runtest_call(item):
# We want to use our own launching function for distributed tests
Expand Down
2 changes: 2 additions & 0 deletions tests/unit/checkpoint/test_moe_checkpoint.py
Original file line number Diff line number Diff line change
Expand Up @@ -8,12 +8,14 @@

from unit.common import DistributedTest
from unit.simple_model import *
from deepspeed.accelerator import get_accelerator

from unit.checkpoint.common import checkpoint_correctness_verification

import pytest


@pytest.mark.skipif(not get_accelerator().is_fp16_supported(), reason="fp16 is not supported on this accelerator")
class TestMoECheckpoint(DistributedTest):
world_size = 4

Expand Down
5 changes: 5 additions & 0 deletions tests/unit/checkpoint/test_pipeline.py
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@
from unit.simple_model import *
from unit.checkpoint.common import checkpoint_correctness_verification
from unit.util import skip_on_arch
from deepspeed.accelerator import get_accelerator

import pytest

Expand All @@ -18,6 +19,10 @@ class TestPipelineCheckpoint(DistributedTest):
@pytest.mark.parametrize("zero_stage", [0, 1])
def test_checkpoint_pipe_engine(self, zero_stage, tmpdir):
skip_on_arch(min_arch=7)
# fp16 is only enabled for zero_stage > 0; skip that parametrization on
# accelerators without fp16 support instead of failing the sanity check.
if zero_stage > 0 and not get_accelerator().is_fp16_supported():
pytest.skip("fp16 is not supported on this accelerator")

config_dict = {
"train_batch_size": 2,
Expand Down
5 changes: 5 additions & 0 deletions tests/unit/checkpoint/test_universal_checkpoint.py
Original file line number Diff line number Diff line change
Expand Up @@ -162,6 +162,11 @@ class _baseline(DistributedFixture):
world_size = None

def run(self, tmpdir, ds_config, zero_stage, dtype, load_optim, use_torch_adam):
# fp16 configs crash deepspeed.initialize's sanity check on accelerators
# without fp16 support, surfacing as a setup error for every dependent
# test instead of a skip.
if dtype == torch.float16 and not get_accelerator().is_fp16_supported():
pytest.skip("fp16 is not supported on this accelerator")
hidden_dim = 10
train_save_convert(ds_config, hidden_dim, load_optim, use_torch_adam, dtype, tmpdir, self.world_size)

Expand Down
49 changes: 49 additions & 0 deletions tests/unit/common.py
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@

# DeepSpeed Team

import itertools
import os
import re
import time
Expand Down Expand Up @@ -44,12 +45,21 @@ def get_xdist_worker_id():
return None


_master_port_counter = itertools.count()


def get_master_port(base_port=29500, port_range_size=1000):
xdist_worker_id = get_xdist_worker_id()
if xdist_worker_id is not None:
# Make xdist workers use different port ranges to avoid race conditions
base_port += port_range_size * xdist_worker_id

# The bind-and-release probe below hands every pool the same first-free port
# when each test runs in a fresh process (--forked), but this port also keys
# the shm allreduce segments, so pools must not share it. Offset the probe
# start per process and per pool.
base_port += (os.getpid() + next(_master_port_counter)) % (port_range_size - 100)

# Select first open port in range
port = base_port
max_port = base_port + port_range_size
Expand All @@ -64,6 +74,25 @@ def get_master_port(base_port=29500, port_range_size=1000):
raise IOError('no free ports')


def _warm_shm_comm_op():
# Rank workers JIT-load the shm comm op on first init_distributed; concurrent
# loads serialize on torch's FileBaton, which has no stale-lock handling, so a
# baton left by a killed process deadlocks every later load (the CI worker
# wedges). Drop a stale baton before it blocks anything, then build once here
# so workers only hit the disk cache.
try:
from torch.utils.cpp_extension import _get_build_directory
build_dir = Path(_get_build_directory("deepspeed_shm_comm", verbose=False))
lock = build_dir / "lock"
if (build_dir /
"deepspeed_shm_comm.so").exists() and lock.exists() and time.time() - lock.stat().st_mtime > 600:
lock.unlink(missing_ok=True)
from deepspeed.comm.torch import build_shm_op
build_shm_op()
except Exception:
pass


def _get_cpu_socket_count():
import shlex
p1 = subprocess.Popen(shlex.split("cat /proc/cpuinfo"), stdout=subprocess.PIPE)
Expand Down Expand Up @@ -182,6 +211,7 @@ def _get_fixture_kwargs(self, request, func):
def _launch_daemonic_procs(self, num_procs, init_method):
# Create process pool or use cached one
master_port = None
_warm_shm_comm_op()

if get_accelerator().device_name() == 'hpu':
if self.reuse_dist_env:
Expand All @@ -198,6 +228,7 @@ def _launch_daemonic_procs(self, num_procs, init_method):
master_port = get_master_port()

# Run the test
self._master_port = master_port
args = [(local_rank, num_procs, master_port, init_method) for local_rank in range(num_procs)]
skip_msgs_async = pool.starmap_async(self._dist_run, args)

Expand All @@ -222,6 +253,7 @@ def _launch_non_daemonic_procs(self, num_procs, init_method):
assert not self.reuse_dist_env, "Cannot reuse distributed environment with non-daemonic processes"

master_port = get_master_port()
self._master_port = master_port
skip_msg = mp.Queue() # Allows forked processes to share pytest.skip reason
processes = []
prev_start_method = mp.get_start_method()
Expand Down Expand Up @@ -252,6 +284,7 @@ def _launch_non_daemonic_procs(self, num_procs, init_method):
# Wait for all other processes to complete
for p in processes:
p.join(self.exec_timeout)
self._remove_shm_segments(num_procs)

failed = [(rank, p) for rank, p in enumerate(processes) if p.exitcode != 0]
for rank, p in failed:
Expand Down Expand Up @@ -363,6 +396,22 @@ def _close_pool(self, pool, num_procs, force=False):
pool.starmap(self._dist_destroy, [() for _ in range(num_procs)])
pool.close()
pool.join()
self._remove_shm_segments(num_procs)

def _remove_shm_segments(self, num_procs):
# The shm-based allreduce (active when LOCAL_SIZE matches the pool size)
# leaves one ~66MB segment per rank under /dev/shm keyed by the master
# port. The op never unlinks them, which exhausts /dev/shm over a full
# run, so remove them once the pool is gone.
master_port = getattr(self, '_master_port', None)
if master_port is None or not hasattr(os, 'getuid'):
return
for rank in range(num_procs):
seg = f"/dev/shm/deepspeed_allreduce_buffer_{os.getuid()}_127.0.0.1_{master_port}_{rank}"
try:
os.remove(seg)
except OSError:
pass


class DistributedFixture(DistributedExec):
Expand Down
13 changes: 13 additions & 0 deletions tests/unit/runtime/test_no_sync_ctxt.py
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@

import deepspeed
import deepspeed.comm as dist
from deepspeed.accelerator import get_accelerator
from deepspeed.utils import safe_get_full_grad


Expand All @@ -22,6 +23,10 @@ class TestNoSyncCtxt(DistributedTest):
@pytest.mark.parametrize("dtype", [torch.float16, torch.bfloat16, torch.float32])
@pytest.mark.parametrize("zero_stage", [0, 1, 2, 3])
def test_zero_stage(self, zero_stage, dtype):
# The fp16 parametrization crashes initialize's sanity check on accelerators
# without fp16 support (#8398's hardware lottery); skip instead of failing.
if dtype == torch.float16 and not get_accelerator().is_fp16_supported():
pytest.skip("fp16 is not supported on this accelerator")
config_dict = {
"train_micro_batch_size_per_gpu": 1,
"gradient_accumulation_steps": 1,
Expand Down Expand Up @@ -65,6 +70,10 @@ def test_zero_stage(self, zero_stage, dtype):
@pytest.mark.parametrize("dtype", [torch.float16, torch.bfloat16, torch.float32])
@pytest.mark.parametrize("zero_stage", [0, 1])
def test_engine_step(self, zero_stage, dtype):
# The fp16 parametrization crashes initialize's sanity check on accelerators
# without fp16 support (#8398's hardware lottery); skip instead of failing.
if dtype == torch.float16 and not get_accelerator().is_fp16_supported():
pytest.skip("fp16 is not supported on this accelerator")
config_dict = {
"train_micro_batch_size_per_gpu": 1,
"gradient_accumulation_steps": 1,
Expand Down Expand Up @@ -107,6 +116,10 @@ def test_engine_step(self, zero_stage, dtype):
@pytest.mark.parametrize("dtype", [torch.float16, torch.bfloat16, torch.float32])
@pytest.mark.parametrize("zero_stage", [0, 1])
def test_multiple_ctxts(self, zero_stage, dtype):
# The fp16 parametrization crashes initialize's sanity check on accelerators
# without fp16 support (#8398's hardware lottery); skip instead of failing.
if dtype == torch.float16 and not get_accelerator().is_fp16_supported():
pytest.skip("fp16 is not supported on this accelerator")
config_dict = {
"train_micro_batch_size_per_gpu": 1,
"gradient_accumulation_steps": 1,
Expand Down
2 changes: 2 additions & 0 deletions tests/unit/runtime/zero/test_zero_offloadpp.py
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@
import deepspeed
import torch
from deepspeed.runtime.zero.offload_config import DeepSpeedZeroOffloadOptimizerConfig
from deepspeed.accelerator import get_accelerator

import torch.nn as nn

Expand All @@ -33,6 +34,7 @@ def test_zero_partial_offload_config():


#Large sweep along hidden dim, num_layers of different sizes
@pytest.mark.skipif(not get_accelerator().is_fp16_supported(), reason="fp16 is not supported on this accelerator")
@pytest.mark.parametrize("h_dim", [1024])
@pytest.mark.parametrize("n_layers", [4, 8])
class TestZeroPartialOffloadConfigSweep(DistributedTest):
Expand Down
Loading
Loading