From 6d032589f45a34c7b68757df636551fcda3a8fe4 Mon Sep 17 00:00:00 2001 From: fzyzcjy Date: Tue, 28 Jul 2026 09:59:37 +0800 Subject: [PATCH] Refactor to remove allocate_rollout_engine_addr_and_ports_normal allocate_rollout_engine_addr_and_ports_normal is inlined into ServerGroup.start_engines: the cell list is derived from new_engine_indices, ports are probed via the underlying actor handle, and the padding that filled addressing entries up to the end of each node is dropped, since only the newly launched engines are ever initialized from it. addr_and_ports stays keyed by global rank, as before. The addressing tests drive ServerGroup.start_engines now that the allocation function they called is gone. disaggregation_bootstrap_port stays prefill-only, so a regular or decode engine neither reserves a port for it nor reports a bootstrap_port in its AddrInfo. --- miles/ray/rollout/addr_allocator.py | 69 ------- miles/ray/rollout/server_group.py | 33 +++- tests/fast/ray/rollout/conftest.py | 5 +- .../ray/rollout/real_ray/test_server_group.py | 6 +- tests/fast/ray/rollout/test_addr_allocator.py | 172 +++++++++--------- 5 files changed, 120 insertions(+), 165 deletions(-) diff --git a/miles/ray/rollout/addr_allocator.py b/miles/ray/rollout/addr_allocator.py index 1e0870acce8..a7ab12d2205 100644 --- a/miles/ray/rollout/addr_allocator.py +++ b/miles/ray/rollout/addr_allocator.py @@ -1,4 +1,3 @@ -import functools import logging from dataclasses import dataclass @@ -27,71 +26,3 @@ def alloc(self, *, engine, node_ip: str, consecutive: int = 1) -> int: ) self._values[node_ip] = port + consecutive return port - - -# NOTE: May re-implement this in a potentially easier way if needed -def allocate_rollout_engine_addr_and_ports_normal( - *, - args, - port_allocator: PortAllocator, - rollout_engines, - worker_type="regular", - num_gpus_per_engine=None, - rank_offset=0, -): - # get ports - # there are 4 ports we need to allocate - # 1. server port - # 2. nccl port - # 3. dist_init_addr port - # 4. other ports for dp_attention, which is of size 4 + dp_size - _gpus_per_engine = num_gpus_per_engine or args.rollout_num_gpus_per_engine - num_engines_per_node = max(1, args.num_gpus_per_node // _gpus_per_engine) - addr_and_ports: dict[int, dict] = {} - - visited_nodes = set() - for rank, engine in rollout_engines: - local_rank = rank - rank_offset - node_index = local_rank // num_engines_per_node - if node_index in visited_nodes: - continue - visited_nodes.add(node_index) - # TODO: currently when restarting engines, we will set port for all engines on this node starting with this rank. - # e.g. for 8 gpus, if we are restarting engine on gpu 3, we will set port for engine 3,4,5,6,7 on this node. - num_engines_on_this_node = num_engines_per_node - (local_rank % num_engines_per_node) - - node_ip, _ = ray.get(engine._get_current_node_ip_and_free_port.remote()) - - get_port = functools.partial(port_allocator.alloc, engine=engine, node_ip=node_ip) - - for i in range(num_engines_on_this_node): - current_rank = rank + i - addr_and_ports.setdefault(current_rank, {}) - addr_and_ports[current_rank]["host"] = node_ip - addr_and_ports[current_rank]["port"] = get_port() - addr_and_ports[current_rank]["nccl_port"] = get_port() - # Always allocate a unique engine_info_bootstrap_port per engine - addr_and_ports[current_rank]["engine_info_bootstrap_port"] = get_port() - - if worker_type == "prefill": - addr_and_ports[current_rank]["disaggregation_bootstrap_port"] = get_port() - - if _gpus_per_engine > args.num_gpus_per_node: - num_node_per_engine = _gpus_per_engine // args.num_gpus_per_node - if local_rank % num_node_per_engine == 0: - dist_init_addr = f"{node_ip}:{get_port(consecutive=30 + args.sglang_dp_size)}" - for i in range(num_node_per_engine): - addr_and_ports.setdefault(rank + i, {}) - addr_and_ports[rank + i]["dist_init_addr"] = dist_init_addr - else: - for i in range(num_engines_on_this_node): - addr_and_ports[rank + i][ - "dist_init_addr" - ] = f"{node_ip}:{get_port(consecutive=30 + args.sglang_dp_size)}" - - for i, _ in rollout_engines: - for key in ["port", "nccl_port", "dist_init_addr"]: - assert key in addr_and_ports[i], f"Engine {i} {key} is not set." - logger.info(f"Ports for engine {i}: {addr_and_ports[i]}") - - return addr_and_ports diff --git a/miles/ray/rollout/server_group.py b/miles/ray/rollout/server_group.py index 5b1707489a2..c7488b101aa 100644 --- a/miles/ray/rollout/server_group.py +++ b/miles/ray/rollout/server_group.py @@ -1,13 +1,15 @@ import asyncio import dataclasses +import functools import logging from typing import Any +import ray from sglang.srt.constants import GPU_MEMORY_TYPE_WEIGHTS from miles.backends.sglang_utils.sglang_engine import build_server_url from miles.backends.sglang_utils.sglang_router_api_client import SGLangRouterApiClient, use_legacy_router_api -from miles.ray.rollout.addr_allocator import PortAllocator, allocate_rollout_engine_addr_and_ports_normal +from miles.ray.rollout.addr_allocator import PortAllocator from miles.ray.rollout.server_cell import SHUTDOWN_TIMEOUT, ServerCell, flatten_cells, launch_sglang_ray_actor from miles.ray.rollout.server_engine import AddrInfo, ServerEngine from miles.utils import async_utils @@ -117,14 +119,27 @@ def start_engines( if curr_num_new_engines == 0: return [], [] - addr_and_ports = allocate_rollout_engine_addr_and_ports_normal( - args=self.args, - port_allocator=port_allocator, - rollout_engines=new_engines, - worker_type=self.worker_type, - num_gpus_per_engine=self.num_gpus_per_engine, - rank_offset=self.rank_offset, - ) + addr_and_ports: dict[int, dict[str, Any]] = {} + for cell_index in sorted({index // self.nodes_per_engine for index in new_engine_indices}): + dist_init_addr = None + for engine_in_cell_index in range(self.nodes_per_engine): + actor = self.cells[cell_index].engines[engine_in_cell_index].actor_handle + node_ip, _ = ray.get(actor._get_current_node_ip_and_free_port.remote()) + alloc = functools.partial(port_allocator.alloc, engine=actor, node_ip=node_ip) + + if engine_in_cell_index == 0: + dist_init_addr = f"{node_ip}:{alloc(consecutive=30 + self.args.sglang_dp_size)}" + + rank = self.rank_offset + cell_index * self.nodes_per_engine + engine_in_cell_index + addr_and_ports[rank] = dict( + host=node_ip, + port=alloc(), + nccl_port=alloc(), + engine_info_bootstrap_port=alloc(), + dist_init_addr=dist_init_addr, + ) + if self.worker_type == "prefill": + addr_and_ports[rank]["disaggregation_bootstrap_port"] = alloc() for index, _ in new_engines: engine_addr_and_ports = addr_and_ports[index] diff --git a/tests/fast/ray/rollout/conftest.py b/tests/fast/ray/rollout/conftest.py index 0ab54ebec56..96a81af4e33 100644 --- a/tests/fast/ray/rollout/conftest.py +++ b/tests/fast/ray/rollout/conftest.py @@ -296,8 +296,11 @@ def fake_engine(host: str = "10.0.0.1", port_seed: int = 30000) -> MagicMock: Mocks ``_get_current_node_ip_and_free_port.remote(start_port, consecutive)`` with a deterministic ``max(seq, start_port)`` counter so allocator tests - can predict and assert on port assignment.""" + can predict and assert on port assignment. It also passes + ``isinstance(x, ray.actor.ActorHandle)`` so it can be handed to + ``ServerEngine.mark_allocated_uninitialized`` (see ``fake_actor_handle``).""" e = MagicMock() + e._spec_class = ray.actor.ActorHandle e._port_cursor = port_seed def _alloc(start_port: int = 15000, consecutive: int = 1): diff --git a/tests/fast/ray/rollout/real_ray/test_server_group.py b/tests/fast/ray/rollout/real_ray/test_server_group.py index 8766ca8d5a4..196be7f43d0 100644 --- a/tests/fast/ray/rollout/real_ray/test_server_group.py +++ b/tests/fast/ray/rollout/real_ray/test_server_group.py @@ -173,10 +173,8 @@ def test_stop_handles_shutdown_failure_gracefully(self, patched_sglang_engine, p class TestStartEnginesRealAllocator: - """Drive ``start_engines`` with the real - ``allocate_rollout_engine_addr_and_ports_normal`` (no stub) so that the - actor → driver port round-trip via - ``_get_current_node_ip_and_free_port.remote`` actually runs.""" + """Drive ``start_engines`` with real actors so that the actor → driver port + round-trip via ``_get_current_node_ip_and_free_port.remote`` actually runs.""" def test_real_allocator_assigns_distinct_ports_via_remote_calls( self, diff --git a/tests/fast/ray/rollout/test_addr_allocator.py b/tests/fast/ray/rollout/test_addr_allocator.py index 4068b93f856..04469262c42 100644 --- a/tests/fast/ray/rollout/test_addr_allocator.py +++ b/tests/fast/ray/rollout/test_addr_allocator.py @@ -1,14 +1,53 @@ from __future__ import annotations -from unittest.mock import MagicMock - -from tests.fast.ray.rollout.conftest import fake_engine, make_args - -from miles.ray.rollout.addr_allocator import ( - PortAllocator, - allocate_rollout_engine_addr_and_ports_external, - allocate_rollout_engine_addr_and_ports_normal, -) +from unittest.mock import patch + +from tests.fast.ray.rollout.conftest import chunk_engines_into_cells, fake_actor_handle, fake_engine, make_args + +import miles.ray.rollout.server_group as server_group_module +from miles.ray.rollout.addr_allocator import PortAllocator +from miles.ray.rollout.server_engine import ServerEngine +from miles.ray.rollout.server_group import ServerGroup + + +def _start_engines_and_collect_addressing( + *, + args, + port_allocator: PortAllocator, + rollout_engines, + worker_type: str = "regular", + num_gpus_per_engine: int | None = None, + rank_offset: int = 0, +) -> dict[int, dict]: + """Run ``ServerGroup.start_engines`` against the given actor mocks and return, + per global rank, the kwargs its ``init`` was called with.""" + gpus_per_engine = num_gpus_per_engine or args.rollout_num_gpus_per_engine + nodes_per_engine = max(1, gpus_per_engine // args.num_gpus_per_node) + requested = dict(rollout_engines) + slots = [ServerEngine() for _ in range(max(requested) - rank_offset + 1)] + group = ServerGroup( + args=args, + pg=None, + cells=chunk_engines_into_cells( + slots, num_gpus_per_engine=gpus_per_engine, num_gpus_per_node=args.num_gpus_per_node + ), + num_gpus_per_engine=gpus_per_engine, + has_new_engines=False, + worker_type=worker_type, + rank_offset=rank_offset, + ) + for index, slot in enumerate(slots): + if rank_offset + index not in requested: + slot.mark_allocated_uninitialized(fake_actor_handle()) + started_cell_indices = sorted({(rank - rank_offset) // nodes_per_engine for rank in requested}) + + def _launch(*, global_rank, **kwargs): + return requested[global_rank] + + with patch.object(server_group_module, "launch_sglang_ray_actor", side_effect=_launch): + group.start_engines(port_allocator, start_cell_indices=started_cell_indices) + + return {rank: dict(engine.init.remote.call_args.kwargs) for rank, engine in requested.items()} class TestPortAllocator: @@ -63,12 +102,13 @@ def _all_ports(addr_and_ports: dict) -> list[int]: return out -class TestAllocateNormal: +class TestAddressingOfStartedEngines: def test_single_node_8_cards_tp1(self, patch_ray_get): + """Eight single-gpu engines on one node get complete, mutually distinct addressing.""" args = make_args(num_gpus_per_node=8, sglang_dp_size=1) engines = [(rank, fake_engine(host="10.0.0.1", port_seed=30000)) for rank in range(8)] cursors = PortAllocator.empty() - addr_and_ports = allocate_rollout_engine_addr_and_ports_normal( + addr_and_ports = _start_engines_and_collect_addressing( args=args, port_allocator=cursors, rollout_engines=engines, num_gpus_per_engine=1 ) @@ -82,7 +122,7 @@ def test_single_node_8_cards_tp1(self, patch_ray_get): host, _, port_str = addr_and_ports[rank]["dist_init_addr"].partition(":") assert host == "10.0.0.1" assert int(port_str) >= 30000 - # No same-rank collisions among the four port fields. + # No same-rank collisions among the port fields. same_rank_ports = { addr_and_ports[rank]["port"], addr_and_ports[rank]["nccl_port"], @@ -101,20 +141,19 @@ def test_single_node_8_cards_tp1(self, patch_ray_get): assert len(all_ports) == len(set(all_ports)), f"port collision across engines on the same node: {all_ports}" def test_prefill_worker_gets_disagg_bootstrap_port(self, patch_ray_get): + """A prefill engine's disaggregation bootstrap port is distinct from its other ports.""" args = make_args(num_gpus_per_node=8, sglang_dp_size=1) engines = [(rank, fake_engine()) for rank in range(2)] - addr_and_ports = allocate_rollout_engine_addr_and_ports_normal( + addr_and_ports = _start_engines_and_collect_addressing( args=args, port_allocator=PortAllocator.empty(), rollout_engines=engines, worker_type="prefill", num_gpus_per_engine=1, ) - for rank in range(2): - assert isinstance(addr_and_ports[rank]["disaggregation_bootstrap_port"], int) - # The disagg port must be distinct from the other ports on the same rank. for rank in range(2): entry = addr_and_ports[rank] + assert isinstance(entry["disaggregation_bootstrap_port"], int) assert entry["disaggregation_bootstrap_port"] not in ( entry["port"], entry["nccl_port"], @@ -122,45 +161,42 @@ def test_prefill_worker_gets_disagg_bootstrap_port(self, patch_ray_get): ) def test_regular_worker_does_not_get_disagg_bootstrap_port(self, patch_ray_get): + """Only prefill engines carry a disaggregation bootstrap port; others must not reserve one.""" args = make_args(num_gpus_per_node=8, sglang_dp_size=1) engines = [(rank, fake_engine()) for rank in range(2)] - addr_and_ports = allocate_rollout_engine_addr_and_ports_normal( + addr_and_ports = _start_engines_and_collect_addressing( args=args, port_allocator=PortAllocator.empty(), rollout_engines=engines, num_gpus_per_engine=1 ) for rank in range(2): assert "disaggregation_bootstrap_port" not in addr_and_ports[rank] def test_gpus_per_engine_greater_than_node_shares_dist_init_addr(self, patch_ray_get): - """When `_gpus_per_engine > num_gpus_per_node`, all ranks of one engine + """When `gpus_per_engine > num_gpus_per_node`, all ranks of one engine share a single ``dist_init_addr`` (multi-node engine).""" args = make_args(num_gpus_per_node=8, sglang_dp_size=1) # 2-node engine: 16 gpus total, 8 per node, 2 ranks share dist_init_addr engines = [(rank, fake_engine(host="10.0.0.42")) for rank in range(2)] - addr_and_ports = allocate_rollout_engine_addr_and_ports_normal( + addr_and_ports = _start_engines_and_collect_addressing( args=args, port_allocator=PortAllocator.empty(), rollout_engines=engines, num_gpus_per_engine=16 ) - # Same string for both ranks — not just equal, identity of representation. assert addr_and_ports[0]["dist_init_addr"] == addr_and_ports[1]["dist_init_addr"] host, _, port_str = addr_and_ports[0]["dist_init_addr"].partition(":") assert host == "10.0.0.42" assert int(port_str) > 0 def test_rank_offset_does_not_break_indexing(self, patch_ray_get): + """A group starting at rank 4 populates exactly ranks 4..7 with collision-free ports.""" args = make_args(num_gpus_per_node=8, sglang_dp_size=1) engines = [(rank, fake_engine(host="10.0.0.7", port_seed=40000)) for rank in (4, 5, 6, 7)] - addr_and_ports = allocate_rollout_engine_addr_and_ports_normal( + addr_and_ports = _start_engines_and_collect_addressing( args=args, port_allocator=PortAllocator.empty(), rollout_engines=engines, num_gpus_per_engine=1, rank_offset=4, ) - # Allocator fills remaining slots on the node starting from rank 4 - # (see source comment: "we will set port for engine 3,4,5,6,7 on this - # node"); so the requested ranks {4,5,6,7} must be a subset, with no - # leakage into ranks 0..3. - assert {4, 5, 6, 7} <= set(addr_and_ports.keys()) - assert set(addr_and_ports.keys()).isdisjoint({0, 1, 2, 3}) + # Exactly the requested ranks are populated; no leakage into 0..3. + assert set(addr_and_ports.keys()) == {4, 5, 6, 7} for r in (4, 5, 6, 7): assert addr_and_ports[r]["host"] == "10.0.0.7" assert isinstance(addr_and_ports[r]["port"], int) @@ -169,68 +205,44 @@ def test_rank_offset_does_not_break_indexing(self, patch_ray_get): all_ports = _all_ports(addr_and_ports) assert len(all_ports) == len(set(all_ports)) - def test_mid_rank_restart_fills_remaining_slots_on_node(self, patch_ray_get): - """Restarting starting from rank 3 on an 8-card node should populate - addr_and_ports[3..7], i.e. ``num_engines_on_this_node`` accounts for - the offset within the node.""" + def test_mid_rank_restart_populates_only_the_requested_rank(self, patch_ray_get): + """Restarting rank 3 must allocate for rank 3 alone; other slots on the + node keep their existing engines and ports.""" args = make_args(num_gpus_per_node=8, sglang_dp_size=1) engines = [(3, fake_engine())] - addr_and_ports = allocate_rollout_engine_addr_and_ports_normal( + addr_and_ports = _start_engines_and_collect_addressing( args=args, port_allocator=PortAllocator.empty(), rollout_engines=engines, num_gpus_per_engine=1 ) - # Exact key set: mid-rank fill must cover [3..7] and ONLY those. - assert set(addr_and_ports.keys()) == {3, 4, 5, 6, 7} - # Each filled slot has a complete addr/port set. - for r in range(3, 8): - for k in ("host", "port", "nccl_port", "engine_info_bootstrap_port", "dist_init_addr"): - assert k in addr_and_ports[r] + assert set(addr_and_ports.keys()) == {3} + for k in ("host", "port", "nccl_port", "engine_info_bootstrap_port", "dist_init_addr"): + assert k in addr_and_ports[3] def test_cursor_ends_past_every_issued_port(self, patch_ray_get): + """The node cursor is left beyond every port handed out, including reserved blocks.""" args = make_args(num_gpus_per_node=8, sglang_dp_size=1) engines = [(0, fake_engine(port_seed=22000))] cursors = PortAllocator.empty() - addr_and_ports = allocate_rollout_engine_addr_and_ports_normal( + addr_and_ports = _start_engines_and_collect_addressing( args=args, port_allocator=cursors, rollout_engines=engines, num_gpus_per_engine=1 ) - # Cursor must sit strictly past every port we handed out (the allocator + # Cursor must sit strictly past every port we handed out (the allocation # also reserves consecutive blocks for dist_init_addr that aren't all # visible in the output, so we can't pin to max_issued + 1). max_issued = max(_all_ports(addr_and_ports)) assert cursors._values["10.0.0.1"] > max_issued -class TestAllocateExternal: - def test_basic_split(self): - args = make_args(rollout_external_engine_addrs=["10.0.0.1:30000", "10.0.0.2:30001"]) - engines = [(0, MagicMock()), (1, MagicMock())] - result = allocate_rollout_engine_addr_and_ports_external(args=args, rollout_engines=engines) - # Whole-dict equality (not just a key check) — pins the entire shape. - assert result == { - 0: dict(dist_init_addr="10.0.0.1:30000", nccl_port=None, host="10.0.0.1", port=30000), - 1: dict(dist_init_addr="10.0.0.2:30001", nccl_port=None, host="10.0.0.2", port=30001), - } - - def test_ipv4_addr_split_is_consistent(self): - args = make_args(rollout_external_engine_addrs=["192.168.1.10:31000"]) - engines = [(0, MagicMock())] - result = allocate_rollout_engine_addr_and_ports_external(args=args, rollout_engines=engines) - assert result == { - 0: dict(dist_init_addr="192.168.1.10:31000", nccl_port=None, host="192.168.1.10", port=31000), - } - # `port` is an int, not a string accidentally split through. - assert isinstance(result[0]["port"], int) - - class TestSharedPortAllocatorAcrossGroups: """Two ``ServerGroup``s sharing one ``PortAllocator`` must produce disjoint port allocations across nodes — required for parallel recover.""" def test_sequential_groups_share_cursor_and_avoid_overlap(self, patch_ray_get): + """Groups started one after another off a shared allocator never reuse a port.""" args = make_args(num_gpus_per_node=8, sglang_dp_size=1) cursors = PortAllocator.empty() engines_a = [(rank, fake_engine(port_seed=0)) for rank in range(4)] - addrs_a = allocate_rollout_engine_addr_and_ports_normal( + addrs_a = _start_engines_and_collect_addressing( args=args, port_allocator=cursors, rollout_engines=engines_a, @@ -238,7 +250,7 @@ def test_sequential_groups_share_cursor_and_avoid_overlap(self, patch_ray_get): ) engines_b = [(rank, fake_engine(port_seed=0)) for rank in range(4, 8)] - addrs_b = allocate_rollout_engine_addr_and_ports_normal( + addrs_b = _start_engines_and_collect_addressing( args=args, port_allocator=cursors, rollout_engines=engines_b, @@ -252,22 +264,17 @@ def test_sequential_groups_share_cursor_and_avoid_overlap(self, patch_ray_get): class TestRankPortConsistency: - """rank ↔ addr_and_ports index consistency in ``ServerGroup.start_engines``. + """rank ↔ addr_and_ports consistency inside ``ServerGroup.start_engines``. - The init-handles loop iterates ``new_engines`` as ``(global_rank, engine)`` - pairs while the allocator keys its output dict on ``rank + i``. When - ``rank_offset != 0`` or ``nodes_per_engine > 1`` the two index spaces - must still agree.""" + The init loop iterates ``new_engines`` as ``(global_rank, engine)`` pairs, so + the addressing must be keyed by global rank even when ``rank_offset != 0`` or + ``nodes_per_engine > 1``.""" def test_rank_offset_kwargs_keyed_by_global_rank(self, patch_ray_get): - """When rank_offset=4, addr_and_ports must be keyed by ranks 4..7, - not 0..3. - - ``num_gpus_per_node=4`` so a 4-engine group exactly fills one node — - otherwise the allocator pads up to ``num_engines_per_node``.""" + """When rank_offset=4, addr_and_ports must be keyed by ranks 4..7, not 0..3.""" args = make_args(num_gpus_per_node=4, sglang_dp_size=1) engines = [(rank, fake_engine(port_seed=0)) for rank in range(4, 8)] - addr_and_ports = allocate_rollout_engine_addr_and_ports_normal( + addr_and_ports = _start_engines_and_collect_addressing( args=args, port_allocator=PortAllocator.empty(), rollout_engines=engines, @@ -277,9 +284,10 @@ def test_rank_offset_kwargs_keyed_by_global_rank(self, patch_ray_get): assert set(addr_and_ports.keys()) == {4, 5, 6, 7} def test_each_global_rank_has_complete_kwargs(self, patch_ray_get): + """Every started rank receives the full addressing kwarg set.""" args = make_args(num_gpus_per_node=8, sglang_dp_size=1) engines = [(rank, fake_engine(port_seed=0)) for rank in range(4)] - addr_and_ports = allocate_rollout_engine_addr_and_ports_normal( + addr_and_ports = _start_engines_and_collect_addressing( args=args, port_allocator=PortAllocator.empty(), rollout_engines=engines, @@ -295,7 +303,7 @@ def test_multinode_engine_shares_dist_init_addr_across_node_ranks(self, patch_ra multi-node engine MUST get the same dist_init_addr.""" args = make_args(num_gpus_per_node=8, sglang_dp_size=1) engines = [(0, fake_engine(port_seed=0)), (1, fake_engine(port_seed=0))] - addr_and_ports = allocate_rollout_engine_addr_and_ports_normal( + addr_and_ports = _start_engines_and_collect_addressing( args=args, port_allocator=PortAllocator.empty(), rollout_engines=engines, @@ -303,12 +311,11 @@ def test_multinode_engine_shares_dist_init_addr_across_node_ranks(self, patch_ra ) assert addr_and_ports[0]["dist_init_addr"] == addr_and_ports[1]["dist_init_addr"] - def test_init_handles_iteration_pairs_match_addr_dict(self, patch_ray_get): - """For every (index, engine) pair in server_group.py's loop, - addr_and_ports[index] must contain all required kwargs.""" + def test_init_kwargs_exist_for_every_started_rank(self, patch_ray_get): + """For every (rank, engine) pair the init loop walks, the addressing dict has an entry.""" args = make_args(num_gpus_per_node=8, sglang_dp_size=1) new_engines = [(rank, fake_engine(port_seed=0)) for rank in range(2, 6)] - addr_and_ports = allocate_rollout_engine_addr_and_ports_normal( + addr_and_ports = _start_engines_and_collect_addressing( args=args, port_allocator=PortAllocator.empty(), rollout_engines=new_engines, @@ -321,9 +328,10 @@ def test_init_handles_iteration_pairs_match_addr_dict(self, patch_ray_get): assert key in addr_and_ports[index] def test_ports_are_unique_within_a_node(self, patch_ray_get): + """No two engines on the same node share any of their allocated ports.""" args = make_args(num_gpus_per_node=8, sglang_dp_size=1) engines = [(rank, fake_engine(port_seed=0)) for rank in range(8)] - addr_and_ports = allocate_rollout_engine_addr_and_ports_normal( + addr_and_ports = _start_engines_and_collect_addressing( args=args, port_allocator=PortAllocator.empty(), rollout_engines=engines,