项目文件夹

文件
wehub-resource-sync 94057c3d3e
PR Test (NPU) / check-changes (push) Has been cancelled
PR Test (NPU) / pr-gate (push) Has been cancelled
PR Test (NPU) / set-image-config (push) Has been cancelled
PR Test (NPU) / stage-b-test-1-npu-a2 (0) (push) Has been cancelled
PR Test (NPU) / stage-b-test-1-npu-a2 (1) (push) Has been cancelled
PR Test (NPU) / stage-b-test-2-npu-a2 (0) (push) Has been cancelled
PR Test (NPU) / stage-b-test-2-npu-a2 (1) (push) Has been cancelled
PR Test (NPU) / stage-b-test-4-npu-a3 (push) Has been cancelled
PR Test (NPU) / stage-b-test-16-npu-a3 (push) Has been cancelled
PR Test (NPU) / multimodal-gen-test-1-npu-a3 (push) Has been cancelled
PR Test (NPU) / multimodal-gen-test-2-npu-a3 (push) Has been cancelled
PR Test (Arm64) / pr-gate (push) Has been cancelled
PR Test (Arm64) / check-changes (push) Has been cancelled
PR Test (Arm64) / build-test (push) Has been cancelled
PR Test (sgl-router) / gate (push) Has been cancelled
PR Test (sgl-router) / tier-1 — lint (push) Has been cancelled
PR Test (sgl-router) / tier-2 — build + test (push) Has been cancelled
PR Test (sgl-router) / tier-3 — docker (placeholder) (push) Has been cancelled
PR Test (sgl-router) / tier-3 — k8s integration (push) Has been cancelled
PR Test (sgl-router) / tier-3 — e2e (push) Has been cancelled
PR Test (sgl-router) / finish (push) Has been cancelled
PR Test (NPU) / single-node-poc (map[name:qwen3_6_27b_w8a8_1p_in64k_out1k_50ms runner:linux-aarch64-a3-2 test_case:test/registered/ascend/performance/qwen3_6_27b/test_npu_qwen3_6_27b_w8a8_1p_in64k_out1k_50ms.py test_type:perf]) (push) Has been cancelled
PR Test (NPU) / pr-test-npu-finish (push) Has been cancelled
PR Test (Xeon) / pr-gate (push) Has been cancelled
PR Test (Xeon) / check-changes (push) Has been cancelled
PR Test (Xeon) / build-test (, xeon-gnr, base-b-test-cpu) (push) Has been cancelled
PR Test (XPU) / check-changes (push) Has been cancelled
PR Test (XPU) / pr-gate (push) Has been cancelled
PR Test (XPU) / stage-a-test-1-gpu-xpu (push) Has been cancelled
PR Test (XPU) / wait-for-stage-a (push) Has been cancelled
PR Test (XPU) / stage-b-test-1-gpu-xpu (push) Has been cancelled
PR Test (XPU) / finish (push) Has been cancelled
CI Model Inventory / build-inventory (push) Has been cancelled
Lint / lint (push) Has been cancelled
PR Benchmark (SMG Components) / Benchmark Compilation Check (push) Has been cancelled
PR Benchmark (SMG Components) / Benchmark - Manual Policy (push) Has been cancelled
PR Benchmark (SMG Components) / Benchmark - Request Processing (push) Has been cancelled
PR Benchmark (SMG Components) / Benchmark Summary (push) Has been cancelled
PR Test (SMG) / build-wheel (push) Has been cancelled
Release SGLang Model Gateway to PyPI / build on windows (x86_64 - auto) (push) Has been cancelled
Release SGLang Model Gateway to PyPI / build on macos (x86_64 - auto) (push) Has been cancelled
PR Test (SMG) / python-unit-tests (push) Has been cancelled
PR Test (SMG) / unit-tests (push) Has been cancelled
PR Test (SMG) / benchmarks (push) Has been cancelled
PR Test (SMG) / chat-completions (push) Has been cancelled
PR Test (SMG) / chat-completions-4gpu (push) Has been cancelled
PR Test (SMG) / e2e (push) Has been cancelled
PR Test (SMG) / docker-build-test (push) Has been cancelled
PR Test (SMG) / k8s-integration (push) Has been cancelled
PR Test (SMG) / finish (push) Has been cancelled
PR Test (SMG) / summarize-benchmarks (push) Has been cancelled
Release SGLang Model Gateway Docker Image / publish (push) Has been cancelled
Release SGLang Model Gateway to PyPI / build on macos (aarch64 - auto) (push) Has been cancelled
Release SGLang Model Gateway to PyPI / build on linux (aarch64 - auto) (push) Has been cancelled
Release SGLang Model Gateway to PyPI / build on linux (x86_64 - auto) (push) Has been cancelled
Release SGLang Model Gateway to PyPI / build on linux (aarch64 - musllinux_1_1) (push) Has been cancelled
Release SGLang Model Gateway to PyPI / build on linux (x86_64 - musllinux_1_1) (push) Has been cancelled
Release SGLang Model Gateway to PyPI / Build SDist (push) Has been cancelled
Release SGLang Model Gateway to PyPI / Upload to PyPI (push) Has been cancelled
Release SGLang Kernels / build-cu129-matrix (aarch64, 12.9, 3.10, arm-kernel-build-node) (push) Has been cancelled
Release SGLang Kernels / build-cu129-matrix (x86_64, 12.9, 3.10, x64-kernel-build-node) (push) Has been cancelled
Release SGLang Kernels / release-cu129 (push) Has been cancelled
Release SGLang Kernels / build-cu130-matrix (aarch64, 13.0, 3.10, arm-kernel-build-node) (push) Has been cancelled
Release SGLang Kernels / build-cu130-matrix (x86_64, 13.0, 3.10, x64-kernel-build-node) (push) Has been cancelled
Release SGLang Kernels / release-cu130 (push) Has been cancelled
Release SGLang Kernels / build-rocm-matrix (3.10, 700) (push) Has been cancelled
Release SGLang Kernels / build-rocm-matrix (3.10, 720) (push) Has been cancelled
Release SGLang Kernels / release-rocm700 (push) Has been cancelled
Release SGLang Kernels / release-rocm720 (push) Has been cancelled
Release SGLang Kernels / build-musa43 (43, 3.10) (push) Has been cancelled
Release SGLang Kernels / release-musa43 (push) Has been cancelled
chore: import upstream snapshot with attribution
2026-07-13 12:38:16 +08:00

