-
Notifications
You must be signed in to change notification settings - Fork 8.1k
[serve] Add request_routing_timeout_s to bound replica assignment #66512
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
base: master
Are you sure you want to change the base?
Changes from all commits
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change | ||||||||||||||||
|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|---|
|
|
@@ -495,6 +495,7 @@ def __init__( | |||||||||||||||||
| initial_backoff_s: float = RAY_SERVE_ROUTER_RETRY_INITIAL_BACKOFF_S, | ||||||||||||||||||
| backoff_multiplier: float = RAY_SERVE_ROUTER_RETRY_BACKOFF_MULTIPLIER, | ||||||||||||||||||
| max_backoff_s: float = RAY_SERVE_ROUTER_RETRY_MAX_BACKOFF_S, | ||||||||||||||||||
| request_routing_timeout_s: Optional[float] = None, | ||||||||||||||||||
| *args, | ||||||||||||||||||
| **kwargs, | ||||||||||||||||||
| ): | ||||||||||||||||||
|
|
@@ -510,6 +511,8 @@ def __init__( | |||||||||||||||||
| self.backoff_multiplier = backoff_multiplier | ||||||||||||||||||
| self.max_backoff_s = max_backoff_s | ||||||||||||||||||
|
|
||||||||||||||||||
| self.request_routing_timeout_s = request_routing_timeout_s | ||||||||||||||||||
|
|
||||||||||||||||||
| # Current replicas available to be routed. | ||||||||||||||||||
| # Updated via `update_replicas`. | ||||||||||||||||||
| self._replica_id_set: Set[ReplicaID] = set() | ||||||||||||||||||
|
|
@@ -1405,13 +1408,34 @@ async def _choose_replica_for_request( | |||||||||||||||||
|
|
||||||||||||||||||
| self._add_pending_request_to_indices(pending_request) | ||||||||||||||||||
| self._maybe_start_routing_tasks() | ||||||||||||||||||
| replica = await pending_request.future | ||||||||||||||||||
| if self.request_routing_timeout_s is None: | ||||||||||||||||||
| replica = await pending_request.future | ||||||||||||||||||
| else: | ||||||||||||||||||
| replica = await asyncio.wait_for( | ||||||||||||||||||
| pending_request.future, self.request_routing_timeout_s | ||||||||||||||||||
| ) | ||||||||||||||||||
| except asyncio.CancelledError as e: | ||||||||||||||||||
| pending_request.future.cancel() | ||||||||||||||||||
| self._cancel_routing_task_for_pending_request(pending_request) | ||||||||||||||||||
| self._remove_pending_request_from_indices(pending_request) | ||||||||||||||||||
|
|
||||||||||||||||||
| raise e from None | ||||||||||||||||||
| except asyncio.TimeoutError: | ||||||||||||||||||
| self._cancel_routing_task_for_pending_request(pending_request) | ||||||||||||||||||
| self._remove_pending_request_from_indices(pending_request) | ||||||||||||||||||
| # No routing task runs while the deployment has no replicas, so the | ||||||||||||||||||
| # lazy cleanup in `_fulfill_pending_requests` can't drop the request. | ||||||||||||||||||
| for queue in ( | ||||||||||||||||||
| self._pending_requests_to_fulfill, | ||||||||||||||||||
| self._pending_requests_to_route, | ||||||||||||||||||
| ): | ||||||||||||||||||
| while queue and queue[0].future.done(): | ||||||||||||||||||
| queue.popleft() | ||||||||||||||||||
|
|
||||||||||||||||||
| raise TimeoutError( | ||||||||||||||||||
| f"Failed to route request to a replica of {self._deployment_id} " | ||||||||||||||||||
| f"within {self.request_routing_timeout_s}s." | ||||||||||||||||||
| ) from None | ||||||||||||||||||
|
Comment on lines
+1435
to
+1438
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. In Python versions prior to 3.11,
Suggested change
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Builtin TimeoutError is intentional. The HTTP/gRPC proxies map There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Timeout path leaves pending request futureHigh Severity The Reviewed by Cursor Bugbot for commit ecff197. Configure here.
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. wait_for already cancels the awaited future on timeout. 3.10/3.11 call fut.cancel(). On 3.12+ asyncio.timeout cancels the task, and Task.cancel() cancels its _fut_waiter, which is this future. Checked on 3.10-3.13: future.cancelled() is True after the timeout. Routing tasks skip done futures, so a later replica goes to a live request. test_request_routing_timeout_no_replicas relies on this when it asserts the route queue is empty. |
||||||||||||||||||
|
|
||||||||||||||||||
| return replica | ||||||||||||||||||
|
|
||||||||||||||||||
|
|
||||||||||||||||||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -274,6 +274,16 @@ class RequestRouterConfig(BaseModel): | |
| ), | ||
| ) | ||
|
|
||
| request_routing_timeout_s: Optional[PositiveFloat] = Field( | ||
| default=None, | ||
| description=( | ||
| "Maximum duration in seconds that a request waits to be assigned a " | ||
| "replica before it fails with a TimeoutError. This bounds routing " | ||
| "only, not the time the replica takes to handle the request. " | ||
| "Defaults to None, meaning a request waits indefinitely." | ||
| ), | ||
| ) | ||
|
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. None timeout breaks proto serializationHigh Severity
Reviewed by Cursor Bugbot for commit ecff197. Configure here.
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. The proto field is |
||
|
|
||
| @field_validator("request_router_kwargs") | ||
| @classmethod | ||
| def request_router_kwargs_json_serializable(cls, v): | ||
|
|
@@ -308,6 +318,7 @@ def __eq__(self, other): | |
| and self.initial_backoff_s == other.initial_backoff_s | ||
| and self.backoff_multiplier == other.backoff_multiplier | ||
| and self.max_backoff_s == other.max_backoff_s | ||
| and self.request_routing_timeout_s == other.request_routing_timeout_s | ||
| ) | ||
|
|
||
| def __hash__(self): | ||
|
|
@@ -329,6 +340,7 @@ def __hash__(self): | |
| self.initial_backoff_s, | ||
| self.backoff_multiplier, | ||
| self.max_backoff_s, | ||
| self.request_routing_timeout_s, | ||
| ) | ||
| ) | ||
|
|
||
|
|
||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -2232,5 +2232,41 @@ def test_compute_backoff_s_does_not_overflow(): | |
| assert router._compute_backoff_s(2048) == router.max_backoff_s | ||
|
|
||
|
|
||
| @pytest.mark.asyncio | ||
| async def test_request_routing_timeout_no_replicas(pow_2_router): | ||
| """A request that isn't assigned a replica within the timeout fails.""" | ||
| s = pow_2_router | ||
| s.request_routing_timeout_s = 0.05 | ||
| loop = get_or_create_event_loop() | ||
|
|
||
| task = loop.create_task(s._choose_replica_for_request(fake_pending_request())) | ||
| with pytest.raises(TimeoutError, match="Failed to route request to a replica"): | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more.
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. The router raises the builtin TimeoutError on purpose (see the thread on request_router.py), so the test stays as is. |
||
| await asyncio.wait_for(task, timeout=10) | ||
|
|
||
| assert s.num_pending_requests == 0 | ||
| assert len(s._pending_requests_to_route) == 0 | ||
|
|
||
|
|
||
| @pytest.mark.asyncio | ||
| async def test_request_routing_timeout_when_replicas_maxed(pow_2_router): | ||
| """A request that times out stops its routing task.""" | ||
| s = pow_2_router | ||
| s.request_routing_timeout_s = 0.05 | ||
| loop = get_or_create_event_loop() | ||
|
|
||
| r1 = FakeRunningReplica("r1") | ||
| r1.set_queue_len_response(DEFAULT_MAX_ONGOING_REQUESTS) | ||
| s.update_replicas([r1]) | ||
|
|
||
| task = loop.create_task(s._choose_replica_for_request(fake_pending_request())) | ||
| with pytest.raises(TimeoutError, match="Failed to route request to a replica"): | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more.
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Same as above: the builtin TimeoutError is intentional, so the test stays as is. |
||
| await asyncio.wait_for(task, timeout=10) | ||
|
|
||
| await async_wait_for_condition( | ||
| lambda: s.curr_num_routing_tasks == 0, retry_interval_ms=1 | ||
| ) | ||
| assert s.num_pending_requests == 0 | ||
|
|
||
|
|
||
| if __name__ == "__main__": | ||
| sys.exit(pytest.main(["-v", "-s", __file__])) | ||
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -135,6 +135,9 @@ message RequestRouterConfig { | |
|
|
||
| // Maximum backoff time (in seconds) between retries. | ||
| double max_backoff_s = 8; | ||
|
|
||
| // Maximum time (in seconds) a request waits to be assigned a replica. | ||
| optional double request_routing_timeout_s = 9; | ||
|
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Proto change needs fault-tolerance reviewLow Severity
This is required by the RPC Fault Tolerance Standards Guide rule because Triggered by project rule: Bugbot Rules Reviewed by Cursor Bugbot for commit 0d9909b. Configure here.
Author
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. No RPC is added or changed. The PR only adds an optional config field to the RequestRouterConfig message. An unset field reads back as None, and older readers ignore it. |
||
| } | ||
| //[End] ROUTING CONFIG | ||
|
|
||
|
|
||


There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
The current cleanup logic only removes done/timed-out requests from the front of the queues (
while queue and queue[0].future.done(): queue.popleft()). If there is a pending request at the front of the queue that has not timed out (e.g., because it has no timeout set), any timed-out requests behind it will remain in the queue indefinitely. This head-of-line blocking leads to a memory leak of timed-out requests and their associated arguments/metadata.To prevent this memory leak, we should filter the queues in-place to remove all done/timed-out requests regardless of their position in the queue.
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
With a timeout set, requests mostly expire in FIFO order, and each expiry pops the done prefix. 50 staggered timeouts with no replicas leave both deques empty. Entries can sit behind a pending head that has a later deadline: a retry (re-inserted by created_at with a fresh timeout), or a head from before the timeout was lowered. They stay only until that head resolves or expires. They're held indefinitely only if the head was enqueued before any timeout was set, and that head waits forever regardless. Cancelled requests already get this lazy cleanup today. Rebuilding both deques on every timeout would be O(n) per timeout under overload, so I kept the lazy pop.