Add timeout abort kits for normal / eagle. (#18815)
This commit is contained in:
@@ -0,0 +1,144 @@
|
|||||||
|
import time
|
||||||
|
from concurrent.futures import ThreadPoolExecutor, as_completed
|
||||||
|
|
||||||
|
import requests
|
||||||
|
|
||||||
|
# Safety timeout for all HTTP requests to prevent CI from hanging forever.
|
||||||
|
_REQUEST_TIMEOUT = 60
|
||||||
|
|
||||||
|
|
||||||
|
class AbortAllMixin:
|
||||||
|
"""Test /abort_request with abort_all=True.
|
||||||
|
|
||||||
|
Server needs sufficient --max-running-requests.
|
||||||
|
"""
|
||||||
|
|
||||||
|
abort_all_num_requests: int = 32
|
||||||
|
abort_all_max_new_tokens: int = 16000
|
||||||
|
abort_all_sleep: float = 2
|
||||||
|
|
||||||
|
def test_abort_all(self):
|
||||||
|
num_requests = self.abort_all_num_requests
|
||||||
|
|
||||||
|
def run_decode():
|
||||||
|
response = requests.post(
|
||||||
|
self.base_url + "/generate",
|
||||||
|
json={
|
||||||
|
"text": "The capital of France is",
|
||||||
|
"sampling_params": {
|
||||||
|
"temperature": 0,
|
||||||
|
"max_new_tokens": self.abort_all_max_new_tokens,
|
||||||
|
"ignore_eos": True,
|
||||||
|
},
|
||||||
|
},
|
||||||
|
timeout=_REQUEST_TIMEOUT,
|
||||||
|
)
|
||||||
|
return response.json()
|
||||||
|
|
||||||
|
with ThreadPoolExecutor(num_requests) as executor:
|
||||||
|
futures = [executor.submit(run_decode) for _ in range(num_requests)]
|
||||||
|
|
||||||
|
time.sleep(self.abort_all_sleep)
|
||||||
|
|
||||||
|
requests.post(
|
||||||
|
self.base_url + "/abort_request",
|
||||||
|
json={"abort_all": True},
|
||||||
|
timeout=10,
|
||||||
|
).raise_for_status()
|
||||||
|
|
||||||
|
for future in as_completed(futures):
|
||||||
|
result = future.result()
|
||||||
|
self.assertEqual(result["meta_info"]["finish_reason"]["type"], "abort")
|
||||||
|
|
||||||
|
self.assertIsNone(self.process.poll())
|
||||||
|
|
||||||
|
|
||||||
|
class WaitingTimeoutMixin:
|
||||||
|
"""Test waiting queue timeout.
|
||||||
|
|
||||||
|
Server needs SGLANG_REQ_WAITING_TIMEOUT and --max-running-requests=1.
|
||||||
|
"""
|
||||||
|
|
||||||
|
waiting_timeout_num_requests: int = 2
|
||||||
|
waiting_timeout_max_new_tokens: int = 512
|
||||||
|
|
||||||
|
def test_waiting_timeout(self):
|
||||||
|
num_requests = self.waiting_timeout_num_requests
|
||||||
|
|
||||||
|
def run_decode():
|
||||||
|
response = requests.post(
|
||||||
|
self.base_url + "/generate",
|
||||||
|
json={
|
||||||
|
"text": "Today is ",
|
||||||
|
"sampling_params": {
|
||||||
|
"temperature": 0,
|
||||||
|
"max_new_tokens": self.waiting_timeout_max_new_tokens,
|
||||||
|
"ignore_eos": True,
|
||||||
|
},
|
||||||
|
},
|
||||||
|
timeout=_REQUEST_TIMEOUT,
|
||||||
|
)
|
||||||
|
return response.json()
|
||||||
|
|
||||||
|
with ThreadPoolExecutor(num_requests) as executor:
|
||||||
|
futures = [executor.submit(run_decode) for _ in range(num_requests)]
|
||||||
|
|
||||||
|
error_count = 0
|
||||||
|
for future in as_completed(futures):
|
||||||
|
result = future.result()
|
||||||
|
if result.get("object") == "error":
|
||||||
|
error_count += 1
|
||||||
|
self.assertEqual(result["code"], 503)
|
||||||
|
|
||||||
|
self.assertEqual(error_count, 1)
|
||||||
|
self.assertIsNone(self.process.poll())
|
||||||
|
|
||||||
|
|
||||||
|
class RunningTimeoutTwoWaveMixin:
|
||||||
|
"""Test running timeout with a two-wave pattern.
|
||||||
|
|
||||||
|
Sends two waves with different forward_entry_time so that timeouts are
|
||||||
|
triggered in separate batches. Regression test for
|
||||||
|
https://github.com/sgl-project/sglang/pull/18760
|
||||||
|
|
||||||
|
Server needs SGLANG_REQ_RUNNING_TIMEOUT and sufficient --max-running-requests
|
||||||
|
to hold both waves.
|
||||||
|
"""
|
||||||
|
|
||||||
|
running_timeout_num_wave1: int = 8
|
||||||
|
running_timeout_num_wave2: int = 8
|
||||||
|
running_timeout_sleep: float = 3
|
||||||
|
running_timeout_max_new_tokens: int = 1024
|
||||||
|
|
||||||
|
def test_running_timeout_no_crash(self):
|
||||||
|
num_wave1 = self.running_timeout_num_wave1
|
||||||
|
num_wave2 = self.running_timeout_num_wave2
|
||||||
|
|
||||||
|
def run_decode():
|
||||||
|
response = requests.post(
|
||||||
|
self.base_url + "/generate",
|
||||||
|
json={
|
||||||
|
"text": "Write a long story about a magical kingdom.",
|
||||||
|
"sampling_params": {
|
||||||
|
"temperature": 0,
|
||||||
|
"max_new_tokens": self.running_timeout_max_new_tokens,
|
||||||
|
"ignore_eos": True,
|
||||||
|
},
|
||||||
|
},
|
||||||
|
timeout=_REQUEST_TIMEOUT,
|
||||||
|
)
|
||||||
|
return response.json()
|
||||||
|
|
||||||
|
with ThreadPoolExecutor(num_wave1 + num_wave2) as executor:
|
||||||
|
futures1 = [executor.submit(run_decode) for _ in range(num_wave1)]
|
||||||
|
|
||||||
|
time.sleep(self.running_timeout_sleep)
|
||||||
|
|
||||||
|
futures2 = [executor.submit(run_decode) for _ in range(num_wave2)]
|
||||||
|
|
||||||
|
for future in as_completed(futures1 + futures2):
|
||||||
|
result = future.result()
|
||||||
|
if result.get("object") == "error":
|
||||||
|
self.assertEqual(result["code"], 503)
|
||||||
|
|
||||||
|
self.assertIsNone(self.process.poll())
|
||||||
@@ -8,6 +8,7 @@ import requests
|
|||||||
from sglang.srt.environ import envs
|
from sglang.srt.environ import envs
|
||||||
from sglang.srt.utils import kill_process_tree
|
from sglang.srt.utils import kill_process_tree
|
||||||
from sglang.test.ci.ci_register import register_amd_ci, register_cuda_ci
|
from sglang.test.ci.ci_register import register_amd_ci, register_cuda_ci
|
||||||
|
from sglang.test.kits.abort_timeout_kit import AbortAllMixin, WaitingTimeoutMixin
|
||||||
from sglang.test.test_utils import (
|
from sglang.test.test_utils import (
|
||||||
DEFAULT_MODEL_NAME_FOR_TEST,
|
DEFAULT_MODEL_NAME_FOR_TEST,
|
||||||
DEFAULT_TIMEOUT_FOR_SERVER_LAUNCH,
|
DEFAULT_TIMEOUT_FOR_SERVER_LAUNCH,
|
||||||
@@ -108,7 +109,7 @@ class TestAbortWithApiKey(CustomTestCase):
|
|||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
class TestAbortAll(CustomTestCase):
|
class TestAbortAll(AbortAllMixin, CustomTestCase):
|
||||||
@classmethod
|
@classmethod
|
||||||
def setUpClass(cls):
|
def setUpClass(cls):
|
||||||
cls.model = DEFAULT_MODEL_NAME_FOR_TEST
|
cls.model = DEFAULT_MODEL_NAME_FOR_TEST
|
||||||
@@ -124,40 +125,6 @@ class TestAbortAll(CustomTestCase):
|
|||||||
def tearDownClass(cls):
|
def tearDownClass(cls):
|
||||||
kill_process_tree(cls.process.pid)
|
kill_process_tree(cls.process.pid)
|
||||||
|
|
||||||
def _run_decode(self):
|
|
||||||
response = requests.post(
|
|
||||||
self.base_url + "/generate",
|
|
||||||
json={
|
|
||||||
"text": "The capital of France is",
|
|
||||||
"sampling_params": {
|
|
||||||
"temperature": 0,
|
|
||||||
"max_new_tokens": 16000,
|
|
||||||
"ignore_eos": True,
|
|
||||||
},
|
|
||||||
},
|
|
||||||
)
|
|
||||||
return response.json()
|
|
||||||
|
|
||||||
def test_abort_all(self):
|
|
||||||
num_requests = 32
|
|
||||||
with ThreadPoolExecutor(num_requests) as executor:
|
|
||||||
futures = [executor.submit(self._run_decode) for _ in range(num_requests)]
|
|
||||||
|
|
||||||
# ensure the decode has been started
|
|
||||||
time.sleep(2)
|
|
||||||
|
|
||||||
requests.post(
|
|
||||||
self.base_url + "/abort_request",
|
|
||||||
json={
|
|
||||||
"abort_all": True,
|
|
||||||
},
|
|
||||||
)
|
|
||||||
|
|
||||||
for future in as_completed(futures):
|
|
||||||
self.assertEqual(
|
|
||||||
future.result()["meta_info"]["finish_reason"]["type"], "abort"
|
|
||||||
)
|
|
||||||
|
|
||||||
|
|
||||||
class TestAbortAllWithRetraction(CustomTestCase):
|
class TestAbortAllWithRetraction(CustomTestCase):
|
||||||
@classmethod
|
@classmethod
|
||||||
@@ -236,7 +203,7 @@ class TestAbortAllWithRetraction(CustomTestCase):
|
|||||||
print("Finished test_abort_all_with_retraction")
|
print("Finished test_abort_all_with_retraction")
|
||||||
|
|
||||||
|
|
||||||
class TestAbortWithWaitingTimeout(CustomTestCase):
|
class TestAbortWithWaitingTimeout(WaitingTimeoutMixin, CustomTestCase):
|
||||||
@classmethod
|
@classmethod
|
||||||
def setUpClass(cls):
|
def setUpClass(cls):
|
||||||
cls.model = DEFAULT_MODEL_NAME_FOR_TEST
|
cls.model = DEFAULT_MODEL_NAME_FOR_TEST
|
||||||
@@ -255,33 +222,6 @@ class TestAbortWithWaitingTimeout(CustomTestCase):
|
|||||||
def tearDownClass(cls):
|
def tearDownClass(cls):
|
||||||
kill_process_tree(cls.process.pid)
|
kill_process_tree(cls.process.pid)
|
||||||
|
|
||||||
def _run_decode(self):
|
|
||||||
response = requests.post(
|
|
||||||
self.base_url + "/generate",
|
|
||||||
json={
|
|
||||||
"text": "Today is ",
|
|
||||||
"sampling_params": {
|
|
||||||
"temperature": 0,
|
|
||||||
"max_new_tokens": 512,
|
|
||||||
"ignore_eos": True,
|
|
||||||
},
|
|
||||||
},
|
|
||||||
)
|
|
||||||
return response.json()
|
|
||||||
|
|
||||||
def test_waiting_timeout(self):
|
|
||||||
num_requests = 2
|
|
||||||
with ThreadPoolExecutor(num_requests) as executor:
|
|
||||||
futures = [executor.submit(self._run_decode) for _ in range(num_requests)]
|
|
||||||
|
|
||||||
error_count = 0
|
|
||||||
for future in as_completed(futures):
|
|
||||||
result = future.result()
|
|
||||||
if result.get("object") == "error":
|
|
||||||
error_count += 1
|
|
||||||
self.assertEqual(result["code"], 503)
|
|
||||||
self.assertEqual(error_count, 1)
|
|
||||||
|
|
||||||
|
|
||||||
class TestAbortWithRunningTimeout(CustomTestCase):
|
class TestAbortWithRunningTimeout(CustomTestCase):
|
||||||
@classmethod
|
@classmethod
|
||||||
|
|||||||
@@ -13,6 +13,11 @@ import requests
|
|||||||
from sglang.srt.environ import envs
|
from sglang.srt.environ import envs
|
||||||
from sglang.test.ci.ci_register import register_cuda_ci
|
from sglang.test.ci.ci_register import register_cuda_ci
|
||||||
from sglang.test.few_shot_gsm8k import run_eval as run_gsm8k_eval
|
from sglang.test.few_shot_gsm8k import run_eval as run_gsm8k_eval
|
||||||
|
from sglang.test.kits.abort_timeout_kit import (
|
||||||
|
AbortAllMixin,
|
||||||
|
RunningTimeoutTwoWaveMixin,
|
||||||
|
WaitingTimeoutMixin,
|
||||||
|
)
|
||||||
from sglang.test.kits.radix_cache_server_kit import run_radix_attention_test
|
from sglang.test.kits.radix_cache_server_kit import run_radix_attention_test
|
||||||
from sglang.test.server_fixtures.eagle_fixture import EagleServerBase
|
from sglang.test.server_fixtures.eagle_fixture import EagleServerBase
|
||||||
from sglang.test.test_utils import DEFAULT_TARGET_MODEL_EAGLE, run_logprob_check
|
from sglang.test.test_utils import DEFAULT_TARGET_MODEL_EAGLE, run_logprob_check
|
||||||
@@ -347,5 +352,29 @@ class TestEAGLEServerPageSizeTopkFA3(TestEAGLEServerBasic):
|
|||||||
]
|
]
|
||||||
|
|
||||||
|
|
||||||
|
class TestEAGLEAbortAll(AbortAllMixin, EagleServerBase):
|
||||||
|
abort_all_max_new_tokens = 4000
|
||||||
|
extra_args = ["--max-running-requests=8"]
|
||||||
|
|
||||||
|
|
||||||
|
class TestEAGLEWaitingTimeout(WaitingTimeoutMixin, EagleServerBase):
|
||||||
|
extra_args = ["--max-running-requests=1"]
|
||||||
|
|
||||||
|
@classmethod
|
||||||
|
def setUpClass(cls):
|
||||||
|
with envs.SGLANG_REQ_WAITING_TIMEOUT.override(0.001):
|
||||||
|
super().setUpClass()
|
||||||
|
|
||||||
|
|
||||||
|
class TestEAGLERunningTimeout(RunningTimeoutTwoWaveMixin, EagleServerBase):
|
||||||
|
# Regression test for https://github.com/sgl-project/sglang/pull/18760
|
||||||
|
extra_args = ["--max-running-requests=16"]
|
||||||
|
|
||||||
|
@classmethod
|
||||||
|
def setUpClass(cls):
|
||||||
|
with envs.SGLANG_REQ_RUNNING_TIMEOUT.override(3):
|
||||||
|
super().setUpClass()
|
||||||
|
|
||||||
|
|
||||||
if __name__ == "__main__":
|
if __name__ == "__main__":
|
||||||
unittest.main()
|
unittest.main()
|
||||||
|
|||||||
Reference in New Issue
Block a user