diff --git a/.github/workflows/pr-test-rust.yml b/.github/workflows/pr-test-rust.yml index ddd10339e..c741f1a6c 100644 --- a/.github/workflows/pr-test-rust.yml +++ b/.github/workflows/pr-test-rust.yml @@ -196,14 +196,6 @@ jobs: python3 -m pip install pytest pytest-cov pytest-xdist pytest -q tests --cov=sglang_router --cov-config=.coveragerc --cov-report=term-missing --cov-fail-under=80 - - name: Run Python integration tests - run: | - cd sgl-model-gateway - source "$HOME/.cargo/env" - # Integration tests use FastAPI/uvicorn for mock workers - python3 -m pip install fastapi uvicorn orjson - pytest -q py_test/integration_mock - - name: Run Python E2E tests run: | bash scripts/killall_sglang.sh "nuk_gpus" diff --git a/sgl-model-gateway/py_test/integration_mock/__init__.py b/sgl-model-gateway/py_test/integration_mock/__init__.py deleted file mode 100644 index 1e342eca0..000000000 --- a/sgl-model-gateway/py_test/integration_mock/__init__.py +++ /dev/null @@ -1 +0,0 @@ -"""Integration test package for the router.""" diff --git a/sgl-model-gateway/py_test/integration_mock/conftest.py b/sgl-model-gateway/py_test/integration_mock/conftest.py deleted file mode 100644 index 25bd7c2bc..000000000 --- a/sgl-model-gateway/py_test/integration_mock/conftest.py +++ /dev/null @@ -1,128 +0,0 @@ -import shutil -import subprocess -import time -from pathlib import Path -from typing import Iterable, List, Optional, Tuple - -import pytest -import requests - -from ..fixtures.generate_test_certs import generate_all_certificates -from ..fixtures.ports import find_free_port -from ..fixtures.router_manager import RouterManager - - -def pytest_configure(config): - config.addinivalue_line("markers", "integration: mark as router integration test") - - -@pytest.fixture -def router_manager() -> Iterable[RouterManager]: - mgr = RouterManager() - try: - yield mgr - finally: - mgr.stop_all() - - -def _spawn_mock_worker(args: List[str]) -> Tuple[subprocess.Popen, str, str]: - repo_root = Path(__file__).resolve().parents[2] - script = repo_root / "py_test" / "fixtures" / "mock_worker.py" - port = find_free_port() - worker_id = f"worker-{port}" - base_cmd = [ - "python3", - str(script), - "--port", - str(port), - "--worker-id", - worker_id, - ] - cmd = base_cmd + args - proc = subprocess.Popen(cmd) - url = f"http://127.0.0.1:{port}" - _wait_health(url) - return proc, url, worker_id - - -def _wait_health(url: str, timeout: float = 10.0): - start = time.time() - with requests.Session() as s: - while time.time() - start < timeout: - try: - r = s.get(f"{url}/health", timeout=1) - if r.status_code == 200: - return - except requests.RequestException: - pass - time.sleep(0.1) - raise TimeoutError(f"Mock worker at {url} did not become healthy") - - -@pytest.fixture -def mock_worker(): - """Start a single healthy mock worker; yields (process, url, worker_id).""" - proc, url, worker_id = _spawn_mock_worker([]) - try: - yield proc, url, worker_id - finally: - if proc.poll() is None: - proc.terminate() - try: - proc.wait(timeout=3) - except subprocess.TimeoutExpired: - proc.kill() - - -@pytest.fixture -def mock_workers(): - """Factory to start N workers with custom args. - - Usage: - procs, urls, ids = mock_workers(n=3, args=["--latency-ms", "5"]) # same args for all - ... - """ - - procs: List[subprocess.Popen] = [] - - def _start(n: int, args: Optional[List[str]] = None): - args = args or [] - new_procs: List[subprocess.Popen] = [] - urls: List[str] = [] - ids: List[str] = [] - for _ in range(n): - p, url, wid = _spawn_mock_worker(args) - procs.append(p) - new_procs.append(p) - urls.append(url) - ids.append(wid) - return new_procs, urls, ids - - try: - yield _start - finally: - for p in procs: - if p.poll() is None: - p.terminate() - try: - p.wait(timeout=3) - except subprocess.TimeoutExpired: - p.kill() - - -@pytest.fixture(scope="session") -def test_certificates(): - """Generate test certificates for mTLS tests, clean up after session.""" - # Get the test_certs directory path - fixtures_dir = Path(__file__).parent.parent / "fixtures" - certs_dir = fixtures_dir / "test_certs" - - # Generate certificates - generate_all_certificates(certs_dir) - - # Yield the path to the certificates directory - yield certs_dir - - # Cleanup: remove the generated certificates - if certs_dir.exists(): - shutil.rmtree(certs_dir) diff --git a/sgl-model-gateway/py_test/integration_mock/load_balancing/__init__.py b/sgl-model-gateway/py_test/integration_mock/load_balancing/__init__.py deleted file mode 100644 index 77b8c2460..000000000 --- a/sgl-model-gateway/py_test/integration_mock/load_balancing/__init__.py +++ /dev/null @@ -1 +0,0 @@ -"""Load balancing integration tests.""" diff --git a/sgl-model-gateway/py_test/integration_mock/load_balancing/test_cache_aware.py b/sgl-model-gateway/py_test/integration_mock/load_balancing/test_cache_aware.py deleted file mode 100644 index acbbd3682..000000000 --- a/sgl-model-gateway/py_test/integration_mock/load_balancing/test_cache_aware.py +++ /dev/null @@ -1,73 +0,0 @@ -import collections -import concurrent.futures -import uuid - -import pytest -import requests - - -@pytest.mark.integration -def test_cache_aware_affinity(mock_workers, router_manager): - # Two workers; same prompt should stick to one due to cache tree - _, urls, ids = mock_workers(n=2) - rh = router_manager.start_router(worker_urls=urls, policy="cache_aware") - - counts = collections.Counter() - with requests.Session() as s: - for i in range(12): - r = s.post( - f"{rh.url}/v1/completions", - json={ - "model": "test-model", - "prompt": "repeated prompt for cache", - "max_tokens": 1, - "stream": False, - }, - ) - assert r.status_code == 200 - wid = r.headers.get("X-Worker-Id") or r.json().get("worker_id") - counts[wid] += 1 - - # Expect strong skew toward one worker (tree match); majority > 80% - top = max(counts.values()) - assert top >= 10, counts - - -@pytest.mark.integration -def test_cache_aware_diverse_prompts_balances(mock_workers, router_manager): - # Add latency so concurrent requests overlap and influence load-based selection - _, urls, ids = mock_workers(n=3, args=["--latency-ms", "30"]) - rh = router_manager.start_router( - worker_urls=urls, - policy="cache_aware", - extra={ - "cache_threshold": 0.99, - "balance_abs_threshold": 0, - "balance_rel_threshold": 1.0, - }, - ) - - counts = collections.Counter() - - def call(i): - # Use diverse, unrelated prompts to avoid prefix matches entirely - prompt = str(uuid.uuid4()) - r = requests.post( - f"{rh.url}/v1/completions", - json={ - "model": "test-model", - "prompt": prompt, - "max_tokens": 1, - "stream": False, - }, - timeout=5, - ) - assert r.status_code == 200 - return r.headers.get("X-Worker-Id") or r.json().get("worker_id") - - with concurrent.futures.ThreadPoolExecutor(max_workers=16) as ex: - for wid in ex.map(call, range(40)): - counts[wid] += 1 - - # Expect participation of at least two workers - assert sum(1 for v in counts.values() if v > 0) >= 2, counts diff --git a/sgl-model-gateway/py_test/integration_mock/load_balancing/test_manual.py b/sgl-model-gateway/py_test/integration_mock/load_balancing/test_manual.py deleted file mode 100644 index bbc1e07de..000000000 --- a/sgl-model-gateway/py_test/integration_mock/load_balancing/test_manual.py +++ /dev/null @@ -1,52 +0,0 @@ -import collections - -import pytest -import requests - -ROUTING_KEY_HEADER = "X-SMG-Routing-Key" - - -@pytest.mark.integration -def test_manual_routing_with_header(mock_workers, router_manager): - """With X-SMG-Routing-Key header: sticky routing + distribution across workers.""" - _, urls, _ = mock_workers(n=2) - rh = router_manager.start_router(worker_urls=urls, policy="manual") - - # Send requests: 5 keys × 4 requests each - results = collections.defaultdict(set) - with requests.Session() as s: - for key_id in range(5): - for _ in range(4): - worker = send_completion(s, rh.url, f"user-{key_id}") - results[f"user-{key_id}"].add(worker) - - # Verify sticky: each key should route to exactly one worker - for key, workers in results.items(): - assert len(workers) == 1, f"Key {key} routed to multiple workers: {workers}" - - # Verify distribution: different keys should use multiple workers - all_workers = {list(w)[0] for w in results.values()} - assert len(all_workers) > 1, f"Should distribute across workers: {results}" - - -@pytest.mark.integration -def test_manual_routing_without_header(mock_workers, router_manager): - """Without X-SMG-Routing-Key header: random fallback distribution.""" - _, urls, _ = mock_workers(n=2) - rh = router_manager.start_router(worker_urls=urls, policy="manual") - - with requests.Session() as s: - counts = collections.Counter(send_completion(s, rh.url) for _ in range(20)) - - assert len(counts) > 1, f"Random fallback should distribute: {counts}" - - -def send_completion(session, base_url, routing_key=None): - headers = {ROUTING_KEY_HEADER: routing_key} if routing_key is not None else {} - r = session.post( - f"{base_url}/v1/completions", - json={"model": "test", "prompt": "hi", "max_tokens": 1, "stream": False}, - headers=headers, - ) - assert r.status_code == 200 - return r.headers.get("X-Worker-Id") or r.json().get("worker_id") diff --git a/sgl-model-gateway/py_test/integration_mock/load_balancing/test_power_of_two.py b/sgl-model-gateway/py_test/integration_mock/load_balancing/test_power_of_two.py deleted file mode 100644 index 0a8d9eab1..000000000 --- a/sgl-model-gateway/py_test/integration_mock/load_balancing/test_power_of_two.py +++ /dev/null @@ -1,99 +0,0 @@ -import collections -import concurrent.futures -import time - -import pytest -import requests - - -@pytest.mark.integration -def test_power_of_two_prefers_less_loaded(mock_workers, router_manager): - # Start two workers: one slow (higher inflight), one fast - # Router monitors /get_load and Power-of-Two uses cached loads to choose - # Start one slow and one fast worker using the fixture factory - procs_slow, urls_slow, ids_slow = mock_workers(n=1, args=["--latency-ms", "200"]) - procs_fast, urls_fast, ids_fast = mock_workers(n=1, args=["--latency-ms", "0"]) - procs = procs_slow + procs_fast - urls = urls_slow + urls_fast - ids = ids_slow + ids_fast - slow_id = ids_slow[0] - slow_url = urls_slow[0] - - rh = router_manager.start_router( - worker_urls=urls, - policy="power_of_two", - extra={"worker_startup_check_interval": 1}, - ) - - # Prime: fire a burst to create measurable load on slow worker, then wait for monitor tick - - def _prime_call(i): - try: - requests.post( - f"{rh.url}/v1/completions", - json={ - "model": "test-model", - "prompt": f"warm-{i}", - "max_tokens": 1, - "stream": False, - }, - timeout=5, - ) - except Exception: - pass - - with concurrent.futures.ThreadPoolExecutor(max_workers=32) as ex: - list(ex.map(_prime_call, range(128))) - time.sleep(2) - - # Apply direct background load on the slow worker to amplify load diff - def _direct_load(i): - try: - requests.post( - f"{slow_url}/v1/completions", - json={ - "model": "test-model", - "prompt": f"bg-{i}", - "max_tokens": 1, - "stream": False, - }, - timeout=5, - ) - except Exception: - pass - - # Start background load in a non-blocking way to keep slow worker busy - background_executor = concurrent.futures.ThreadPoolExecutor(max_workers=8) - background_futures = [] - for i in range(32): - future = background_executor.submit(_direct_load, i) - background_futures.append(future) - - # Wait longer for the load monitor to update (at least 2 monitor intervals) - time.sleep(3) - - def call(i): - r = requests.post( - f"{rh.url}/v1/completions", - json={ - "model": "test-model", - "prompt": f"p{i}", - "max_tokens": 1, - "stream": False, - }, - timeout=5, - ) - assert r.status_code == 200 - return r.headers.get("X-Worker-Id") or r.json().get("worker_id") - - counts = collections.Counter() - with concurrent.futures.ThreadPoolExecutor(max_workers=32) as ex: - for wid in ex.map(call, range(200)): - counts[wid] += 1 - - # Clean up background executor - background_executor.shutdown(wait=False) - - # Expect the slow worker (higher latency/inflight) to receive fewer requests - fast_worker_id = [i for i in ids if i != slow_id][0] - assert counts[slow_id] < counts[fast_worker_id], counts diff --git a/sgl-model-gateway/py_test/integration_mock/load_balancing/test_random.py b/sgl-model-gateway/py_test/integration_mock/load_balancing/test_random.py deleted file mode 100644 index 4662dbce0..000000000 --- a/sgl-model-gateway/py_test/integration_mock/load_balancing/test_random.py +++ /dev/null @@ -1,32 +0,0 @@ -import collections - -import pytest -import requests - - -@pytest.mark.integration -def test_random_distribution(mock_workers, router_manager): - procs, urls, ids = mock_workers(n=4) - rh = router_manager.start_router(worker_urls=urls, policy="random") - - counts = collections.Counter() - N = 200 - with requests.Session() as s: - for i in range(N): - r = s.post( - f"{rh.url}/v1/completions", - json={ - "model": "test-model", - "prompt": f"p{i}", - "max_tokens": 1, - "stream": False, - }, - ) - assert r.status_code == 200 - wid = r.headers.get("X-Worker-Id") or r.json().get("worker_id") - counts[wid] += 1 - - # simple statistical tolerance: each worker should be within ±50% of mean - mean = N / len(ids) - for wid in ids: - assert 0.5 * mean <= counts[wid] <= 1.5 * mean, counts diff --git a/sgl-model-gateway/py_test/integration_mock/load_balancing/test_round_robin.py b/sgl-model-gateway/py_test/integration_mock/load_balancing/test_round_robin.py deleted file mode 100644 index 13f149635..000000000 --- a/sgl-model-gateway/py_test/integration_mock/load_balancing/test_round_robin.py +++ /dev/null @@ -1,33 +0,0 @@ -import collections - -import pytest -import requests - - -@pytest.mark.integration -def test_round_robin_distribution(mock_workers, router_manager): - procs, urls, ids = mock_workers(n=3) - - rh = router_manager.start_router(worker_urls=urls, policy="round_robin") - - counts = collections.Counter() - with requests.Session() as s: - for i in range(30): - r = s.post( - f"{rh.url}/v1/completions", - json={ - "model": "test-model", - "prompt": f"hello {i}", - "max_tokens": 1, - "stream": False, - }, - ) - assert r.status_code == 200 - wid = r.headers.get("X-Worker-Id") or r.json().get("worker_id") - assert wid in ids - counts[wid] += 1 - - # Expect near-even distribution across 3 workers - # 30 requests -> ideally 10 each; allow small tolerance ±3 - for wid in ids: - assert 7 <= counts[wid] <= 13, counts diff --git a/sgl-model-gateway/py_test/integration_mock/test_api_auth.py b/sgl-model-gateway/py_test/integration_mock/test_api_auth.py deleted file mode 100644 index b8ba5c670..000000000 --- a/sgl-model-gateway/py_test/integration_mock/test_api_auth.py +++ /dev/null @@ -1,38 +0,0 @@ -import pytest -import requests - - -@pytest.mark.integration -def test_router_api_key_enforcement(router_manager, mock_workers): - # Start backend requiring API key; router should forward Authorization header transparently - _, urls, _ = mock_workers( - n=1, args=["--require-api-key", "--api-key", "correct_api_key"] - ) - rh = router_manager.start_router( - worker_urls=urls, - policy="round_robin", - extra={}, - ) - - # No auth -> 401 - r = requests.post( - f"{rh.url}/v1/completions", - json={"model": "test-model", "prompt": "x", "max_tokens": 1, "stream": False}, - ) - assert r.status_code == 401 - - # Invalid auth -> 401 - r = requests.post( - f"{rh.url}/v1/completions", - json={"model": "test-model", "prompt": "x", "max_tokens": 1, "stream": False}, - headers={"Authorization": "Bearer wrong"}, - ) - assert r.status_code == 401 - - # Correct auth -> 200 - r = requests.post( - f"{rh.url}/v1/completions", - json={"model": "test-model", "prompt": "x", "max_tokens": 1, "stream": False}, - headers={"Authorization": "Bearer correct_api_key"}, - ) - assert r.status_code == 200 diff --git a/sgl-model-gateway/py_test/integration_mock/test_circuit_breaker.py b/sgl-model-gateway/py_test/integration_mock/test_circuit_breaker.py deleted file mode 100644 index 16c619924..000000000 --- a/sgl-model-gateway/py_test/integration_mock/test_circuit_breaker.py +++ /dev/null @@ -1,228 +0,0 @@ -import time - -import pytest -import requests - - -@pytest.mark.integration -def test_circuit_breaker_opens_and_recovers(router_manager, mock_workers): - # A single worker that fails first 3 requests, then succeeds - _, [wurl], _ = mock_workers(n=1, args=["--fail-first-n", "3"]) # fails first 3 - rh = router_manager.start_router( - worker_urls=[wurl], - policy="round_robin", - extra={ - "cb_failure_threshold": 3, - "cb_success_threshold": 2, - "cb_timeout_duration_secs": 3, - "cb_window_duration_secs": 10, - "disable_retries": True, - }, - ) - - def post_once(): - return requests.post( - f"{rh.url}/v1/completions", - json={ - "model": "test-model", - "prompt": "trigger", - "max_tokens": 1, - "stream": False, - }, - timeout=3, - ) - - # should see 500 when worker actually starts, before that should see 503 - saw_500 = False - for _ in range(8): - r = post_once() - if r.status_code == 500: - # Worker starts, continue to circuit breaker test - saw_500 = True - break - assert ( - r.status_code == 503 - ), "Should only see 503 when waiting for worker to start" - assert saw_500, "Worker didn't start after 8 requests" - - saw_503 = False - for _ in range(4): - r = post_once() - if r.status_code == 503: - saw_503 = True - break - assert saw_503, "circuit breaker did not open to return 503" - - time.sleep(4) - r1 = post_once() - r2 = post_once() - assert r1.status_code == 200 and r2.status_code == 200 - - -@pytest.mark.integration -def test_circuit_breaker_half_open_failure_reopens(router_manager, mock_workers): - _, [wurl], _ = mock_workers(n=1, args=["--status-code", "500"]) # always fail - rh = router_manager.start_router( - worker_urls=[wurl], - policy="round_robin", - extra={ - "cb_failure_threshold": 2, - "cb_success_threshold": 2, - "cb_timeout_duration_secs": 2, - "cb_window_duration_secs": 5, - "disable_retries": True, - }, - ) - - def post_once(): - return requests.post( - f"{rh.url}/v1/completions", - json={ - "model": "test-model", - "prompt": "x", - "max_tokens": 1, - "stream": False, - }, - timeout=3, - ) - - # should see 500 when worker actually starts, before that should see 503 - saw_500 = False - for _ in range(8): - r = post_once() - if r.status_code == 500: - # Worker starts, continue to circuit breaker test - saw_500 = True - break - assert ( - r.status_code == 503 - ), "Should only see 503 when waiting for worker to start" - assert saw_500, "Worker didn't start after 8 requests" - - opened = False - for _ in range(8): - r = post_once() - if r.status_code == 503: - opened = True - break - assert opened, "circuit breaker did not open" - - time.sleep(3) - r = post_once() - assert r.status_code == 500 - r2 = post_once() - assert r2.status_code == 503 - - -@pytest.mark.integration -def test_circuit_breaker_disable_flag(router_manager, mock_workers): - _, [wurl], _ = mock_workers(n=1, args=["--status-code", "500"]) # always fail - rh = router_manager.start_router( - worker_urls=[wurl], - policy="round_robin", - extra={ - "disable_circuit_breaker": True, - "disable_retries": True, - }, - ) - - saw_500 = False - for _ in range(8): - r = requests.post( - f"{rh.url}/v1/completions", - json={ - "model": "test-model", - "prompt": "x", - "max_tokens": 1, - "stream": False, - }, - timeout=3, - ) - if r.status_code == 500: - # Worker starts, continue to check - saw_500 = True - break - assert ( - r.status_code == 503 - ), "Should only see 503 when waiting for worker to start" - - assert saw_500 - - -@pytest.mark.integration -def test_circuit_breaker_per_worker_isolation(router_manager, mock_workers): - _, [fail_url], _ = mock_workers(n=1, args=["--status-code", "500"]) # always fail - _, [ok_url], _ = mock_workers(n=1) - rh = router_manager.start_router( - worker_urls=[fail_url, ok_url], - policy="round_robin", - extra={ - "cb_failure_threshold": 2, - "cb_success_threshold": 1, - "cb_timeout_duration_secs": 2, - "cb_window_duration_secs": 10, - "disable_retries": True, - }, - ) - - def post_once(): - return requests.post( - f"{rh.url}/v1/completions", - json={ - "model": "test-model", - "prompt": "y", - "max_tokens": 1, - "stream": False, - }, - timeout=3, - ) - - failures = 0 - successes_after_open = 0 - opened = False - for _ in range(30): - r = post_once() - if not opened: - if r.status_code == 500: - failures += 1 - if failures >= 2: - _ = post_once() - _ = post_once() - opened = True - else: - if r.status_code == 200: - successes_after_open += 1 - else: - assert False, f"Unexpected non-200 after CB open: {r.status_code}" - assert opened and successes_after_open >= 5 - - -@pytest.mark.integration -def test_circuit_breaker_with_retries(router_manager, mock_workers): - _, [fail_url], _ = mock_workers(n=1, args=["--status-code", "500"]) # always fail - _, [ok_url], _ = mock_workers(n=1) - rh = router_manager.start_router( - worker_urls=[fail_url, ok_url], - policy="round_robin", - extra={ - "retry_max_retries": 3, - "retry_initial_backoff_ms": 10, - "retry_max_backoff_ms": 50, - "cb_failure_threshold": 2, - "cb_success_threshold": 1, - "cb_timeout_duration_secs": 2, - "cb_window_duration_secs": 10, - }, - ) - - r = requests.post( - f"{rh.url}/v1/completions", - json={ - "model": "test-model", - "prompt": "z", - "max_tokens": 1, - "stream": False, - }, - timeout=5, - ) - assert r.status_code == 200 diff --git a/sgl-model-gateway/py_test/integration_mock/test_fault_tolerance.py b/sgl-model-gateway/py_test/integration_mock/test_fault_tolerance.py deleted file mode 100644 index 6cadf1fae..000000000 --- a/sgl-model-gateway/py_test/integration_mock/test_fault_tolerance.py +++ /dev/null @@ -1,32 +0,0 @@ -import pytest -import requests - - -@pytest.mark.integration -def test_worker_crash_reroute_with_retries(router_manager, mock_workers): - # Start one healthy and one that will crash on first request - _, [ok_url], _ = mock_workers(n=1) - _, [crash_url], _ = mock_workers(n=1, args=["--crash-on-request"]) - rh = router_manager.start_router( - worker_urls=[crash_url, ok_url], - policy="round_robin", - extra={ - "retry_max_retries": 3, - "retry_initial_backoff_ms": 10, - "retry_max_backoff_ms": 50, - }, - ) - - # A single request should succeed via retry to the healthy worker - r = requests.post( - f"{rh.url}/v1/completions", - json={ - "model": "test-model", - "prompt": "crash", - "max_tokens": 1, - "stream": False, - }, - timeout=5, - ) - assert r.status_code == 200 - # mock_workers fixture handles cleanup diff --git a/sgl-model-gateway/py_test/integration_mock/test_header_forwarding.py b/sgl-model-gateway/py_test/integration_mock/test_header_forwarding.py deleted file mode 100644 index 4af28a4a8..000000000 --- a/sgl-model-gateway/py_test/integration_mock/test_header_forwarding.py +++ /dev/null @@ -1,36 +0,0 @@ -import pytest -import requests - - -@pytest.mark.integration -def test_header_forwarding_whitelist(mock_workers, router_manager): - _, urls, _ = mock_workers(n=1) - rh = router_manager.start_router(worker_urls=urls) - - with requests.Session() as s: - r = s.post( - f"{rh.url}/v1/completions", - json={"model": "test", "prompt": "hi", "max_tokens": 1, "stream": False}, - headers={ - "Authorization": "Bearer test-token", - "X-SMG-Routing-Key": "routing-123", - "X-Request-Id": "req-456", - "X-Correlation-Id": "corr-789", - "traceparent": "00-trace-span-01", - "tracestate": "vendor=value", - "X-Custom-Header": "should-not-forward", - "Cookie": "session=abc", - }, - ) - assert r.status_code == 200 - h = r.json().get("received_headers", {}) - - assert h.get("authorization") == "Bearer test-token" - assert h.get("x-request-id") == "req-456" - assert h.get("x-correlation-id") == "corr-789" - assert h.get("traceparent") == "00-trace-span-01" - assert h.get("tracestate") == "vendor=value" - assert h.get("x-smg-routing-key") == "routing-123" - - assert "x-custom-header" not in h - assert "cookie" not in h diff --git a/sgl-model-gateway/py_test/integration_mock/test_mtls.py b/sgl-model-gateway/py_test/integration_mock/test_mtls.py deleted file mode 100644 index caa0db924..000000000 --- a/sgl-model-gateway/py_test/integration_mock/test_mtls.py +++ /dev/null @@ -1,332 +0,0 @@ -""" -Integration tests for mTLS (mutual TLS) authentication between router and workers. - -Tests verify that: -1. Router can successfully connect to TLS-enabled workers with proper certificates -2. Router fails to connect to mTLS-required workers without client certificates -3. Router with CA certs can connect to TLS-only workers (server auth only) -""" - -import subprocess -import time -from pathlib import Path -from typing import Tuple - -import pytest -import requests - -from ..fixtures.ports import find_free_port - - -def get_test_certs_dir() -> Path: - """Get the path to the test certificates directory.""" - return Path(__file__).parent.parent / "fixtures" / "test_certs" - - -def _spawn_tls_worker( - port: int, - worker_id: str, - ssl_certfile: str, - ssl_keyfile: str, - ssl_ca_certs: str = None, -) -> Tuple[subprocess.Popen, str]: - """Spawn a mock worker with TLS/mTLS enabled.""" - repo_root = Path(__file__).resolve().parents[2] - script = repo_root / "py_test" / "fixtures" / "mock_worker.py" - - cmd = [ - "python3", - str(script), - "--port", - str(port), - "--worker-id", - worker_id, - "--ssl-certfile", - ssl_certfile, - "--ssl-keyfile", - ssl_keyfile, - ] - - if ssl_ca_certs: - cmd.extend(["--ssl-ca-certs", ssl_ca_certs]) - - # Use DEVNULL for stdout to avoid blocking, but keep stderr for debugging - proc = subprocess.Popen( - cmd, stdout=subprocess.DEVNULL, stderr=subprocess.PIPE, text=True - ) - url = f"https://127.0.0.1:{port}" - - # Give worker a moment to start or fail - import time - - time.sleep(3) # Increased delay to ensure TLS server is fully initialized - - # Check if process died immediately - if proc.poll() is not None: - _, stderr = proc.communicate() - raise RuntimeError(f"Worker failed to start.\nStderr: {stderr}") - - # Wait for worker to be ready (with retries for SSL startup) - # For mTLS workers (with ssl_ca_certs), provide client cert for health check - certs_dir = get_test_certs_dir() - client_cert = certs_dir / "client-cert.pem" if ssl_ca_certs else None - client_key = certs_dir / "client-key.pem" if ssl_ca_certs else None - - try: - _wait_tls_health(url, certs_dir / "ca-cert.pem", client_cert, client_key) - except TimeoutError: - # If health check times out, capture stderr for debugging - if proc.poll() is not None: - _, stderr = proc.communicate() - raise RuntimeError(f"Worker died during health check.\nStderr: {stderr}") - raise - return proc, url - - -def _wait_tls_health( - url: str, - ca_cert_path: Path = None, - client_cert_path: Path = None, - client_key_path: Path = None, - timeout: float = 10.0, -): - """Wait for TLS-enabled worker to become healthy. - - Args: - url: HTTPS URL of the worker - ca_cert_path: Path to CA certificate for verifying server cert - client_cert_path: Path to client certificate for mTLS - client_key_path: Path to client private key for mTLS - timeout: Maximum time to wait in seconds - """ - start = time.time() - last_error = None - with requests.Session() as s: - while time.time() - start < timeout: - try: - # Verify server cert with CA if provided, otherwise skip verification - verify = str(ca_cert_path) if ca_cert_path else False - - # Provide client cert for mTLS if specified - cert = None - if client_cert_path and client_key_path: - cert = (str(client_cert_path), str(client_key_path)) - - r = s.get(f"{url}/health", timeout=1, verify=verify, cert=cert) - if r.status_code == 200: - return - except requests.RequestException as e: - # Save last error for debugging - last_error = e - time.sleep(0.2) - raise TimeoutError( - f"TLS worker at {url} did not become healthy. Last error: {last_error}" - ) - - -@pytest.mark.integration -def test_mtls_successful_communication(router_manager, test_certificates): - """Test that router can successfully communicate with mTLS-enabled worker.""" - certs_dir = test_certificates - - # Start worker with mTLS (requires client certificate) - port = find_free_port() - worker_id = f"tls-worker-{port}" - worker_proc, worker_url = _spawn_tls_worker( - port=port, - worker_id=worker_id, - ssl_certfile=str(certs_dir / "server-cert.pem"), - ssl_keyfile=str(certs_dir / "server-key.pem"), - ssl_ca_certs=str(certs_dir / "ca-cert.pem"), # Require client cert - ) - - try: - # Start router with mTLS configuration - rh = router_manager.start_router( - worker_urls=[worker_url], - policy="round_robin", - extra={ - "client_cert_path": str(certs_dir / "client-cert.pem"), - "client_key_path": str(certs_dir / "client-key.pem"), - "ca_cert_paths": [str(certs_dir / "ca-cert.pem")], - }, - ) - - # Make request through router - should succeed - r = requests.post( - f"{rh.url}/v1/completions", - json={ - "model": "test-model", - "prompt": "hello", - "max_tokens": 1, - "stream": False, - }, - timeout=5, - ) - - assert r.status_code == 200, f"Request failed: {r.status_code} {r.text}" - data = r.json() - assert "choices" in data - assert data.get("worker_id") == worker_id - - finally: - if worker_proc.poll() is None: - worker_proc.terminate() - try: - worker_proc.wait(timeout=3) - except subprocess.TimeoutExpired: - worker_proc.kill() - - -@pytest.mark.integration -def test_mtls_failure_without_client_cert(router_manager, test_certificates): - """Test that router fails to connect to mTLS worker without client certificates.""" - certs_dir = test_certificates - - # Start worker with mTLS (requires client certificate) - port = find_free_port() - worker_id = f"tls-worker-{port}" - worker_proc, worker_url = _spawn_tls_worker( - port=port, - worker_id=worker_id, - ssl_certfile=str(certs_dir / "server-cert.pem"), - ssl_keyfile=str(certs_dir / "server-key.pem"), - ssl_ca_certs=str(certs_dir / "ca-cert.pem"), # Require client cert - ) - - try: - # Start router WITHOUT client certificates (but with CA to verify server) - rh = router_manager.start_router( - worker_urls=[worker_url], - policy="round_robin", - extra={ - "ca_cert_paths": [str(certs_dir / "ca-cert.pem")], - # Note: no client_cert_path or client_key_path - }, - ) - - # Make request through router - should fail because worker requires client cert - r = requests.post( - f"{rh.url}/v1/completions", - json={ - "model": "test-model", - "prompt": "hello", - "max_tokens": 1, - "stream": False, - }, - timeout=5, - ) - - # Router should return 503 (service unavailable) or 500 because it can't connect to worker - assert r.status_code in [500, 503], f"Expected 500/503 but got {r.status_code}" - - finally: - if worker_proc.poll() is None: - worker_proc.terminate() - try: - worker_proc.wait(timeout=3) - except subprocess.TimeoutExpired: - worker_proc.kill() - - -@pytest.mark.integration -def test_tls_server_auth_only(router_manager, test_certificates): - """Test router can connect to TLS worker that doesn't require client certificates.""" - certs_dir = test_certificates - - # Start worker with TLS but WITHOUT requiring client certificates - port = find_free_port() - worker_id = f"tls-worker-{port}" - worker_proc, worker_url = _spawn_tls_worker( - port=port, - worker_id=worker_id, - ssl_certfile=str(certs_dir / "server-cert.pem"), - ssl_keyfile=str(certs_dir / "server-key.pem"), - ssl_ca_certs=None, # Don't require client cert - ) - - try: - # Start router with only CA cert (to verify server), no client cert - rh = router_manager.start_router( - worker_urls=[worker_url], - policy="round_robin", - extra={ - "ca_cert_paths": [str(certs_dir / "ca-cert.pem")], - # Note: no client_cert_path or client_key_path needed - }, - ) - - # Make request through router - should succeed with server-only TLS - r = requests.post( - f"{rh.url}/v1/completions", - json={ - "model": "test-model", - "prompt": "hello", - "max_tokens": 1, - "stream": False, - }, - timeout=5, - ) - - assert r.status_code == 200, f"Request failed: {r.status_code} {r.text}" - data = r.json() - assert "choices" in data - assert data.get("worker_id") == worker_id - - finally: - if worker_proc.poll() is None: - worker_proc.terminate() - try: - worker_proc.wait(timeout=3) - except subprocess.TimeoutExpired: - worker_proc.kill() - - -@pytest.mark.integration -def test_tls_failure_without_ca_cert(router_manager, test_certificates): - """Test that router fails to connect to TLS worker without CA certificate.""" - certs_dir = test_certificates - - # Start worker with TLS - port = find_free_port() - worker_id = f"tls-worker-{port}" - worker_proc, worker_url = _spawn_tls_worker( - port=port, - worker_id=worker_id, - ssl_certfile=str(certs_dir / "server-cert.pem"), - ssl_keyfile=str(certs_dir / "server-key.pem"), - ssl_ca_certs=None, - ) - - try: - # Start router WITHOUT CA certificate (can't verify server cert) - rh = router_manager.start_router( - worker_urls=[worker_url], - policy="round_robin", - extra={ - # Note: no ca_cert_paths - router won't trust self-signed cert - }, - ) - - # Make request through router - should fail because router can't verify server cert - r = requests.post( - f"{rh.url}/v1/completions", - json={ - "model": "test-model", - "prompt": "hello", - "max_tokens": 1, - "stream": False, - }, - timeout=5, - ) - - # Router should return 503 (service unavailable) or 500 because it can't verify worker cert - assert r.status_code in [500, 503], f"Expected 500/503 but got {r.status_code}" - - finally: - if worker_proc.poll() is None: - worker_proc.terminate() - try: - worker_proc.wait(timeout=3) - except subprocess.TimeoutExpired: - worker_proc.kill() diff --git a/sgl-model-gateway/py_test/integration_mock/test_payload_size.py b/sgl-model-gateway/py_test/integration_mock/test_payload_size.py deleted file mode 100644 index b3583ab28..000000000 --- a/sgl-model-gateway/py_test/integration_mock/test_payload_size.py +++ /dev/null @@ -1,33 +0,0 @@ -import pytest -import requests - - -@pytest.mark.integration -def test_payload_size_limit(router_manager, mock_workers): - # Start one backend and a router with a 1MB payload limit - _, urls, _ = mock_workers(n=1) - rh = router_manager.start_router( - worker_urls=urls, - policy="round_robin", - extra={"max_payload_size": 1 * 1024 * 1024}, # 1MB - ) - - # Payload just under 1MB should succeed - payload_small = { - "model": "test-model", - "prompt": "x" * int(0.5 * 1024 * 1024), # ~0.5MB - "max_tokens": 1, - "stream": False, - } - r = requests.post(f"{rh.url}/v1/completions", json=payload_small) - assert r.status_code == 200 - - # Payload over 1MB should fail with 413 - payload_large = { - "model": "test-model", - "prompt": "x" * int(1.2 * 1024 * 1024), # ~1.2MB - "max_tokens": 1, - "stream": False, - } - r = requests.post(f"{rh.url}/v1/completions", json=payload_large) - assert r.status_code == 413 diff --git a/sgl-model-gateway/py_test/integration_mock/test_pd_routing.py b/sgl-model-gateway/py_test/integration_mock/test_pd_routing.py deleted file mode 100644 index 00919868d..000000000 --- a/sgl-model-gateway/py_test/integration_mock/test_pd_routing.py +++ /dev/null @@ -1,126 +0,0 @@ -import collections -import concurrent.futures -import time - -import pytest -import requests - - -@pytest.mark.integration -def test_pd_power_of_two_decode_attribution(router_manager, mock_workers): - # Start two prefill and three decode mock workers via fixture - _, prefill_urls_raw, prefill_ids = mock_workers(n=2) - _, decode_urls_raw, decode_ids_list = mock_workers(n=3) - prefill_urls = [(u, None) for u in prefill_urls_raw] - decode_urls = list(decode_urls_raw) - decode_ids = set(decode_ids_list) - - rh = router_manager.start_router( - policy="power_of_two", - pd_disaggregation=True, - prefill_urls=prefill_urls, - decode_urls=decode_urls, - extra={"worker_startup_check_interval": 1}, - ) - - counts = collections.Counter() - with requests.Session() as s: - for i in range(30): - r = s.post( - f"{rh.url}/v1/completions", - json={ - "model": "test-model", - "prompt": f"p{i}", - "max_tokens": 1, - "stream": False, - }, - ) - assert r.status_code == 200 - wid = r.headers.get("X-Worker-Id") or r.json().get("worker_id") - assert wid in decode_ids - counts[wid] += 1 - - assert sum(1 for v in counts.values() if v > 0) >= 2 - - -@pytest.mark.integration -def test_pd_power_of_two_skews_to_faster_decode(router_manager, mock_workers): - # Start two prefill workers (fast) - _, prefill_urls_raw, _ = mock_workers(n=2) - - # Start two decode workers: one slow, one fast - _, [decode_slow_url], [slow_id] = mock_workers( - n=1, args=["--latency-ms", "300"] - ) # slower decode - _, [decode_fast_url], [fast_id] = mock_workers(n=1) - decode_urls_raw = [decode_slow_url, decode_fast_url] - - prefill_urls = [(u, None) for u in prefill_urls_raw] - decode_urls = list(decode_urls_raw) - - rh = router_manager.start_router( - policy="power_of_two", - pd_disaggregation=True, - prefill_urls=prefill_urls, - decode_urls=decode_urls, - extra={"worker_startup_check_interval": 1}, - ) - - def _prime_call(i): - try: - requests.post( - f"{rh.url}/v1/completions", - json={ - "model": "test-model", - "prompt": f"warm-{i}", - "max_tokens": 1, - "stream": False, - }, - timeout=8, - ) - except Exception: - pass - - with concurrent.futures.ThreadPoolExecutor(max_workers=32) as ex: - list(ex.map(_prime_call, range(128))) - time.sleep(2) - - def _direct_decode_load(i): - try: - requests.post( - f"{decode_slow_url}/v1/completions", - json={ - "model": "test-model", - "prompt": f"bg-{i}", - "max_tokens": 1, - "stream": False, - }, - timeout=8, - ) - except Exception: - pass - - with concurrent.futures.ThreadPoolExecutor(max_workers=32) as ex: - list(ex.map(_direct_decode_load, range(128))) - time.sleep(1) - - def call(i): - r = requests.post( - f"{rh.url}/v1/completions", - json={ - "model": "test-model", - "prompt": f"p{i}", - "max_tokens": 1, - "stream": False, - }, - timeout=8, - ) - assert r.status_code == 200 - return r.headers.get("X-Worker-Id") or r.json().get("worker_id") - - counts = collections.Counter() - with concurrent.futures.ThreadPoolExecutor(max_workers=32) as ex: - for wid in ex.map(call, range(200)): - counts[wid] += 1 - - assert counts[slow_id] < counts[fast_id], counts diff --git a/sgl-model-gateway/py_test/integration_mock/test_rate_limiting.py b/sgl-model-gateway/py_test/integration_mock/test_rate_limiting.py deleted file mode 100644 index 960c67a91..000000000 --- a/sgl-model-gateway/py_test/integration_mock/test_rate_limiting.py +++ /dev/null @@ -1,90 +0,0 @@ -import concurrent.futures - -import pytest -import requests - - -@pytest.mark.integration -def test_rate_limit_and_queue(router_manager, mock_workers): - # One fast backend - _, urls, _ = mock_workers(n=1) - rh = router_manager.start_router( - worker_urls=urls, - policy="round_robin", - extra={ - "max_concurrent_requests": 2, - "queue_size": 0, # no queue -> immediate 429 when limit exceeded - }, - ) - - def call_once(i): - try: - r = requests.post( - f"{rh.url}/v1/completions", - json={ - "model": "test-model", - "prompt": f"p{i}", - "max_tokens": 1, - "stream": False, - }, - timeout=3, - ) - return r.status_code - except Exception: - return 599 - - # Fire a burst of concurrent requests - with concurrent.futures.ThreadPoolExecutor(max_workers=16) as ex: - results = list(ex.map(call_once, range(16))) - - # Expect some to succeed and some to be rate limited (429) - assert any(code == 200 for code in results) - assert any(code == 429 for code in results) - - -@pytest.mark.integration -def test_rate_limit_queue_and_timeout(router_manager, mock_workers): - # Slow backend: ~2s per request ensures queue wait > timeout - _, urls, _ = mock_workers(n=1, args=["--latency-ms", "2000"]) # 2.0s per request - - # Allow 1 concurrent, queue up to 1, with 1s queue timeout - rh = router_manager.start_router( - worker_urls=urls, - policy="round_robin", - extra={ - "max_concurrent_requests": 1, - "queue_size": 1, - "queue_timeout_secs": 1, - }, - ) - - def call_once(i): - try: - r = requests.post( - f"{rh.url}/v1/completions", - json={ - "model": "test-model", - "prompt": f"q{i}", - "max_tokens": 1, - "stream": False, - }, - timeout=5, - ) - return r.status_code - except Exception: - return 599 - - # Fire 4 concurrent requests: 1 runs (~2s), 1 queued (times out at 1s -> 408), 2 overflow -> 429 - import concurrent.futures - - with concurrent.futures.ThreadPoolExecutor(max_workers=8) as ex: - results = list(ex.map(call_once, range(4))) - - # We expect: - # - Some 200s (processed) - # - At least one 408 (queued too long and timed out) - # - Remaining non-200s are either 429 (queue overflow) or additional 408s depending on scheduling - assert any(code == 200 for code in results) - assert any(code == 408 for code in results), results - non200 = [c for c in results if c != 200] - assert len(non200) >= 2 and all(c in (408, 429) for c in non200), results diff --git a/sgl-model-gateway/py_test/integration_mock/test_retries.py b/sgl-model-gateway/py_test/integration_mock/test_retries.py deleted file mode 100644 index 0c88ca7d6..000000000 --- a/sgl-model-gateway/py_test/integration_mock/test_retries.py +++ /dev/null @@ -1,71 +0,0 @@ -import pytest -import requests - - -@pytest.mark.integration -def test_retry_reroutes_to_healthy_worker(router_manager, mock_workers): - # Worker A always 500; Worker B healthy - # Worker A always 500; Worker B/C healthy - _, [url_a], [id_a] = mock_workers(n=1, args=["--status-code", "500"]) # fail - _, [url_b], [id_b] = mock_workers(n=1) - _, [url_c], [id_c] = mock_workers(n=1) - rh = router_manager.start_router( - worker_urls=[url_a, url_b, url_c], - policy="round_robin", - extra={ - "retry_max_retries": 3, - "retry_initial_backoff_ms": 10, - "retry_max_backoff_ms": 50, - }, - ) - - r = requests.post( - f"{rh.url}/v1/completions", - json={ - "model": "test-model", - "prompt": "x", - "max_tokens": 1, - "stream": False, - }, - timeout=5, - ) - assert r.status_code == 200 - wid = r.headers.get("X-Worker-Id") or r.json().get("worker_id") - assert wid in [id_b, id_c] # should have retried onto a healthy worker (B or C) - # mock_workers fixture handles cleanup - - -@pytest.mark.integration -def test_disable_retries_surfaces_failure(router_manager, mock_workers): - # Single failing worker, retries disabled -> should return 500 - _, [url], [wid] = mock_workers(n=1, args=["--status-code", "500"]) # always fail - rh = router_manager.start_router( - worker_urls=[url], - policy="round_robin", - extra={ - "disable_retries": True, - }, - ) - - saw_500 = False - for _ in range(8): - r = requests.post( - f"{rh.url}/v1/completions", - json={ - "model": "test-model", - "prompt": "x", - "max_tokens": 1, - "stream": False, - }, - timeout=5, - ) - if r.status_code == 500: - # Worker starts, continue to check - saw_500 = True - break - assert ( - r.status_code == 503 - ), "Should only see 503 when waiting for worker to start" - - assert saw_500 - # mock_workers fixture handles cleanup diff --git a/sgl-model-gateway/py_test/integration_mock/test_service_discovery_shim.py b/sgl-model-gateway/py_test/integration_mock/test_service_discovery_shim.py deleted file mode 100644 index 5cc1d6734..000000000 --- a/sgl-model-gateway/py_test/integration_mock/test_service_discovery_shim.py +++ /dev/null @@ -1,36 +0,0 @@ -import pytest -import requests - - -@pytest.mark.integration -def test_discovery_shim_add_remove(router_manager, mock_workers): - # Start router without workers - rh = router_manager.start_router(worker_urls=[], policy="round_robin") - - # Initially empty - urls = router_manager.list_workers(rh.url) - assert urls == [] - - # Add a worker (simulate discovery event) - _, [wurl], [wid] = mock_workers(n=1) - router_manager.add_worker(rh.url, wurl) - urls = router_manager.list_workers(rh.url) - assert wurl in urls - - # Can serve a request - r = requests.post( - f"{rh.url}/v1/completions", - json={ - "model": "test-model", - "prompt": "hi", - "max_tokens": 1, - "stream": False, - }, - ) - assert r.status_code == 200 - - # Remove worker (simulate pod deletion) - router_manager.remove_worker(rh.url, wurl) - urls = router_manager.list_workers(rh.url) - assert wurl not in urls - # mock_workers fixture handles cleanup diff --git a/sgl-model-gateway/py_test/integration_mock/test_worker_management.py b/sgl-model-gateway/py_test/integration_mock/test_worker_management.py deleted file mode 100644 index 4eace76e0..000000000 --- a/sgl-model-gateway/py_test/integration_mock/test_worker_management.py +++ /dev/null @@ -1,57 +0,0 @@ -import pytest -import requests - - -@pytest.mark.integration -def test_add_and_remove_worker(mock_worker, router_manager, mock_workers): - # Start with a single worker - proc1, url1, id1 = mock_worker - rh = router_manager.start_router(worker_urls=[url1], policy="round_robin") - - # Add a second worker - - procs2, urls2, ids2 = mock_workers(n=1) - url2 = urls2[0] - id2 = ids2[0] - router_manager.add_worker(rh.url, url2) - - # Send some requests and ensure both workers are seen - seen = set() - with requests.Session() as s: - for i in range(20): - r = s.post( - f"{rh.url}/v1/completions", - json={ - "model": "test-model", - "prompt": f"x{i}", - "max_tokens": 1, - "stream": False, - }, - ) - assert r.status_code == 200 - wid = r.headers.get("X-Worker-Id") or r.json().get("worker_id") - seen.add(wid) - if len(seen) == 2: - break - - assert id1 in seen and id2 in seen - - # Now remove the second worker - router_manager.remove_worker(rh.url, url2) - - # After removal, subsequent requests should only come from first worker - with requests.Session() as s: - for i in range(10): - r = s.post( - f"{rh.url}/v1/completions", - json={ - "model": "test-model", - "prompt": f"y{i}", - "max_tokens": 1, - "stream": False, - }, - ) - assert r.status_code == 200 - wid = r.headers.get("X-Worker-Id") or r.json().get("worker_id") - assert wid == id1 - # mock_workers fixture handles cleanup