diff --git a/tensorrt_llm/_torch/disaggregation/executor/admission.py b/tensorrt_llm/_torch/disaggregation/executor/admission.py deleted file mode 100644 index f940a7e09378..000000000000 --- a/tensorrt_llm/_torch/disaggregation/executor/admission.py +++ /dev/null @@ -1,114 +0,0 @@ -# SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. -# SPDX-License-Identifier: Apache-2.0 - -import dataclasses -from typing import Iterable, List, Optional - -from tensorrt_llm._torch.pyexecutor.llm_request import LlmRequest - - -@dataclasses.dataclass -class DisaggTransferAdmissionResult: - admitted_requests: List[LlmRequest] - active_transfer_blocks: int = 0 - admitted_transfer_blocks: int = 0 - deferred_request_count: int = 0 - limited_by_budget: bool = False - - def is_blocked_by_active_transfers(self) -> bool: - return ( - self.limited_by_budget - and not self.admitted_requests - and self.active_transfer_blocks > 0 - ) - - -class DisaggTransferAdmissionController: - """FCFS admission gate for disaggregated generation KV transfers.""" - - def __init__( - self, max_tokens_in_buffer: Optional[int], tokens_per_block: Optional[int] - ) -> None: - self.max_transfer_blocks = self._to_block_budget(max_tokens_in_buffer, tokens_per_block) - self.tokens_per_block = tokens_per_block or 0 - - def enabled(self) -> bool: - return self.max_transfer_blocks is not None - - @staticmethod - def _to_block_budget( - max_tokens_in_buffer: Optional[int], tokens_per_block: Optional[int] - ) -> Optional[int]: - if ( - max_tokens_in_buffer is None - or max_tokens_in_buffer == 0 - or tokens_per_block is None - or tokens_per_block <= 0 - ): - return None - return (max_tokens_in_buffer + tokens_per_block - 1) // tokens_per_block - - @staticmethod - def _to_nonnegative_int(value) -> Optional[int]: - try: - return max(int(value), 0) - except (TypeError, ValueError): - return None - - def _get_request_transfer_token_count(self, request: LlmRequest) -> int: - for attr_name in ("total_input_len_cp", "py_prompt_len", "prompt_len"): - token_count = self._to_nonnegative_int(getattr(request, attr_name, None)) - if token_count is not None: - return token_count - return 0 - - def _estimate_request_blocks(self, request: LlmRequest) -> int: - if self.tokens_per_block <= 0: - return 0 - prompt_len = self._get_request_transfer_token_count(request) - return (prompt_len + self.tokens_per_block - 1) // self.tokens_per_block - - def _estimate_requests_blocks(self, requests: Iterable[LlmRequest]) -> int: - return sum(self._estimate_request_blocks(request) for request in requests) - - def _estimate_active_transfer_blocks(self, active_requests: Iterable[LlmRequest]) -> int: - return sum( - self._estimate_request_blocks(request) - for request in active_requests - if request.is_disagg_generation_transmission_in_progress - ) - - def select( - self, active_requests: Iterable[LlmRequest], candidates: List[LlmRequest] - ) -> DisaggTransferAdmissionResult: - if not self.enabled(): - return DisaggTransferAdmissionResult( - admitted_requests=list(candidates), - active_transfer_blocks=self._estimate_active_transfer_blocks(active_requests), - admitted_transfer_blocks=self._estimate_requests_blocks(candidates), - ) - - result = DisaggTransferAdmissionResult(admitted_requests=[]) - result.active_transfer_blocks = self._estimate_active_transfer_blocks(active_requests) - - used_blocks = result.active_transfer_blocks - max_transfer_blocks = self.max_transfer_blocks - assert max_transfer_blocks is not None - for request in candidates: - request_blocks = self._estimate_request_blocks(request) - fits_budget = used_blocks + request_blocks <= max_transfer_blocks - admit_oversized_head = ( - not result.admitted_requests - and result.active_transfer_blocks == 0 - and request_blocks > max_transfer_blocks - ) - if not fits_budget and not admit_oversized_head: - result.limited_by_budget = True - break - - result.admitted_requests.append(request) - used_blocks += request_blocks - result.admitted_transfer_blocks += request_blocks - - result.deferred_request_count = len(candidates) - len(result.admitted_requests) - return result diff --git a/tensorrt_llm/_torch/disaggregation/executor/coordinator.py b/tensorrt_llm/_torch/disaggregation/executor/coordinator.py index a2a2b9918901..d9c264fa5dc8 100644 --- a/tensorrt_llm/_torch/disaggregation/executor/coordinator.py +++ b/tensorrt_llm/_torch/disaggregation/executor/coordinator.py @@ -8,7 +8,7 @@ """ from dataclasses import dataclass, fields -from typing import TYPE_CHECKING, Callable, List, Tuple +from typing import TYPE_CHECKING, Callable, List from tensorrt_llm._torch.pyexecutor.llm_request import LlmRequest @@ -28,7 +28,6 @@ class DisaggLoopDelegates: prepare_context_schedulable: Callable[[List[LlmRequest]], None] poll_gen_transfers: Callable[[], None] check_transfer_timeouts: Callable[[], None] - admit: Callable[[List[LlmRequest]], Tuple[List[LlmRequest], bool]] revert_deferred_gen_init: Callable[[List[LlmRequest], List[LlmRequest]], None] receive_gen_init: Callable[[List[LlmRequest]], None] poll_progress_when_idle: Callable[[], None] @@ -71,13 +70,6 @@ def check_transfer_timeouts(self) -> None: # -- scheduling ---------------------------------------------------------- - def admit(self, fitting_gen_init: List[LlmRequest]) -> Tuple[List[LlmRequest], bool]: - """Select the gen-init requests that may start receiving this iteration. - - Returns ``(admitted, blocked_by_active_transfers)``. - """ - return self._d.admit(fitting_gen_init) - def revert_deferred_gen_init( self, candidates: List[LlmRequest], admitted: List[LlmRequest] ) -> None: @@ -121,9 +113,6 @@ def __init__(self) -> None: DisaggLoopDelegates(**{f.name: _noop for f in fields(DisaggLoopDelegates)}) ) - def admit(self, fitting_gen_init: List[LlmRequest]) -> Tuple[List[LlmRequest], bool]: - return fitting_gen_init, False - def _noop(*_args, **_kwargs) -> None: return None diff --git a/tensorrt_llm/_torch/disaggregation/kv_cache_transceiver.py b/tensorrt_llm/_torch/disaggregation/kv_cache_transceiver.py index 9ee46e9f9ef2..b83758bd4c53 100644 --- a/tensorrt_llm/_torch/disaggregation/kv_cache_transceiver.py +++ b/tensorrt_llm/_torch/disaggregation/kv_cache_transceiver.py @@ -302,11 +302,6 @@ def has_retired_send_session(self, req: LlmRequest) -> bool: """Whether the send session closed before its final slice.""" return False - @property - def consumes_transfer_buffer(self) -> bool: - """Return whether this runtime consumes the C++ CacheTransBuffer budget.""" - return True - @abstractmethod def respond_and_send_async(self, req: LlmRequest) -> None: """Start sending ``req``'s KV cache to the requesting instance. diff --git a/tensorrt_llm/_torch/disaggregation/transceiver.py b/tensorrt_llm/_torch/disaggregation/transceiver.py index e2811bf885b6..7ad42b6e2ed6 100644 --- a/tensorrt_llm/_torch/disaggregation/transceiver.py +++ b/tensorrt_llm/_torch/disaggregation/transceiver.py @@ -82,10 +82,6 @@ def _find_consensus_request_ids(request_ids_all_ranks, sync_size): class KvCacheTransceiverV2(KvCacheTransceiver): - @property - def consumes_transfer_buffer(self) -> bool: - return False - def __init__( self, mapping: Mapping, diff --git a/tensorrt_llm/_torch/pyexecutor/py_executor.py b/tensorrt_llm/_torch/pyexecutor/py_executor.py index dcb09edd5814..b78f1b8d1e04 100644 --- a/tensorrt_llm/_torch/pyexecutor/py_executor.py +++ b/tensorrt_llm/_torch/pyexecutor/py_executor.py @@ -53,8 +53,6 @@ get_global_profiler, host_profiler_context) from ..disaggregation.base.transfer import get_unique_rid -from ..disaggregation.executor.admission import \ - DisaggTransferAdmissionController from ..disaggregation.executor.coordinator import (DisaggLoopDelegates, DisaggTransferCoordinator, NoopDisaggCoordinator) @@ -935,24 +933,6 @@ def on_detected(): if kv_cache_transceiver is not None: self.hang_detector.register_status_provider( kv_cache_transceiver.get_status_dump) - cache_transceiver_config = getattr(self.llm_args, - "cache_transceiver_config", None) - max_tokens_in_buffer = getattr(cache_transceiver_config, - "max_tokens_in_buffer", None) - tokens_per_block = getattr(self.kv_cache_manager, "tokens_per_block", - None) - self._disagg_transfer_admission_controller = DisaggTransferAdmissionController( - max_tokens_in_buffer, tokens_per_block) - if (self.global_rank == 0 - and self._disagg_transfer_admission_controller.enabled() - and self._is_disagg_transfer_window_bypass_eligible()): - logger.warning_once( - f"[PyExecutor] Bypassing the executor transfer window " - f"configured by max_tokens_in_buffer={max_tokens_in_buffer} " - "for asynchronous Python generation with KV cache manager " - "V2 and pp_size=1; " - "scheduler KV cache capacity admission remains active.", - key="disagg_transfer_window_bypass") self.is_benchmark_disagg = (self.benchmark_req_queues_size > 0 and self.kv_cache_transceiver is not None) # True while the benchmark disagg fill phase is in progress (waiting @@ -2569,17 +2549,14 @@ def _pp_schedule_and_propagate(self, microbatch_id: int): # For DP cases, the first PP rank schedules the requests. scheduled_batch = None serializable_schedule = None - wait_for_disagg_gen_transfer_progress = False is_dp_broadcast = self.dist.tp_size > 1 and self.enable_attention_dp if self.dist.rank == 0 or (self.dist.is_first_pp_rank and is_dp_broadcast): scheduled_batch, fitting_disagg_gen_init_requests, num_fitting_reqs = self._schedule( ) - fitting_disagg_gen_init_requests, wait_for_disagg_gen_transfer_progress = ( - self.disagg.admit(fitting_disagg_gen_init_requests)) serializable_schedule = SerializableSchedulerOutput.from_scheduler_result( scheduled_batch, fitting_disagg_gen_init_requests, - num_fitting_reqs, wait_for_disagg_gen_transfer_progress) + num_fitting_reqs) # Broadcast within first tp+cp group before send/recv chain to other tp+cp groups if self.dist.is_first_pp_rank: @@ -2611,10 +2588,8 @@ def _pp_schedule_and_propagate(self, microbatch_id: int): if scheduled_batch is None: scheduled_batch, fitting_disagg_gen_init_requests, num_fitting_reqs = serializable_schedule.to_scheduler_result( self.active_requests) - wait_for_disagg_gen_transfer_progress = ( - serializable_schedule.wait_for_disagg_gen_transfer_progress) return (scheduled_batch, fitting_disagg_gen_init_requests, - num_fitting_reqs, wait_for_disagg_gen_transfer_progress) + num_fitting_reqs) def _pp_retry_until_can_schedule(self, scheduled_batch): """ @@ -2704,7 +2679,7 @@ def _executor_loop_pp(self): self._pad_attention_dp_dummy_request() # Stage 0: first PP rank schedules requests and propagates the result to all other PP ranks. - (scheduled_batch, fitting_disagg_gen_init_requests, _, + (scheduled_batch, fitting_disagg_gen_init_requests, _) = self._pp_schedule_and_propagate(microbatch_id) if self.dist.rank != 0: # Retry until current rank can run first PP's schedule result. @@ -3693,7 +3668,6 @@ def _build_disagg_coordinator(self) -> DisaggTransferCoordinator: _check_disagg_ctx_schedulable_status, poll_gen_transfers=self._check_disagg_gen_transfer_status, check_transfer_timeouts=self._check_kv_transfer_timeout, - admit=self._apply_disagg_transfer_admission, revert_deferred_gen_init=self. _revert_deferred_disagg_gen_init_alloc, receive_gen_init=self._prepare_disagg_gen_init, @@ -3706,42 +3680,6 @@ def _build_disagg_coordinator(self) -> DisaggTransferCoordinator: pace_idle=self._pace_idle_disagg_loop, )) - def _get_disagg_transfer_admission_controller( - self) -> DisaggTransferAdmissionController: - controller = getattr(self, "_disagg_transfer_admission_controller", - None) - if controller is not None: - return controller - - cache_transceiver_config = getattr(getattr(self, "llm_args", None), - "cache_transceiver_config", None) - kv_cache_manager = getattr(self, "kv_cache_manager", None) - return DisaggTransferAdmissionController( - getattr(cache_transceiver_config, "max_tokens_in_buffer", None), - getattr(kv_cache_manager, "tokens_per_block", None), - ) - - def _is_disagg_transfer_window_bypass_eligible(self) -> bool: - """Return whether this runtime may bypass an enabled transfer window.""" - transceiver = getattr(self, "kv_cache_transceiver", None) - return (transceiver is not None - and transceiver.consumes_transfer_buffer is False - and self._uses_async_disagg_gen_transfer() - and self.dist.pp_size == 1 and self._is_kv_manager_v2) - - def _disagg_transfer_window_is_active(self) -> bool: - """Return whether the executor-level transfer window is active. - - ``max_tokens_in_buffer`` describes the C++ transceiver's physical - buffer. The asynchronous Python transceiver does not consume that - buffer. With KV cache manager V2 and PP1, its generation requests - remain constrained by inline scheduler KV admission without a second - executor-level budget. Other configurations retain the window. - """ - return (getattr(self, "kv_cache_transceiver", None) is not None - and self._get_disagg_transfer_admission_controller().enabled() - and not self._is_disagg_transfer_window_bypass_eligible()) - @staticmethod def _is_disagg_gen_only_no_context_benchmark() -> bool: """Return whether ``gen_only_no_context`` skips KV transfer.""" @@ -3752,47 +3690,14 @@ def _uses_async_disagg_gen_transfer(self) -> bool: return (not self._is_disagg_gen_only_no_context_benchmark() and os.getenv("TRTLLM_DISABLE_KV_CACHE_TRANSFER_OVERLAP") != "1") - def _apply_disagg_transfer_admission( - self, fitting_disagg_gen_init_requests: List[LlmRequest] - ) -> Tuple[List[LlmRequest], bool]: - # gen_only_no_context has no CTX worker and does not transfer data. - # Real synchronous gen_only transfers still honor the budget to bound - # the number of blocking transfers started in one executor iteration. - if self._is_disagg_gen_only_no_context_benchmark(): - return fitting_disagg_gen_init_requests, False - - controller = self._get_disagg_transfer_admission_controller() - if not (self._disagg_transfer_window_is_active() - and fitting_disagg_gen_init_requests): - return fitting_disagg_gen_init_requests, False - - admission_result = controller.select(self.active_requests, - fitting_disagg_gen_init_requests) - if admission_result.deferred_request_count > 0: - logger.debug("Disagg transfer admission deferred " - f"{admission_result.deferred_request_count} requests; " - f"active transfer blocks=" - f"{admission_result.active_transfer_blocks}, " - f"admitted transfer blocks=" - f"{admission_result.admitted_transfer_blocks}, " - f"budget={controller.max_transfer_blocks}") - - self._revert_deferred_disagg_gen_init_alloc( - fitting_disagg_gen_init_requests, - admission_result.admitted_requests) - - return (admission_result.admitted_requests, - admission_result.is_blocked_by_active_transfers()) - def _revert_deferred_disagg_gen_init_alloc( self, candidates: List[LlmRequest], admitted_requests: List[LlmRequest]) -> None: """Revert Scheduler V2 allocations absent from an admitted request set. Scheduler V2 allocates KV while evaluating generation-init requests. - This reconciliation is required both after transfer-window admission - and when a PP follower's local candidates differ from the canonical - schedule propagated by rank 0. + This reconciliation is required when a PP follower's local candidates + differ from the canonical schedule propagated by rank 0. """ if not (self._is_kv_manager_v2 and candidates): return @@ -3909,8 +3814,7 @@ def _pp_ring_is_drained(self) -> bool: and all(batch is None for batch in self.micro_batches)) def _sync_gen_only_benchmark_has_insufficient_kv( - self, scheduler_fitting_disagg_gen_init_requests: List[LlmRequest], - wait_for_disagg_gen_transfer_progress: bool) -> bool: + self, fitting_disagg_gen_init_requests: List[LlmRequest]) -> bool: """Return whether benchmark fill has terminal KV exhaustion. Model-parallel ranks can make different local scheduling decisions. @@ -3920,18 +3824,13 @@ def _sync_gen_only_benchmark_has_insufficient_kv( to every decode iteration after the gate opens. Args: - scheduler_fitting_disagg_gen_init_requests: Generation INIT - requests that fit KV capacity before transfer admission. A - nonempty list means KV capacity exists even if transfer - admission temporarily defers every request. - wait_for_disagg_gen_transfer_progress: Whether active generation - transfers are consuming the admission budget and transfer - progress can unblock a deferred request. + fitting_disagg_gen_init_requests: Generation INIT requests that + fit KV capacity. A nonempty list means KV capacity exists. Returns: True when every TP+CP rank has fetched its full benchmark queue and - at least one rank has an INIT request that cannot fit KV capacity - and has no transfer progress that can unblock it; otherwise False. + at least one rank has an INIT request that cannot fit KV capacity; + otherwise False. """ if (self.benchmark_req_queues_size <= 0 or self.is_warmup or not self._benchmark_fill_phase_active): @@ -3941,9 +3840,8 @@ def _sync_gen_only_benchmark_has_insufficient_kv( for req in self.active_requests) local_all_fetched = (self.num_fetch_requests >= self.benchmark_req_queues_size) - local_terminal_no_fit = (local_has_stuck and - not scheduler_fitting_disagg_gen_init_requests - and not wait_for_disagg_gen_transfer_progress) + local_terminal_no_fit = (local_has_stuck + and not fitting_disagg_gen_init_requests) local_status = (local_all_fetched, local_terminal_no_fit) all_rank_status = self._allgather_model_parallel_status(local_status) @@ -4043,8 +3941,7 @@ def _prepare_and_schedule_batch(self): request.py_draft_tokens = [0] * self.max_total_draft_tokens request.draft_tokens = [0] * self.max_total_draft_tokens - scheduled_batch, scheduler_fitting_disagg_gen_init_requests, _ = self._schedule( - ) + scheduled_batch, fitting_disagg_gen_init_requests, _ = self._schedule() # Must run after _schedule(): the empty scheduled batch it repairs does # not exist until the capacity scheduler has returned its verdict. @@ -4055,12 +3952,9 @@ def _prepare_and_schedule_batch(self): request.py_disable_speculative_decoding = True if self.kv_cache_transceiver: - wait_for_disagg_gen_transfer_progress = False - admitted_disagg_gen_init_requests, wait_for_disagg_gen_transfer_progress = ( - self.disagg.admit(scheduler_fitting_disagg_gen_init_requests)) - # Prepare KV cache manager resources only for requests admitted - # into the transfer window this iteration. - self.disagg.receive_gen_init(admitted_disagg_gen_init_requests) + # For requests that are fitting disagg gen init, also prepare + # resources for KV cache manager. + self.disagg.receive_gen_init(fitting_disagg_gen_init_requests) self.disagg.poll_progress_when_idle() @@ -4069,12 +3963,8 @@ def _prepare_and_schedule_batch(self): # scheduler could not allocate KV for any of them, the benchmark # will hang forever because in-progress generation requests won't # release their KV cache. - # Check the scheduler result from before transfer admission. An - # empty admitted list can mean that active transfers are - # temporarily consuming the transfer budget. has_insufficient_kv = self._sync_gen_only_benchmark_has_insufficient_kv( - scheduler_fitting_disagg_gen_init_requests, - wait_for_disagg_gen_transfer_progress) + fitting_disagg_gen_init_requests) if has_insufficient_kv: error_msg = ( f"Insufficient KV cache for gen-only benchmark mode: " diff --git a/tensorrt_llm/_torch/pyexecutor/scheduler/scheduler.py b/tensorrt_llm/_torch/pyexecutor/scheduler/scheduler.py index ff180cca75fb..f62018566b37 100644 --- a/tensorrt_llm/_torch/pyexecutor/scheduler/scheduler.py +++ b/tensorrt_llm/_torch/pyexecutor/scheduler/scheduler.py @@ -343,7 +343,6 @@ class SerializableSchedulerOutput: int ] # request ids of fitting disaggregated generation initialization requests num_fitting_requests: int # number of fitting requests - wait_for_disagg_gen_transfer_progress: bool = False # request id -> prompt-ordered indices of its MM items selected for # encoder execution this iteration scheduled_mm_encoder_items: dict[int, list[int]] | None = None @@ -356,7 +355,6 @@ def from_scheduler_result( scheduled_requests: ScheduledRequests, fitting_disagg_gen_init_requests: RequestList, num_fitting_requests: int, - wait_for_disagg_gen_transfer_progress: bool = False, ) -> "SerializableSchedulerOutput": return cls( encoder_requests=[req.request_id for req in scheduled_requests.encoder_requests], @@ -372,7 +370,6 @@ def from_scheduler_result( req.request_id for req in fitting_disagg_gen_init_requests ], num_fitting_requests=num_fitting_requests, - wait_for_disagg_gen_transfer_progress=wait_for_disagg_gen_transfer_progress, scheduled_mm_encoder_items=scheduled_requests.scheduled_mm_encoder_items, recompute_paused_requests=[ req.request_id for req in scheduled_requests.recompute_paused_requests diff --git a/tests/unittest/_torch/disaggregation/test_benchmark_disagg.py b/tests/unittest/_torch/disaggregation/test_benchmark_disagg.py index 79eabcd99193..68a23118db59 100644 --- a/tests/unittest/_torch/disaggregation/test_benchmark_disagg.py +++ b/tests/unittest/_torch/disaggregation/test_benchmark_disagg.py @@ -23,7 +23,7 @@ - ADP dummy suppression during fill vs taper-down - ADP router imbalance regression (nvbug 6071070) - Non-blocking behaviour of `_prepare_and_schedule_batch` -- Insufficient-KV fail-fast vs transfer-admission backpressure +- Insufficient-KV fail-fast during the benchmark fill phase """ from unittest.mock import Mock, patch @@ -1358,20 +1358,16 @@ def test_healthy_fill_phase_does_not_kill(self): ) ex._handle_errors.assert_not_called() - def test_partial_transfer_admission_uses_only_admitted_requests(self) -> None: - """The admitted subset is prepared and passed to the idle check.""" - admitted_req = _make_active_request(in_init=True) - deferred_req = _make_active_request(in_init=True) - candidates = [admitted_req, deferred_req] - ex = self._make_executor(fill_phase_active=True, fitting_init_requests=candidates) - ex._apply_disagg_transfer_admission = Mock(return_value=([admitted_req], False)) + def test_fitting_init_requests_are_prepared_for_transfer(self) -> None: + """Every scheduler-fitting INIT request is prepared for KV transfer.""" + fitting = [_make_active_request(in_init=True), _make_active_request(in_init=True)] + ex = self._make_executor(fill_phase_active=True, fitting_init_requests=fitting) ex._check_disagg_transfer_progress_when_idle = Mock() result, _ = ex._prepare_and_schedule_batch() assert result is not None - ex._apply_disagg_transfer_admission.assert_called_once_with(candidates) - ex._prepare_disagg_gen_init.assert_called_once_with([admitted_req]) + ex._prepare_disagg_gen_init.assert_called_once_with(fitting) ex._check_disagg_transfer_progress_when_idle.assert_called_once_with() ex._handle_errors.assert_not_called() @@ -1384,30 +1380,6 @@ def test_fill_with_no_init_requests_does_not_kill(self): assert result is not None ex._handle_errors.assert_not_called() - def test_transfer_admission_backpressure_does_not_kill(self, monkeypatch): - """NVBug 6438658: admission backpressure is not KV exhaustion. - - Args: - monkeypatch: Pytest fixture used to select asynchronous transfer - behavior. - """ - monkeypatch.delenv("TRTLLM_DISAGG_BENCHMARK_GEN_ONLY", raising=False) - monkeypatch.delenv("TRTLLM_DISABLE_KV_CACHE_TRANSFER_OVERLAP", raising=False) - fitting_req = _make_active_request(in_init=True) - ex = self._make_executor(fill_phase_active=True, fitting_init_requests=[fitting_req]) - ex._apply_disagg_transfer_admission = Mock(return_value=([], True)) - - result, _ = ex._prepare_and_schedule_batch() - - assert result is not None, ( - "Fail-fast should NOT fire when the scheduler fit an INIT request " - "that transfer admission temporarily deferred" - ) - ex._apply_disagg_transfer_admission.assert_called_once_with([fitting_req]) - ex._prepare_disagg_gen_init.assert_called_once_with([]) - ex._check_disagg_ctx_cache_transfer_status.assert_called_once_with(0) - ex._handle_errors.assert_not_called() - @pytest.mark.parametrize( "enable_attention_dp, tp_size, cp_size, gather_name", [ @@ -1437,7 +1409,6 @@ def test_model_parallel_peer_terminal_no_fit_kills_all_ranks( all_rank_status[-1] = (True, True) gather = getattr(ex.dist, gather_name) gather.return_value = all_rank_status - ex._apply_disagg_transfer_admission = Mock(return_value=([], True)) ex._check_disagg_transfer_progress_when_idle = Mock() result, _ = ex._prepare_and_schedule_batch() @@ -1449,8 +1420,8 @@ def test_model_parallel_peer_terminal_no_fit_kills_all_ranks( ex._handle_errors.assert_called_once() assert "one or more requests" in ex._handle_errors.call_args.args[0] - def test_attention_dp_backpressure_without_terminal_peer_does_not_kill(self): - """Admission backpressure stays non-terminal on every rank.""" + def test_attention_dp_without_terminal_peer_does_not_kill(self): + """A fitting INIT request stays non-terminal on every rank.""" fitting_req = _make_active_request(in_init=True) ex = self._make_executor(fill_phase_active=True, fitting_init_requests=[fitting_req]) ex.enable_attention_dp = True @@ -1460,7 +1431,6 @@ def test_attention_dp_backpressure_without_terminal_peer_does_not_kill(self): (True, False), (True, False), ] - ex._apply_disagg_transfer_admission = Mock(return_value=([], True)) ex._check_disagg_transfer_progress_when_idle = Mock() result, _ = ex._prepare_and_schedule_batch() diff --git a/tests/unittest/_torch/disaggregation/test_disagg_coordinator.py b/tests/unittest/_torch/disaggregation/test_disagg_coordinator.py index 35d69425753e..eabf51ec48af 100644 --- a/tests/unittest/_torch/disaggregation/test_disagg_coordinator.py +++ b/tests/unittest/_torch/disaggregation/test_disagg_coordinator.py @@ -63,21 +63,13 @@ def test_real_coordinator_forwards_each_method_to_its_delegate(name: str) -> Non for other in fields(DisaggLoopDelegates): if other.name != name: getattr(delegates, other.name).assert_not_called() - expected = getattr(delegates, name).return_value if name == "admit" else None - assert result is expected - - -def test_noop_coordinator_admits_everything_unchanged() -> None: - """Without a transceiver, scheduler-fitting gen-init requests must pass - through unfiltered and never report a transfer-budget block.""" - fitting = [object(), object()] - assert NoopDisaggCoordinator().admit(fitting) == (fitting, False) + assert result is None def test_noop_coordinator_accepts_every_loop_call() -> None: """Loops call the coordinator unconditionally, so the no-op variant must accept every call the real one does.""" noop = NoopDisaggCoordinator() - for name in _public_methods(DisaggTransferCoordinator) - {"admit"}: + for name in _public_methods(DisaggTransferCoordinator): params = inspect.signature(getattr(noop, name)).parameters getattr(noop, name)(*[Mock() for _ in params]) diff --git a/tests/unittest/_torch/disaggregation/test_disagg_inflight_cancel_gate.py b/tests/unittest/_torch/disaggregation/test_disagg_inflight_cancel_gate.py index 557aee105ed8..3f56cf2a0eee 100644 --- a/tests/unittest/_torch/disaggregation/test_disagg_inflight_cancel_gate.py +++ b/tests/unittest/_torch/disaggregation/test_disagg_inflight_cancel_gate.py @@ -686,7 +686,6 @@ def test_cpp_capability_is_config_scoped( transceiver = BindKvCacheTransceiver(Mock(), dist, kv_cache_manager, Mock(), config) - assert transceiver.consumes_transfer_buffer assert transceiver.supports_inflight_request_cancellation() is expected constructor.assert_called_once() @@ -696,7 +695,6 @@ def test_python_transceiver_capability_defaults_to_unsupported() -> None: transceiver = object.__new__(KvCacheTransceiverV2) - assert not transceiver.consumes_transfer_buffer assert not transceiver.supports_inflight_request_cancellation() assert not transceiver.has_poisoned_transfer_buffer() diff --git a/tests/unittest/_torch/disaggregation/test_disagg_loop_transcript.py b/tests/unittest/_torch/disaggregation/test_disagg_loop_transcript.py index 29cafc578760..55ef5605ba0e 100644 --- a/tests/unittest/_torch/disaggregation/test_disagg_loop_transcript.py +++ b/tests/unittest/_torch/disaggregation/test_disagg_loop_transcript.py @@ -76,11 +76,6 @@ def record(*args, _name=name): setattr(coordinator, name, record) - def admit(fitting): - calls.append(("admit", fitting)) - return fitting, False - - coordinator.admit = admit return coordinator @@ -205,7 +200,6 @@ def _pp_executor(monkeypatch, calls: list, *, rank: int) -> PyExecutor: ("prepare_context_schedulable", []), ("poll_gen_transfers",), ("check_transfer_timeouts",), - ("admit", []), ("receive_gen_init", []), ("poll_progress_when_idle",), ] @@ -249,9 +243,9 @@ def test_executor_loop_overlap_transcript(monkeypatch) -> None: def test_executor_loop_pp_transcript_on_first_rank(monkeypatch) -> None: - """The PP loop admits inside schedule propagation, checks transfer timeouts - only on the retry and executed-batch paths, and flushes responses only from - executed-batch handling, so an idle iteration has none of those.""" + """The PP loop checks transfer timeouts only on the retry and executed-batch + paths, and flushes responses only from executed-batch handling, so an idle + iteration has none of those.""" calls = [] PyExecutor._executor_loop_pp(_pp_executor(monkeypatch, calls, rank=0)) @@ -261,7 +255,6 @@ def test_executor_loop_pp_transcript_on_first_rank(monkeypatch) -> None: ("handle_errors_synced",), ("prepare_context_schedulable", []), ("poll_gen_transfers",), - ("admit", []), ("receive_gen_init", []), ("poll_progress_when_idle",), ("pace_idle",), @@ -271,8 +264,8 @@ def test_executor_loop_pp_transcript_on_first_rank(monkeypatch) -> None: def test_executor_loop_pp_transcript_on_non_first_rank(monkeypatch) -> None: - """A non-first rank does not admit; it reverts KV for candidates its local - scheduler picked but the first rank did not admit.""" + """A non-first rank reverts KV for candidates its local scheduler picked + but the first rank's canonical schedule did not include.""" calls = [] PyExecutor._executor_loop_pp(_pp_executor(monkeypatch, calls, rank=1)) diff --git a/tests/unittest/_torch/executor/test_py_executor.py b/tests/unittest/_torch/executor/test_py_executor.py index 849f6ec2af3d..7139d53a234f 100644 --- a/tests/unittest/_torch/executor/test_py_executor.py +++ b/tests/unittest/_torch/executor/test_py_executor.py @@ -25,7 +25,6 @@ import torch.distributed as torch_dist import torch.multiprocessing as torch_mp -from tensorrt_llm._torch.disaggregation.executor.admission import DisaggTransferAdmissionController from tensorrt_llm._torch.distributed.communicator import ReduceOp from tensorrt_llm._torch.pyexecutor.executor_request_queue import ( SHUTDOWN_REQUEST_ID, @@ -48,7 +47,6 @@ FCFSWaitingQueue, RequestScheduler, ScheduledRequests, - SerializableSchedulerOutput, ) from tensorrt_llm.llmapi.llm_args import EncodeCudaGraphConfig, MTPDecodingConfig from tensorrt_llm.runtime.kv_cache_manager_v2 import OutOfPagesError @@ -869,268 +867,12 @@ def _make_disagg_transfer_request( return req -def _set_disagg_transceiver_capability( - executor: PyExecutor, *, consumes_transfer_buffer: bool -) -> None: - executor.kv_cache_transceiver = Mock() - executor.kv_cache_transceiver.consumes_transfer_buffer = consumes_transfer_buffer - - @pytest.fixture def _clear_disagg_transfer_mode_env(monkeypatch: pytest.MonkeyPatch) -> None: monkeypatch.delenv("TRTLLM_DISAGG_BENCHMARK_GEN_ONLY", raising=False) monkeypatch.delenv("TRTLLM_DISABLE_KV_CACHE_TRANSFER_OVERLAP", raising=False) -@pytest.mark.usefixtures("_clear_disagg_transfer_mode_env") -class TestDisaggTransferAdmissionController: - def test_disabled_preserves_candidates(self): - controller = DisaggTransferAdmissionController( - max_tokens_in_buffer=None, tokens_per_block=32 - ) - candidate = _make_disagg_transfer_request(1, 64) - - result = controller.select(active_requests=[], candidates=[candidate]) - - assert result.admitted_requests == [candidate] - assert result.deferred_request_count == 0 - assert not result.is_blocked_by_active_transfers() - - def test_fcfs_budget_counts_active_transfers(self): - controller = DisaggTransferAdmissionController(max_tokens_in_buffer=64, tokens_per_block=32) - active = _make_disagg_transfer_request(1, 32, in_progress=True) - admitted = _make_disagg_transfer_request(2, 32) - deferred = _make_disagg_transfer_request(3, 32) - - result = controller.select(active_requests=[active], candidates=[admitted, deferred]) - - assert result.admitted_requests == [admitted] - assert result.active_transfer_blocks == 1 - assert result.admitted_transfer_blocks == 1 - assert result.deferred_request_count == 1 - assert result.limited_by_budget - assert not result.is_blocked_by_active_transfers() - - def test_reports_active_transfer_budget_block(self): - controller = DisaggTransferAdmissionController(max_tokens_in_buffer=32, tokens_per_block=32) - active = _make_disagg_transfer_request(1, 32, in_progress=True) - candidate = _make_disagg_transfer_request(2, 32) - - result = controller.select(active_requests=[active], candidates=[candidate]) - - assert result.admitted_requests == [] - assert result.active_transfer_blocks == 1 - assert result.deferred_request_count == 1 - assert result.is_blocked_by_active_transfers() - - def test_admits_oversized_head_when_idle(self): - controller = DisaggTransferAdmissionController(max_tokens_in_buffer=32, tokens_per_block=32) - oversized = _make_disagg_transfer_request(1, 96) - deferred = _make_disagg_transfer_request(2, 32) - - result = controller.select(active_requests=[], candidates=[oversized, deferred]) - - assert result.admitted_requests == [oversized] - assert result.admitted_transfer_blocks == 3 - assert result.deferred_request_count == 1 - assert result.limited_by_budget - assert not result.is_blocked_by_active_transfers() - - def test_uses_global_cp_prompt_length_for_transfer_cost(self): - controller = DisaggTransferAdmissionController( - max_tokens_in_buffer=128, tokens_per_block=32 - ) - request = _make_disagg_transfer_request(1, 32, total_input_len_cp=96) - - result = controller.select(active_requests=[], candidates=[request]) - - assert result.admitted_requests == [request] - assert result.admitted_transfer_blocks == 3 - - def test_apply_reverts_deferred_v2_allocations(self): - executor = object.__new__(PyExecutor) - _set_disagg_transceiver_capability(executor, consumes_transfer_buffer=True) - executor._is_kv_manager_v2 = True - executor._revert_ctx_alloc = Mock() - executor.active_requests = [_make_disagg_transfer_request(1, 32, in_progress=True)] - executor._disagg_transfer_admission_controller = DisaggTransferAdmissionController( - max_tokens_in_buffer=32, tokens_per_block=32 - ) - candidate = _make_disagg_transfer_request(2, 32) - - admitted, wait_for_progress = PyExecutor._apply_disagg_transfer_admission( - executor, [candidate] - ) - - assert admitted == [] - assert wait_for_progress - executor._revert_ctx_alloc.assert_called_once_with([candidate]) - - def test_async_python_v2_pp1_bypasses_transfer_budget(self) -> None: - executor = object.__new__(PyExecutor) - _set_disagg_transceiver_capability(executor, consumes_transfer_buffer=False) - executor.dist = Mock(pp_size=1) - executor._is_kv_manager_v2 = True - executor._revert_ctx_alloc = Mock() - executor.active_requests = [_make_disagg_transfer_request(1, 32, in_progress=True)] - executor._disagg_transfer_admission_controller = DisaggTransferAdmissionController( - max_tokens_in_buffer=32, tokens_per_block=32 - ) - candidates = [ - _make_disagg_transfer_request(2, 32), - _make_disagg_transfer_request(3, 32), - ] - - admitted, wait_for_progress = PyExecutor._apply_disagg_transfer_admission( - executor, candidates - ) - - assert admitted == candidates - assert not wait_for_progress - executor._revert_ctx_alloc.assert_not_called() - - def test_async_python_v1_pp1_retains_transfer_budget(self) -> None: - executor = object.__new__(PyExecutor) - _set_disagg_transceiver_capability(executor, consumes_transfer_buffer=False) - executor.dist = Mock(pp_size=1) - executor._is_kv_manager_v2 = False - executor._revert_ctx_alloc = Mock() - executor.active_requests = [_make_disagg_transfer_request(1, 32, in_progress=True)] - executor._disagg_transfer_admission_controller = DisaggTransferAdmissionController( - max_tokens_in_buffer=32, tokens_per_block=32 - ) - candidate = _make_disagg_transfer_request(2, 32) - - admitted, wait_for_progress = PyExecutor._apply_disagg_transfer_admission( - executor, [candidate] - ) - - assert admitted == [] - assert wait_for_progress - executor._revert_ctx_alloc.assert_not_called() - - def test_disabled_transfer_window_is_inactive(self) -> None: - executor = object.__new__(PyExecutor) - _set_disagg_transceiver_capability(executor, consumes_transfer_buffer=False) - executor.dist = Mock(pp_size=1) - executor._disagg_transfer_admission_controller = DisaggTransferAdmissionController( - max_tokens_in_buffer=0, tokens_per_block=32 - ) - - assert not PyExecutor._disagg_transfer_window_is_active(executor) - - def test_transfer_window_without_transceiver_is_inactive(self) -> None: - executor = object.__new__(PyExecutor) - executor.kv_cache_transceiver = None - executor._disagg_transfer_admission_controller = DisaggTransferAdmissionController( - max_tokens_in_buffer=32, tokens_per_block=32 - ) - - assert not PyExecutor._disagg_transfer_window_is_active(executor) - - def test_active_window_check_requires_initialized_dist(self) -> None: - executor = object.__new__(PyExecutor) - _set_disagg_transceiver_capability(executor, consumes_transfer_buffer=False) - executor._disagg_transfer_admission_controller = DisaggTransferAdmissionController( - max_tokens_in_buffer=32, tokens_per_block=32 - ) - - with pytest.raises(AttributeError): - PyExecutor._disagg_transfer_window_is_active(executor) - - def test_active_window_check_requires_pp_size(self) -> None: - executor = object.__new__(PyExecutor) - _set_disagg_transceiver_capability(executor, consumes_transfer_buffer=False) - executor.dist = types.SimpleNamespace() - executor._disagg_transfer_admission_controller = DisaggTransferAdmissionController( - max_tokens_in_buffer=32, tokens_per_block=32 - ) - - with pytest.raises(AttributeError): - PyExecutor._disagg_transfer_window_is_active(executor) - - def test_apply_missing_controller_preserves_candidates(self): - executor = object.__new__(PyExecutor) - executor.kv_cache_transceiver = Mock() - executor.active_requests = [] - candidate = _make_disagg_transfer_request(1, 32) - - admitted, wait_for_progress = PyExecutor._apply_disagg_transfer_admission( - executor, [candidate] - ) - - assert admitted == [candidate] - assert not wait_for_progress - - def test_apply_non_v2_does_not_revert_deferred_allocations(self): - executor = object.__new__(PyExecutor) - _set_disagg_transceiver_capability(executor, consumes_transfer_buffer=True) - executor._is_kv_manager_v2 = False - executor._revert_ctx_alloc = Mock() - executor.active_requests = [_make_disagg_transfer_request(1, 32, in_progress=True)] - executor._disagg_transfer_admission_controller = DisaggTransferAdmissionController( - max_tokens_in_buffer=32, tokens_per_block=32 - ) - candidate = _make_disagg_transfer_request(2, 32) - - admitted, wait_for_progress = PyExecutor._apply_disagg_transfer_admission( - executor, [candidate] - ) - - assert admitted == [] - assert wait_for_progress - executor._revert_ctx_alloc.assert_not_called() - - def test_sync_python_runtime_retains_transfer_budget( - self, monkeypatch: pytest.MonkeyPatch - ) -> None: - monkeypatch.setenv("TRTLLM_DISABLE_KV_CACHE_TRANSFER_OVERLAP", "1") - executor = object.__new__(PyExecutor) - _set_disagg_transceiver_capability(executor, consumes_transfer_buffer=False) - executor.dist = Mock(pp_size=1) - executor._is_kv_manager_v2 = True - executor._revert_ctx_alloc = Mock() - executor.active_requests = [] - executor._disagg_transfer_admission_controller = DisaggTransferAdmissionController( - max_tokens_in_buffer=32, tokens_per_block=32 - ) - candidates = [ - _make_disagg_transfer_request(2, 32), - _make_disagg_transfer_request(3, 32), - ] - - admitted, wait_for_progress = PyExecutor._apply_disagg_transfer_admission( - executor, candidates - ) - - assert admitted == [candidates[0]] - assert not wait_for_progress - executor._revert_ctx_alloc.assert_called_once_with([candidates[1]]) - - def test_gen_only_no_context_bypasses_transfer_budget(self, monkeypatch): - monkeypatch.setenv("TRTLLM_DISAGG_BENCHMARK_GEN_ONLY", "1") - executor = object.__new__(PyExecutor) - executor.kv_cache_transceiver = Mock() - executor._is_kv_manager_v2 = True - executor._revert_ctx_alloc = Mock() - executor.active_requests = [_make_disagg_transfer_request(1, 32, in_progress=True)] - executor._disagg_transfer_admission_controller = DisaggTransferAdmissionController( - max_tokens_in_buffer=32, tokens_per_block=32 - ) - candidates = [ - _make_disagg_transfer_request(2, 32), - _make_disagg_transfer_request(3, 32), - ] - - admitted, wait_for_progress = PyExecutor._apply_disagg_transfer_admission( - executor, candidates - ) - - assert admitted == candidates - assert not wait_for_progress - executor._revert_ctx_alloc.assert_not_called() - - @pytest.mark.usefixtures("_clear_disagg_transfer_mode_env") class TestDisaggTransferIdleProgress: def test_gen_transfer_status_polls_active_transfers(self): @@ -1405,108 +1147,6 @@ def test_pp_ring_drained_only_when_no_microbatch_is_outstanding( assert PyExecutor._pp_ring_is_drained(executor) is expected -@pytest.mark.usefixtures("_clear_disagg_transfer_mode_env") -class TestDisaggTransferAdmissionPP: - def test_pp_schedule_applies_gate_before_serializing(self) -> None: - executor = object.__new__(PyExecutor) - _set_disagg_transceiver_capability(executor, consumes_transfer_buffer=True) - executor._is_kv_manager_v2 = False - executor.dist = Mock( - rank=0, is_first_pp_rank=True, is_last_pp_rank=True, tp_size=1, cp_size=1 - ) - executor.enable_attention_dp = False - executor.active_requests = [_make_disagg_transfer_request(1, 32, in_progress=True)] - executor._disagg_transfer_admission_controller = DisaggTransferAdmissionController( - max_tokens_in_buffer=32, tokens_per_block=32 - ) - scheduled_batch = ScheduledRequests() - candidate = _make_disagg_transfer_request(2, 32) - executor._schedule = Mock(return_value=(scheduled_batch, [candidate], 0)) - - scheduled, fitting, num_fitting, wait_for_progress = PyExecutor._pp_schedule_and_propagate( - executor, microbatch_id=0 - ) - - assert scheduled is scheduled_batch - assert fitting == [] - assert num_fitting == 0 - assert wait_for_progress - - def test_pp_schedule_async_python_retains_transfer_window(self) -> None: - executor = object.__new__(PyExecutor) - _set_disagg_transceiver_capability(executor, consumes_transfer_buffer=False) - executor.dist = Mock( - rank=0, - is_first_pp_rank=True, - is_last_pp_rank=False, - tp_size=1, - cp_size=1, - pp_size=2, - next_pp_rank=1, - ) - executor._is_kv_manager_v2 = True - executor.enable_attention_dp = False - executor.send_schedule_handles = [None] - executor.wait_on_pp_send_handles = Mock() - executor._revert_ctx_alloc = Mock() - executor.active_requests = [_make_disagg_transfer_request(1, 32, in_progress=True)] - executor._disagg_transfer_admission_controller = DisaggTransferAdmissionController( - max_tokens_in_buffer=32, tokens_per_block=32 - ) - scheduled_batch = ScheduledRequests() - candidate = _make_disagg_transfer_request(2, 32) - executor._schedule = Mock(return_value=(scheduled_batch, [candidate], 0)) - - scheduled, fitting, num_fitting, wait_for_progress = PyExecutor._pp_schedule_and_propagate( - executor, microbatch_id=0 - ) - - assert scheduled is scheduled_batch - assert fitting == [] - assert num_fitting == 0 - assert wait_for_progress - executor._revert_ctx_alloc.assert_called_once_with([candidate]) - executor.wait_on_pp_send_handles.assert_called_once() - wait_args = executor.wait_on_pp_send_handles.call_args.args - assert wait_args[0] is executor.send_schedule_handles - assert wait_args[1] == 0 - executor.dist.isend_object.assert_called_once() - - def test_pp_schedule_restores_propagated_gate_decision(self): - executor = object.__new__(PyExecutor) - executor.dist = Mock( - rank=1, - is_first_pp_rank=False, - is_last_pp_rank=True, - prev_pp_rank=0, - tp_size=1, - cp_size=1, - ) - executor.enable_attention_dp = False - executor.active_requests = [ - _make_disagg_transfer_request(1, 32, in_progress=True), - _make_disagg_transfer_request(2, 32), - ] - serializable_schedule = SerializableSchedulerOutput( - encoder_requests=[], - context_requests_chunking=[], - context_requests_last_chunk=[], - generation_requests=[], - paused_requests=[], - fitting_disagg_gen_init_requests=[2], - num_fitting_requests=0, - wait_for_disagg_gen_transfer_progress=True, - ) - executor.dist.recv_object = Mock(return_value=serializable_schedule) - - _, fitting, _, wait_for_progress = PyExecutor._pp_schedule_and_propagate( - executor, microbatch_id=0 - ) - - assert [req.py_request_id for req in fitting] == [2] - assert wait_for_progress - - def test_nonzero_pp_rank_prepares_snapshot_points_before_local_schedule( monkeypatch, ): @@ -1529,7 +1169,7 @@ class StopLocalSchedule(RuntimeError): executor.kv_cache_transceiver = None executor._pad_attention_dp_dummy_request = Mock() scheduled_batch = Mock() - executor._pp_schedule_and_propagate = Mock(return_value=(scheduled_batch, [], 0, False)) + executor._pp_schedule_and_propagate = Mock(return_value=(scheduled_batch, [], 0)) executor._pp_retry_until_can_schedule = Mock() request = Mock() executor.active_requests = [request] @@ -1600,9 +1240,7 @@ class StopAfterReconciliation(RuntimeError): local_only = _make_disagg_transfer_request(2, 32) executor.active_requests = [canonical, local_only] scheduled_batch = ScheduledRequests() - executor._pp_schedule_and_propagate = Mock( - return_value=(scheduled_batch, [canonical], 0, False) - ) + executor._pp_schedule_and_propagate = Mock(return_value=(scheduled_batch, [canonical], 0)) executor.scheduler.schedule_request.return_value = types.SimpleNamespace( fitting_disagg_gen_init_requests=[local_canonical, local_only] ) diff --git a/tests/unittest/_torch/executor/test_scheduler_serializable_output.py b/tests/unittest/_torch/executor/test_scheduler_serializable_output.py index 7901e030c187..56173d3b20c6 100644 --- a/tests/unittest/_torch/executor/test_scheduler_serializable_output.py +++ b/tests/unittest/_torch/executor/test_scheduler_serializable_output.py @@ -42,7 +42,6 @@ def test_serializable_scheduler_output_round_trip(): scheduled_requests, fitting_disagg_gen_init_requests, num_fitting_requests, - wait_for_disagg_gen_transfer_progress=True, ) # Serialize and deserialize the serializable scheduler output @@ -57,7 +56,6 @@ def test_serializable_scheduler_output_round_trip(): # Verify the restored scheduler result is correct assert restored_num_fitting == num_fitting_requests - assert restored_output.wait_for_disagg_gen_transfer_progress assert _request_ids(restored_schedule.encoder_requests) == _request_ids( scheduled_requests.encoder_requests )