216 行
8.2 KiB
Python

import unittest
from types import SimpleNamespace
from unittest.mock import MagicMock, patch
from sglang.srt.disaggregation.base import KVPoll
from sglang.srt.disaggregation.decode import (
DecodePreallocQueue,
DecodeTransferQueue,
HiCacheRestoreResult,
)
from sglang.srt.disaggregation.utils import DisaggregationMode
from sglang.srt.managers.schedule_batch import FINISH_ABORT
from sglang.srt.managers.scheduler import Scheduler
from sglang.test.ci.ci_register import register_cpu_ci
from sglang.test.test_utils import CustomTestCase
register_cpu_ci(est_time=5, suite="base-a-test-cpu")
class FakeReceiver:
def __init__(self):
self.clear_called = False
def clear(self):
self.clear_called = True
def failure_exception(self):
return None
class TestDecodeQueueCleanup(CustomTestCase):
def test_prealloc_abort_clears_receiver_before_removing_request(self):
receiver = FakeReceiver()
req = SimpleNamespace(
rid="abort-prealloc",
finished_reason=FINISH_ABORT("aborted"),
return_logprob=False,
)
decode_req = SimpleNamespace(req=req, kv_receiver=receiver)
queue = DecodePreallocQueue.__new__(DecodePreallocQueue)
queue.queue = [decode_req]
queue.pending_reqs = []
queue.retracted_queue = []
queue._resolve_pending_reqs = MagicMock()
queue._update_handshake_waiters = MagicMock()
queue._uses_swa_tail_prealloc = MagicMock(return_value=False)
queue._allocatable_token_budgets = MagicMock(return_value=0)
queue._hicache_pending_restore_tokens = MagicMock(return_value=0)
scheduler = MagicMock()
scheduler.running_batch.reqs = []
scheduler.enable_priority_scheduling = False
scheduler.enable_hisparse = False
scheduler.output_streamer = MagicMock()
queue.scheduler = scheduler
preallocated, failed = queue.pop_preallocated()
self.assertEqual(preallocated, [])
self.assertEqual(failed, [decode_req])
self.assertEqual(queue.queue, [])
self.assertTrue(receiver.clear_called)
self.assertIsNone(decode_req.kv_receiver)
scheduler.output_streamer.stream_output.assert_called_once_with(
[req], req.return_logprob
)
def test_prealloc_abort_also_drops_from_pending_reqs(self):
# Same DecodeRequest lives in both queue and pending_reqs (add() slow
# path). Aborting must drop it from both, and compare by identity since
# DecodeRequest's dataclass __eq__ would compare the tensor receiver.
class BadEqReceiver(FakeReceiver):
def __eq__(self, other):
raise TypeError("use identity comparison, not value equality")
__hash__ = object.__hash__
receiver = BadEqReceiver()
req = SimpleNamespace(
rid="abort-shared",
finished_reason=FINISH_ABORT("aborted"),
return_logprob=False,
)
decode_req = SimpleNamespace(req=req, kv_receiver=receiver)
queue = DecodePreallocQueue.__new__(DecodePreallocQueue)
queue.queue = [decode_req]
queue.pending_reqs = [decode_req] # same object, dual ownership
queue.retracted_queue = []
queue._resolve_pending_reqs = MagicMock()
queue._update_handshake_waiters = MagicMock()
queue._uses_swa_tail_prealloc = MagicMock(return_value=False)
queue._allocatable_token_budgets = MagicMock(return_value=0)
queue._hicache_pending_restore_tokens = MagicMock(return_value=0)
scheduler = MagicMock()
scheduler.running_batch.reqs = []
scheduler.enable_priority_scheduling = False
scheduler.enable_hisparse = False
scheduler.output_streamer = MagicMock()
queue.scheduler = scheduler
# Must not raise on the receiver __eq__ above.
preallocated, failed = queue.pop_preallocated()
self.assertEqual(preallocated, [])
self.assertEqual(failed, [decode_req])
self.assertEqual(queue.queue, [])
self.assertTrue(all(r is not decode_req for r in queue.pending_reqs))
self.assertIsNone(decode_req.kv_receiver)
def test_ensure_prefill_info_tolerates_cleared_receiver(self):
# A req whose kv_receiver was already cleared must not crash on .abort().
queue = DecodePreallocQueue.__new__(DecodePreallocQueue)
queue._max_ensure_retries = 1
queue._ensure_retry_interval = 0
queue._ensure_retry_count = {"127.0.0.1:11500": 0}
queue._ensure_last_attempt_time = {}
queue.kv_manager = MagicMock()
queue.kv_manager.try_ensure_parallel_info.return_value = False
cleared_req = SimpleNamespace(
req=SimpleNamespace(rid="cleared"), kv_receiver=None
)
addr_to_reqs = {"127.0.0.1:11500": [cleared_req]}
ready, remaining = queue._ensure_prefill_info(addr_to_reqs)
self.assertEqual(ready, {})
self.assertEqual(remaining, [])
@patch("sglang.srt.disaggregation.decode.release_kv_cache")
@patch("sglang.srt.disaggregation.decode.prepare_abort")
@patch("sglang.srt.disaggregation.decode.poll_and_all_reduce")
def test_transfer_failure_clears_receiver_before_removing_request(
self, mock_poll, mock_prepare_abort, mock_release_kv_cache
):
receiver = FakeReceiver()
req = SimpleNamespace(
rid="failed-transfer",
bootstrap_room=7,
return_logprob=False,
)
decode_req = SimpleNamespace(
req=req,
kv_receiver=receiver,
metadata_buffer_index=3,
hicache_restore_status=HiCacheRestoreResult.READY,
)
queue = DecodeTransferQueue.__new__(DecodeTransferQueue)
queue.queue = [decode_req]
queue.enable_staging = False
queue.gloo_group = MagicMock()
queue.req_to_metadata_buffer_idx_allocator = MagicMock()
queue.tp_rank = 0
queue.tree_cache = MagicMock()
queue.metadata_buffers = SimpleNamespace(bootstrap_room=[None] * 4)
queue.spec_algorithm = MagicMock()
queue.spec_algorithm.is_none.return_value = True
queue._clean_hicache_prefetch_resources = MagicMock()
scheduler = MagicMock()
scheduler.enable_decode_hicache = False
scheduler.enable_hisparse = False
scheduler.output_streamer = MagicMock()
scheduler.metrics_reporter.enable_metrics = False
queue.scheduler = scheduler
mock_poll.return_value = [KVPoll.Failed]
transferred = queue.pop_transferred()
self.assertEqual(transferred, [])
self.assertEqual(queue.queue, [])
self.assertTrue(receiver.clear_called)
self.assertIsNone(decode_req.kv_receiver)
queue.req_to_metadata_buffer_idx_allocator.free.assert_called_once_with(3)
scheduler.output_streamer.stream_output.assert_called_once_with(
[req], req.return_logprob
)
mock_prepare_abort.assert_called_once()
mock_release_kv_cache.assert_called_once_with(
req, queue.tree_cache, is_insert=False
)
def test_retracted_decode_requests_keep_scheduler_non_idle(self):
scheduler = Scheduler.__new__(Scheduler)
scheduler.running_batch = MagicMock()
scheduler.running_batch.is_empty.return_value = True
scheduler.chunked_req = None
scheduler.dllm_manager = MagicMock()
scheduler.dllm_manager.any_staging_reqs.return_value = False
scheduler.last_batch = None
scheduler.cur_batch_for_debug = None
scheduler.enable_overlap = False
scheduler.ps = SimpleNamespace(pp_size=1)
scheduler.running_mbs = []
scheduler.waiting_queue = []
scheduler.grammar_manager = SimpleNamespace(grammar_queue=[])
scheduler.disaggregation_mode = DisaggregationMode.DECODE
scheduler.disagg_decode_prealloc_queue = SimpleNamespace(
queue=[], retracted_queue=[object()]
)
scheduler.disagg_decode_transfer_queue = SimpleNamespace(queue=[])
scheduler.decode_offload_manager = None
scheduler.enable_hisparse = False
scheduler.enable_hierarchical_cache = False
self.assertFalse(scheduler.is_fully_idle())
if __name__ == "__main__":
unittest.main()