Skip to content
Open
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
7 changes: 7 additions & 0 deletions vllm/config/scheduler.py
Original file line number Diff line number Diff line change
Expand Up @@ -148,6 +148,13 @@ class SchedulerConfig:
avoid gaps in GPU utilization, leading to better latency and throughput.
"""

max_waiting_time: float = 0
"""Maximum time (in seconds) a request can wait in the scheduler's waiting
queue before being evicted. When set to 0 (default), the feature is
disabled and requests wait indefinitely. When set to a positive value,
requests exceeding this wait time are evicted with FINISHED_OVERLOAD status
and the client receives an HTTP 429 Too Many Requests response."""

stream_interval: int = Field(default=1, ge=1)
"""The interval (or buffer size) for streaming in terms of token length.
A smaller value (1) makes streaming smoother by sending each token immediately,
Expand Down
3 changes: 2 additions & 1 deletion vllm/entrypoints/openai/api_server.py
Original file line number Diff line number Diff line change
Expand Up @@ -29,7 +29,7 @@
from vllm.entrypoints.launcher import serve_http
from vllm.entrypoints.logger import RequestLogger
from vllm.entrypoints.openai.cli_args import make_arg_parser, validate_parsed_serve_args
from vllm.entrypoints.openai.engine.protocol import GenerationError
from vllm.entrypoints.openai.engine.protocol import GenerationError, OverloadError
from vllm.entrypoints.openai.models.protocol import BaseModelPath
from vllm.entrypoints.openai.models.serving import OpenAIServingModels
from vllm.entrypoints.openai.server_utils import (
Expand Down Expand Up @@ -250,6 +250,7 @@ def build_app(
app.exception_handler(EngineGenerateError)(engine_error_handler)
app.exception_handler(EngineDeadError)(engine_error_handler)
app.exception_handler(GenerationError)(generation_error_handler)
app.exception_handler(OverloadError)(generation_error_handler)
app.exception_handler(Exception)(exception_handler)

# Ensure --api-key option from CLI takes precedence over VLLM_API_KEY
Expand Down
3 changes: 3 additions & 0 deletions vllm/entrypoints/openai/chat_completion/serving.py
Original file line number Diff line number Diff line change
Expand Up @@ -50,6 +50,7 @@
from vllm.entrypoints.openai.engine.serving import (
GenerationError,
OpenAIServing,
OverloadError,
clamp_prompt_logprobs,
)
from vllm.entrypoints.openai.models.serving import OpenAIServingModels
Expand Down Expand Up @@ -926,6 +927,8 @@ async def chat_completion_stream_generator(

except GenerationError as e:
yield f"data: {self._convert_generation_error_to_streaming_response(e)}\n\n"
except OverloadError as e:
yield f"data: {self._convert_generation_error_to_streaming_response(e)}\n\n"
except Exception as e:
logger.exception("Error in chat completion stream generator.")
data = self.create_streaming_error_response(e)
Expand Down
3 changes: 3 additions & 0 deletions vllm/entrypoints/openai/completion/serving.py
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,7 @@
from vllm.entrypoints.openai.engine.serving import (
GenerationError,
OpenAIServing,
OverloadError,
clamp_prompt_logprobs,
)
from vllm.entrypoints.openai.models.serving import OpenAIServingModels
Expand Down Expand Up @@ -467,6 +468,8 @@ async def completion_stream_generator(

except GenerationError as e:
yield f"data: {self._convert_generation_error_to_streaming_response(e)}\n\n"
except OverloadError as e:
yield f"data: {self._convert_generation_error_to_streaming_response(e)}\n\n"
except Exception as e:
logger.exception("Error in completion stream generator.")
data = self.create_streaming_error_response(e)
Expand Down
9 changes: 9 additions & 0 deletions vllm/entrypoints/openai/engine/protocol.py
Original file line number Diff line number Diff line change
Expand Up @@ -352,3 +352,12 @@ class GenerationError(Exception):
def __init__(self, message: str = "Internal server error"):
super().__init__(message)
self.status_code = HTTPStatus.INTERNAL_SERVER_ERROR


class OverloadError(Exception):
"""raised when finish_reason indicates overload (429)"""

def __init__(self, message: str = "Request exceeded maximum queue waiting "
"time. Please retry later."):
super().__init__(message)
self.status_code = HTTPStatus.TOO_MANY_REQUESTS
19 changes: 15 additions & 4 deletions vllm/entrypoints/openai/engine/serving.py
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,7 @@
from vllm.entrypoints.openai.engine.protocol import (
ErrorResponse,
GenerationError,
OverloadError,
)
from vllm.entrypoints.openai.models.serving import OpenAIServingModels
from vllm.entrypoints.openai.responses.protocol import ResponsesRequest
Expand Down Expand Up @@ -190,21 +191,31 @@ def create_streaming_error_response(
return json_str

def _raise_if_error(self, finish_reason: str | None, request_id: str) -> None:
"""Raise GenerationError if finish_reason indicates an error."""
"""Raise appropriate error if finish_reason indicates an error."""
if finish_reason == "error":
logger.error(
"Request %s failed with an internal error during generation",
request_id,
)
raise GenerationError("Internal server error")
if finish_reason == "overload":
logger.warning(
"Request %s evicted from queue due to overload",
request_id,
)
raise OverloadError()

def _convert_generation_error_to_streaming_response(
self, e: GenerationError
self, e: GenerationError | OverloadError
) -> str:
"""Convert GenerationError to streaming error response."""
"""Convert GenerationError/OverloadError to streaming error response."""
if isinstance(e, OverloadError):
err_type = "OverloadError"
else:
err_type = "InternalServerError"
return self.create_streaming_error_response(
str(e),
err_type="InternalServerError",
err_type=err_type,
status_code=e.status_code,
)

Expand Down
7 changes: 7 additions & 0 deletions vllm/entrypoints/openai/responses/serving.py
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,7 @@
from vllm.entrypoints.openai.engine.serving import (
GenerationError,
OpenAIServing,
OverloadError,
)
from vllm.entrypoints.openai.models.serving import OpenAIServingModels
from vllm.entrypoints.openai.parser.harmony_utils import (
Expand Down Expand Up @@ -1553,6 +1554,12 @@ def _increment_sequence_number_and_return(
TypeAdapter(StreamingResponsesResponse).validate_json(error_json)
)
return
except OverloadError as e:
error_json = self._convert_generation_error_to_streaming_response(e)
yield _increment_sequence_number_and_return(
TypeAdapter(StreamingResponsesResponse).validate_json(error_json)
)
return

async def empty_async_generator():
# A hack to trick Python to think this is a generator but
Expand Down
19 changes: 19 additions & 0 deletions vllm/v1/core/sched/scheduler.py
Original file line number Diff line number Diff line change
Expand Up @@ -108,6 +108,7 @@ def __init__(
else self.scheduler_config.max_num_batched_tokens
)
self.max_model_len = vllm_config.model_config.max_model_len
self.max_waiting_time = self.scheduler_config.max_waiting_time
self.enable_kv_cache_events = (
self.kv_events_config is not None
and self.kv_events_config.enable_kv_cache_events
Expand Down Expand Up @@ -558,6 +559,24 @@ def schedule(self) -> SchedulerOutput:
)
assert len(scheduled_loras) <= self.lora_config.max_loras

# Evict requests that have waited too long in the queue.
if self.max_waiting_time > 0:
now = time.monotonic()
evicted_ids: list[str] = []
for queue in (self.waiting, self.skipped_waiting):
remaining = create_request_queue(self.policy)
while queue:
req = queue.pop_request()
if (now - req.arrival_time) > self.max_waiting_time:
evicted_ids.append(req.request_id)
else:
remaining.append_request(req)
while remaining:
queue.append_request(remaining.pop_request())
for req_id in evicted_ids:
self.finish_requests(req_id,
RequestStatus.FINISHED_OVERLOAD)

# Next, schedule the WAITING requests.
if not preempted_reqs and self._pause_state == PauseState.UNPAUSED:
step_skipped_waiting = create_request_queue(self.policy)
Expand Down
3 changes: 2 additions & 1 deletion vllm/v1/engine/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -27,7 +27,7 @@

# These are possible values of RequestOutput.finish_reason,
# so form part of the external API.
FINISH_REASON_STRINGS = ("stop", "length", "abort", "error", "repetition")
FINISH_REASON_STRINGS = ("stop", "length", "abort", "error", "repetition", "overload")

EEP_NOTIFICATION_CALL_ID = -1

Expand Down Expand Up @@ -59,6 +59,7 @@ class FinishReason(enum.IntEnum):
ABORT = 2
ERROR = 3
REPETITION = 4
OVERLOAD = 5

def __str__(self):
return FINISH_REASON_STRINGS[self.value]
Expand Down
2 changes: 2 additions & 0 deletions vllm/v1/request.py
Original file line number Diff line number Diff line change
Expand Up @@ -333,6 +333,7 @@ class RequestStatus(enum.IntEnum):
FINISHED_IGNORED = enum.auto()
FINISHED_ERROR = enum.auto()
FINISHED_REPETITION = enum.auto()
FINISHED_OVERLOAD = enum.auto()

def __str__(self) -> str:
return self.name
Expand All @@ -358,4 +359,5 @@ def get_finished_reason(status: "RequestStatus") -> FinishReason | None:
RequestStatus.FINISHED_ERROR: FinishReason.ERROR,
RequestStatus.WAITING_FOR_STREAMING_REQ: FinishReason.STOP,
RequestStatus.FINISHED_REPETITION: FinishReason.REPETITION,
RequestStatus.FINISHED_OVERLOAD: FinishReason.OVERLOAD,
}
Loading