Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
100 changes: 45 additions & 55 deletions e2e_test/bindings_go/conftest.py
Original file line number Diff line number Diff line change
Expand Up @@ -9,10 +9,7 @@

import logging
import os
import signal
import socket
import subprocess
import time
from collections.abc import Generator
from pathlib import Path
from typing import TYPE_CHECKING
Expand All @@ -22,6 +19,9 @@
if TYPE_CHECKING:
from infra import ModelInstance, ModelPool

from infra import get_open_port, release_port, terminate_process
from infra.process_utils import wait_for_health

logger = logging.getLogger(__name__)

# Paths
Expand All @@ -30,25 +30,6 @@
_GO_OAI_SERVER = _GO_BINDINGS / "examples" / "oai_server"


def _find_free_port() -> int:
"""Find an available port."""
with socket.socket(socket.AF_INET, socket.SOCK_STREAM) as s:
s.bind(("", 0))
return s.getsockname()[1]


def _wait_for_server(host: str, port: int, timeout: float = 30.0) -> bool:
"""Wait for server to become available."""
start = time.time()
while time.time() - start < timeout:
try:
with socket.create_connection((host, port), timeout=1.0):
return True
except (TimeoutError, ConnectionRefusedError, OSError):
time.sleep(0.5)
return False


@pytest.fixture(scope="session")
def go_ffi_library() -> Path:
"""Build the Go FFI library and return its directory path."""
Expand Down Expand Up @@ -241,7 +222,7 @@ def go_oai_server(
grpc_endpoint = f"grpc://localhost:{grpc_worker.port}"

# Find a free port for the Go OAI server
oai_port = _find_free_port()
oai_port = get_open_port()

Comment on lines +225 to 226

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

⚠️ Potential issue | 🟠 Major

Ensure reserved ports are released even when process spawn fails.

At Line 225 and Line 299, get_open_port() reserves immediately, but subprocess.Popen(...) is executed before entering try/finally. If spawn fails, release_port() is skipped and reservations leak.

💡 Proposed fix
@@ def go_oai_server(...):
-    process = subprocess.Popen(
-        cmd,
-        env=env,
-        stdout=subprocess.PIPE,
-        stderr=subprocess.PIPE,
-    )
-
-    try:
+    process: subprocess.Popen | None = None
+    try:
+        process = subprocess.Popen(
+            cmd,
+            env=env,
+            stdout=subprocess.PIPE,
+            stderr=subprocess.PIPE,
+        )
@@
     finally:
         logger.info("Shutting down Go OAI server...")
-        terminate_process(process, timeout=10)
+        if process is not None:
+            terminate_process(process, timeout=10)
         release_port(oai_port)

@@ def go_oai_server_multi(...):
-    process = subprocess.Popen(
-        cmd,
-        env=env,
-        stdout=subprocess.PIPE,
-        stderr=subprocess.PIPE,
-    )
-
-    try:
+    process: subprocess.Popen | None = None
+    try:
+        process = subprocess.Popen(
+            cmd,
+            env=env,
+            stdout=subprocess.PIPE,
+            stderr=subprocess.PIPE,
+        )
@@
     finally:
         logger.info("Shutting down Go OAI server...")
-        terminate_process(process, timeout=10)
+        if process is not None:
+            terminate_process(process, timeout=10)
         release_port(oai_port)

Also applies to: 299-300

🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.

In `@e2e_test/bindings_go/conftest.py` around lines 225 - 226, The port
reservation done by get_open_port() can leak if subprocess.Popen(...) raises
before entering the existing try/finally; wrap the spawn in a try/finally so
release_port(port) always runs on failure: call get_open_port() to get oai_port
(and the other reserved port at lines ~299-300), then immediately enter try,
attempt subprocess.Popen(...) inside that try, and in the finally call
release_port(oai_port) (and the corresponding release for the other reserved
port) when the process object was not successfully created or when cleanup is
needed; alternatively move subprocess.Popen(...) inside the existing try block
so that release_port(...) is guaranteed to run in the finally for both
get_open_port()/oai_port and the other reserved port.

# Set up environment - the Go OAI server uses env vars for config
env = os.environ.copy()
Expand All @@ -261,17 +242,26 @@ def go_oai_server(

cmd = [str(go_oai_binary)]

process = subprocess.Popen(
cmd,
env=env,
stdout=subprocess.PIPE,
stderr=subprocess.PIPE,
)

try:
process = subprocess.Popen(
cmd,
env=env,
stdout=subprocess.PIPE,
stderr=subprocess.PIPE,
)

# Wait for server to start
if not _wait_for_server("localhost", oai_port, timeout=30.0):
stdout, stderr = process.communicate(timeout=5)
try:
wait_for_health(f"http://localhost:{oai_port}", timeout=30.0, check_interval=0.5)
except TimeoutError:
try:
stdout, stderr = process.communicate(timeout=5)
except subprocess.TimeoutExpired:
terminate_process(process, timeout=10)
pytest.fail(
f"Go OAI server failed to start and did not exit cleanly.\n"
f"Command: {' '.join(cmd)}"
)
pytest.fail(
f"Go OAI server failed to start.\n"
f"Command: {' '.join(cmd)}\n"
Expand All @@ -283,14 +273,9 @@ def go_oai_server(
yield ("localhost", oai_port, grpc_worker.model_path)

finally:
# Shutdown the server
logger.info("Shutting down Go OAI server...")
process.send_signal(signal.SIGTERM)
try:
process.wait(timeout=10)
except subprocess.TimeoutExpired:
process.kill()
process.wait()
terminate_process(process, timeout=10)
release_port(oai_port)


@pytest.fixture(scope="class")
Expand Down Expand Up @@ -318,7 +303,7 @@ def go_oai_server_multi(
grpc_endpoints = ",".join(f"grpc://localhost:{w.port}" for w in grpc_workers)

# Find a free port for the Go OAI server
oai_port = _find_free_port()
oai_port = get_open_port()

# Set up environment
env = os.environ.copy()
Expand Down Expand Up @@ -350,17 +335,26 @@ def go_oai_server_multi(

cmd = [str(go_oai_binary)]

process = subprocess.Popen(
cmd,
env=env,
stdout=subprocess.PIPE,
stderr=subprocess.PIPE,
)

try:
process = subprocess.Popen(
cmd,
env=env,
stdout=subprocess.PIPE,
stderr=subprocess.PIPE,
)

# Wait for server to start
if not _wait_for_server("localhost", oai_port, timeout=60.0):
stdout, stderr = process.communicate(timeout=5)
try:
wait_for_health(f"http://localhost:{oai_port}", timeout=60.0, check_interval=0.5)
except TimeoutError:
try:
stdout, stderr = process.communicate(timeout=5)
except subprocess.TimeoutExpired:
terminate_process(process, timeout=10)
pytest.fail(
f"Go OAI server failed to start and did not exit cleanly.\n"
f"Command: {' '.join(cmd)}"
)
pytest.fail(
f"Go OAI server failed to start.\n"
f"Command: {' '.join(cmd)}\n"
Expand All @@ -376,12 +370,8 @@ def go_oai_server_multi(

finally:
logger.info("Shutting down Go OAI server...")
process.send_signal(signal.SIGTERM)
try:
process.wait(timeout=10)
except subprocess.TimeoutExpired:
process.kill()
process.wait()
terminate_process(process, timeout=10)
release_port(oai_port)


@pytest.fixture(scope="class")
Expand Down
2 changes: 0 additions & 2 deletions e2e_test/conftest.py
Original file line number Diff line number Diff line change
Expand Up @@ -213,8 +213,6 @@ def pytest_runtest_logstart(nodeid: str, location: tuple) -> None:
from smg_client import SmgClient
from smg_client._errors import SmgError

logger = logging.getLogger(__name__)


@pytest.fixture
def smg(setup_backend):
Expand Down
3 changes: 0 additions & 3 deletions e2e_test/fixtures/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -6,9 +6,6 @@
- setup_backend.py: Backend setup fixtures (class/function-scoped)
- markers.py: Helper utilities for marker extraction

Legacy modules (to be removed during e2e_response_api migration):
- ports.py: Use infra.get_open_port() instead
- router_manager.py: Use infra.Gateway instead
"""

# Pytest hooks (imported by conftest.py via pytest_plugins)
Expand Down
79 changes: 40 additions & 39 deletions e2e_test/fixtures/pool.py
Original file line number Diff line number Diff line change
Expand Up @@ -16,7 +16,7 @@
import pytest

if TYPE_CHECKING:
from infra import ModelPool
from infra import ModelInstance, ModelPool

from .hooks import get_pool_requirements

Expand Down Expand Up @@ -170,44 +170,60 @@ def model_pool(request: pytest.FixtureRequest) -> ModelPool:
return _model_pool


@pytest.fixture
def model_client(
request: pytest.FixtureRequest, model_pool: ModelPool
) -> Generator[object, None, None]:
"""Get OpenAI client for the model specified by @pytest.mark.model().
def _get_model_instance(
request: pytest.FixtureRequest, model_pool: ModelPool, fixture_name: str
) -> ModelInstance:
"""Extract model from marker and acquire instance from pool.

Usage:
@pytest.mark.model("meta-llama/Llama-3.1-8B-Instruct")
def test_chat(model_client):
response = model_client.chat.completions.create(...)
Args:
request: Pytest fixture request.
model_pool: The model pool.
fixture_name: Name of the calling fixture (for error messages).

Returns:
Acquired ModelInstance.
"""
import openai
from infra import PARAM_MODEL
from infra import PARAM_MODEL, ConnectionMode

marker = request.node.get_closest_marker(PARAM_MODEL)
if marker is None:
if marker is None or not marker.args:
pytest.fail(
f"Test must be marked with @pytest.mark.{PARAM_MODEL}('model-id') "
"to use model_client fixture"
f"to use {fixture_name} fixture"
)

model_id = marker.args[0]

Comment thread
coderabbitai[bot] marked this conversation as resolved.
try:
# get() auto-acquires the returned instance
instance = model_pool.get(model_id)
return model_pool.get(model_id, ConnectionMode.HTTP)
except KeyError:
Comment thread
coderabbitai[bot] marked this conversation as resolved.
pytest.skip(f"Model {model_id} not available in model pool")
Comment on lines 197 to 200

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

⚠️ Potential issue | 🟠 Major

Prefer explicit failure over skip when required model acquisition fails.

At Line 200, pytest.skip(...) can hide infra/config regressions; this path should fail fast with a clear message.

💡 Proposed fix
     try:
         return model_pool.get(model_id, ConnectionMode.HTTP)
     except KeyError:
-        pytest.skip(f"Model {model_id} not available in model pool")
+        pytest.fail(
+            f"Model {model_id} not available in model pool for {fixture_name}; "
+            "check model prelaunch requirements and test markers."
+        )

Based on learnings: tests should fail explicitly when required services are unavailable rather than being skipped to catch misconfigurations.

📝 Committable suggestion

‼️ IMPORTANT
Carefully review the code before committing. Ensure that it accurately replaces the highlighted code, contains no missing lines, and has no issues with indentation. Thoroughly test & benchmark the code to ensure it meets the requirements.

Suggested change
try:
# get() auto-acquires the returned instance
instance = model_pool.get(model_id)
return model_pool.get(model_id, ConnectionMode.HTTP)
except KeyError:
pytest.skip(f"Model {model_id} not available in model pool")
try:
return model_pool.get(model_id, ConnectionMode.HTTP)
except KeyError:
pytest.fail(
f"Model {model_id} not available in model pool for {fixture_name}; "
"check model prelaunch requirements and test markers."
)
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.

In `@e2e_test/fixtures/pool.py` around lines 197 - 200, The current except
KeyError branch in the model acquisition code swallows a missing-model condition
by calling pytest.skip; instead fail fast so infra/config regressions surface.
Replace the pytest.skip(...) call with pytest.fail(...) (or raise a clear
assertion) in the except KeyError handler around the model_pool.get(model_id,
ConnectionMode.HTTP) call so that attempting to acquire a required model
(model_id) triggers an explicit test failure rather than a skip; update
references to pytest.skip -> pytest.fail and keep the original error message.



@pytest.fixture
def model_client(
request: pytest.FixtureRequest, model_pool: ModelPool
) -> Generator[object, None, None]:
"""Get OpenAI client for the model specified by @pytest.mark.model().

Usage:
@pytest.mark.model("meta-llama/Llama-3.1-8B-Instruct")
def test_chat(model_client):
response = model_client.chat.completions.create(...)
"""
import openai

instance = _get_model_instance(request, model_pool, "model_client")

client = openai.OpenAI(
base_url=f"{instance.base_url}/v1",
api_key="not-used",
)

yield client

# Release reference to allow eviction
instance.release()
try:
yield client
finally:
instance.release()


@pytest.fixture
Expand All @@ -221,24 +237,9 @@ def model_base_url(
def test_direct_http(model_base_url):
response = httpx.get(f"{model_base_url}/health")
"""
from infra import PARAM_MODEL

marker = request.node.get_closest_marker(PARAM_MODEL)
if marker is None:
pytest.fail(
f"Test must be marked with @pytest.mark.{PARAM_MODEL}('model-id') "
"to use model_base_url fixture"
)

model_id = marker.args[0]
instance = _get_model_instance(request, model_pool, "model_base_url")

try:
# get() auto-acquires the returned instance
instance = model_pool.get(model_id)
except KeyError:
pytest.skip(f"Model {model_id} not available in model pool")

yield instance.base_url

# Release reference to allow eviction
instance.release()
yield instance.base_url
finally:
instance.release()
14 changes: 0 additions & 14 deletions e2e_test/fixtures/ports.py

This file was deleted.

Loading