Skip to content

feat(router): pre-scoring worker filter stage with label matching - #2118

Closed
pallasathena92 wants to merge 2 commits into
mainfrom
feat/worker-filter-stage
Closed

pallasathena92 wants to merge 2 commits into
mainfrom
feat/worker-filter-stage

Conversation

@pallasathena92

Copy link
Copy Markdown
Collaborator

Description

Problem

The only universal candidate narrowing before policy scoring is health + circuit breaker. Real narrowing concerns — tenant pinning, canary cohorts, hardware classes, adapter residency — have no seam: each would end up hacked into individual policies' select_worker. And when narrowing exhausts the pool, there is no honest status to return: today's paths would report 404 ("model not found" — a lie) or generic unavailability.

Solution

A worker-filter chain on the PolicyRegistry, applied between eligibility and the policy (and the routing-key sticky override) on every policy-based selection path — gRPC regular/PD/EPD, HTTP regular, HTTP PD — including the EPD per-item encode assignment, which previously called the policy directly and bypassed the registry. Design points:

  • Filters are hard requirements. A worker that fails one may not serve the request; preferences (ordering, weighting) stay in policies. Policies are untouched — they already only ever narrow the slice they're handed.
  • Index mapping. The registry hands the policy the narrowed slice and maps the returned index back to the caller's slice. (Documented interplay: with a chain active, the index-pinning X-SMG-Target-Worker header addresses the post-filter slice.)
  • Three-way outcome. SelectionOutcome::{Selected, NotSelected, AllFiltered} replaces Option<usize>. AllFiltered (candidates existed, filters removed all) maps to 503 workers_filtered at every call site — it is unavailability, not backpressure (429 would misreport it, and retrying elsewhere just bounces) and not model absence (404 would lie). NotSelected preserves each path's existing empty-pool/policy-miss behavior exactly.
  • Zero-cost when off. An empty chain short-circuits; nothing is cloned.

Shipped concrete filter — label matching by request header. A new RouterConfig option (--worker-filter-header, unset = off, validated as a header name at startup) names a header whose value is comma-separated key=value pairs; only workers whose spec labels contain ALL pairs survive. Requests without the header pass everything through; malformed pairs are ignored. Worker labels already flow end to end (k8s discovery, /workers PATCH), so tenant pinning / canary cohorts / hardware classes work with zero new worker-side surface.

Out of scope (stated deliberately): the external-proxy routers (OpenAI/Anthropic/Gemini/realtime) use a separate policy-free least-loaded selector and are untouched.

Changes

  • model_gateway/src/policies/filter.rs (new): WorkerFilter trait, LabelHeaderFilter, worker_filters_from_config.
  • model_gateway/src/policies/registry.rs: SelectionOutcome, filter chain + set_worker_filters, select_worker filtering/mapping (sticky override factored into select_with_override).
  • Call sites: gRPC worker_selection.rs (three modes + encode assignment routed through the registry; SelectionMiss → 404 vs 503 mapping), HTTP router.rs (both selection sites), HTTP pd_router.rs.
  • Config: worker_filter_header field + validation (header-name check) + builder + CLI + wiring in app_context.

Test Plan

  • Filter unit tests: missing header keeps all; all-pairs semantics (partial/wrong-value/unlabeled workers drop); malformed pairs ignored; empty-value matching.
  • Registry tests: narrowing + index mapping back to the caller's slice; AllFiltered vs NotSelected distinguished (empty pool still NotSelected with filters installed); empty chain is a passthrough (round-robin still alternates); label filter end-to-end through the registry incl. inert-without-header.
  • Config tests: header-name validation (invalid/empty rejected, valid/None accepted).
  • cargo test -p smg --lib 1485 passed; --test routing_tests 102 passed; cargo clippy -p smg --all-targets clean.
Checklist
  • Format your code: make fmt
  • Run lint checks: cargo clippy -p smg --all-targets -- -D warnings
  • Add unit tests for new functionality
  • Update documentation if needed (CLI help text carries the flag docs)

The only universal candidate narrowing before policy scoring is health
plus circuit breaker. Real narrowing concerns — tenant pinning, canary
cohorts, hardware classes, adapter residency — have no seam: each would
end up hacked into individual policies.

Add a worker-filter chain on the PolicyRegistry, applied between
eligibility and the policy (and the routing-key override) on every
policy-based selection path: gRPC regular/PD/EPD (including the per-item
encode assignment, which previously bypassed the registry), HTTP
regular, and HTTP PD. Filters are hard requirements; preferences stay in
policies. An empty chain costs nothing. The registry hands the policy
the narrowed slice and maps the returned index back to the caller's.

Selection now reports a three-way outcome so filter exhaustion is
distinguishable: candidates-existed-but-all-filtered maps to 503
(workers_filtered) at every call site — unavailability, not
backpressure (429 would misreport it) and not model absence (404 would
lie). The empty-pool / policy-declined case keeps each path's existing
behavior exactly.

Ships one concrete filter: label matching driven by a request header.
A new config option names the header (unset = feature off, validated
as a header name at startup); the header value is comma-separated
key=value pairs and only workers whose spec labels contain all pairs
survive. Requests without the header are unaffected. Worker labels
already flow end to end (k8s discovery, the /workers API), so this
needs no new worker-side surface.

Signed-off-by: yifeng liu <31553858+pallasathena92@users.noreply.github.com>
@coderabbitai

coderabbitai Bot commented Aug 12, 2026 •

Copy link
Copy Markdown

Review Change Stack

Warning

Review limit reached

@pallasathena92, you've reached your PR review limit, so we couldn't start this review.

Next review available in: 48 minutes

You've used all free OSS reviews for now. Wait for the free limit to reset to keep reviewing this public repository.

How can I continue?

After more reviews become available, a review can be triggered using the @coderabbitai review command as a PR comment. Alternatively, push new commits to this PR.

To avoid repeated limits, reduce automatic review volume by pausing incremental auto-reviews earlier, using label-based review opt-in, excluding WIP or generated PR titles, or requesting reviews manually when the PR is ready. If your team needs uninterrupted high-volume reviews, an organization admin can enable usage-based reviews.

How do review limits work?

CodeRabbit enforces per-developer PR review limits for each organization. Most developers receive the normal plan review availability.

For paid Pro and Pro+ PR reviews, CodeRabbit uses adaptive limits for sustained high-volume activity. When a developer's recent PR review activity reaches the 95th percentile or higher among CodeRabbit users, additional reviews become available more gradually as earlier reviews age out of the rolling window.

Please refer docs for additional details.

Review details
⚙️ Run configuration

Configuration used: Organization UI

Review profile: CHILL

Plan: Pro Plus

Run ID: 78fa7c4b-93d0-4bd4-b92f-670094011f35

📥 Commits

Reviewing files that changed from the base of the PR and between caf2b91 and f915f37.

📒 Files selected for processing (3)
  • bindings/python/src/lib.rs
  • bindings/python/src/smg/router_args.py
  • model_gateway/src/routers/http/pd_router.rs
📝 Walkthrough

Summary by CodeRabbit

  • New Features

    • Added optional request-header worker filtering using comma-separated label requirements.
    • Requests now route only to workers matching all valid label filters.
    • Added configuration and command-line options for specifying the filter header.
  • Bug Fixes

    • Improved routing responses by distinguishing unavailable workers from workers excluded by filters.
    • Added validation for filter header names and clearer service-unavailable handling when all workers are filtered out.

Walkthrough

The gateway adds configurable label-based worker filtering from request headers. PolicyRegistry applies filters before policy selection and returns structured outcomes. HTTP and gRPC routes map complete filtering to service-unavailable responses while preserving existing handling for ordinary selection misses.

Changes

Worker filtering

Layer / File(s) Summary
Filter configuration and construction
model_gateway/src/config/*, model_gateway/src/main.rs, model_gateway/src/app_context.rs, model_gateway/src/policies/filter.rs, model_gateway/src/policies/mod.rs
Adds the optional worker-filter header to configuration and CLI wiring. Validates header names, parses label requirements, constructs filters, and installs them in the registry.
Filtered policy selection
model_gateway/src/policies/registry.rs
Applies ordered filters before routing-key and policy selection. Remaps filtered indices and returns Selected, NotSelected, or AllFiltered.
gRPC selection outcome propagation
model_gateway/src/routers/grpc/common/stages/worker_selection.rs
Uses structured selection results for regular, PD, EPD, and encode selection. Centralizes conversion to client responses.
HTTP filtered-worker responses
model_gateway/src/routers/http/router.rs, model_gateway/src/routers/http/pd_router.rs
Returns workers_filtered 503 responses when request filters exclude all candidates. Existing policy and worker-unavailable responses remain distinct.

Estimated code review effort: 4 (Complex) | ~45 minutes

Sequence Diagram(s)

sequenceDiagram
  participant Client
  participant Router
  participant PolicyRegistry
  participant LabelHeaderFilter
  participant Worker
  Client->>Router: Send request with worker-filter header
  Router->>PolicyRegistry: Select worker candidates
  PolicyRegistry->>LabelHeaderFilter: Apply label requirements
  LabelHeaderFilter-->>PolicyRegistry: Return matching candidates
  PolicyRegistry->>Worker: Apply routing-key or load-balancing policy
  Worker-->>Router: Return selected worker or selection outcome
  Router-->>Client: Return response or workers_filtered 503
Loading

Possibly related PRs

Suggested labels: tests

Suggested reviewers: slin1237, catherinesue, key4ng

🚥 Pre-merge checks | ✅ 5
✅ Passed checks (5 passed)
Check name Status Explanation
Title check ✅ Passed The title clearly summarizes the main change: adding a pre-scoring worker filter stage with label matching.
Description check ✅ Passed The description directly explains the worker-filter chain, label matching, selection outcomes, configuration, routing changes, and tests.
Docstring Coverage ✅ Passed Docstring coverage is 100.00% which is sufficient. The required threshold is 80.00%.
Linked Issues check ✅ Passed Check skipped because no linked issues were found for this pull request.
Out of Scope Changes check ✅ Passed Check skipped because no linked issues were found for this pull request.
✨ Finishing Touches
📝 Generate docstrings
  • Create stacked PR
  • Commit on current branch
🧪 Generate unit tests (beta)
  • Create PR with unit tests
  • Commit unit tests in branch feat/worker-filter-stage

Thanks for using CodeRabbit! It's free for OSS, and your support helps us grow. If you like it, consider giving us a shout-out.

❤️ Share

Comment @coderabbitai help to get the list of available commands.

@github-actions github-actions Bot added grpc gRPC client and router changes model-gateway Model gateway crate changes labels Aug 12, 2026
Comment on lines +945 to +949
crate::policies::SelectionOutcome::AllFiltered => {
return Err(format!(
"All candidate {worker_type} workers were excluded by request worker filters"
));
}

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

🟡 Nit: This string error flows through handle_server_selection_error (line 175), which wraps all errors with the code "server_selection_failed". The gRPC path (selection_miss_response) and the HTTP regular path (SelectMiss::AllFiltered) both map filter exhaustion to the "workers_filtered" error code — the three-way outcome was designed precisely for that distinction, but the PD path loses it by collapsing back into an untyped string.

Practically the HTTP status is correct (503 in all cases), so this is about error-code consistency for clients that inspect the body.

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

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

Agreed — the distinction was the point. pick_worker_by_policy_arc / select_pd_pair now return a typed PdSelectionError (Unavailable vs AllFiltered): ordinary unavailability keeps the historical server_selection_failed wrapping, filter exhaustion maps to 503 workers_filtered like the gRPC and HTTP-regular paths.

@claude claude Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

Clean, well-designed PR. The three-way SelectionOutcome is the right abstraction for distinguishing filter exhaustion from empty-pool/policy-miss, and the index-mapping logic in the registry is correct. Test coverage is solid across filter semantics, registry narrowing, and config validation.

One 🟡 nit posted: the HTTP PD path loses the workers_filtered error code because the AllFiltered variant collapses into the same string-error pathway as other selection failures. HTTP status (503) is correct in all paths.

Summary: 0 🔴 Important · 1 🟡 Nit · 0 🟣 Pre-existing

@slin1237
slin1237 marked this pull request as ready for review August 12, 2026 15:05

@coderabbitai coderabbitai Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

Actionable comments posted: 2

Caution

Some comments are outside the diff and can’t be posted inline due to platform limitations.

⚠️ Outside diff range comments (2)
model_gateway/src/routers/grpc/common/stages/worker_selection.rs (2)

703-721: 🎯 Functional Correctness | 🟠 Major | ⚡ Quick win

🔴 Important — Preserve the client filter header for encode selection.

assign_encode_workers replaces the client headers with encode_routing_headers. A configured label filter therefore receives no --worker-filter-header value for encode workers and passes every encode candidate. Prefill and decode still use the client headers.

Pass the request headers into assign_encode_workers. Retain the configured filter header and overlay the synthetic routing-key header for each item. Add an EPD test that proves encode workers also match the requested labels.

As per coding guidelines, “Account for dual-dispatch complexity introduced by PD disaggregation in addition to regular routing.”

🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@model_gateway/src/routers/grpc/common/stages/worker_selection.rs` around
lines 703 - 721, Update assign_encode_workers and its callers to accept the
client request headers, preserving the configured worker-filter header while
overlaying each item’s synthetic content-hash routing header when constructing
SelectWorkerInfo. Keep prefill and decode behavior unchanged, and add an EPD
test proving encode selection honors the requested worker labels under dual
dispatch.

Source: Coding guidelines


333-401: 🎯 Functional Correctness | 🟠 Major | 🏗️ Heavy lift

🔴 Important — Apply worker filters before selecting the shared runtime.

Both paths derive target_runtime from the raw pools, then apply filters only during policy selection. If the first runtime has workers that the request filter excludes, but another runtime has a complete eligible set, the request returns workers_filtered instead of selecting the eligible runtime.

Expose or reuse filtering before computing the runtime intersection. Then run the role policies on those narrowed pools. Add PD and EPD tests with two runtimes where only the second runtime satisfies the request labels.

  • model_gateway/src/routers/grpc/common/stages/worker_selection.rs#L333-L401: derive the PD runtime from filter-eligible prefill and decode pools.
  • model_gateway/src/routers/grpc/common/stages/worker_selection.rs#L495-L595: derive the EPD runtime from filter-eligible encode, prefill, and decode pools.

As per coding guidelines, “Account for dual-dispatch complexity introduced by PD disaggregation in addition to regular routing.”

🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

In `@model_gateway/src/routers/grpc/common/stages/worker_selection.rs` around
lines 333 - 401, The PD selection flow at worker_selection.rs:333-401 must apply
request worker filters before deriving target_runtime, then select policies from
the filtered prefill and decode pools; add coverage for two runtimes where only
the second has eligible labels. Apply the same root-cause fix to the EPD flow at
worker_selection.rs:495-595 by filtering encode, prefill, and decode pools
before computing their shared runtime, and add the corresponding two-runtime
label test.

Source: Coding guidelines

🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.

Inline comments:
In `@model_gateway/src/main.rs`:
- Around line 235-240: Expose the existing worker_filter_header option through
the Python router configuration: add the corresponding field to RouterArgs, pass
it from Router::to_router_config() into the router builder/configuration, and
add coverage verifying propagation. Do not add this option to the minimal Go
example.

In `@model_gateway/src/routers/http/pd_router.rs`:
- Around line 926-950: Preserve SelectionOutcome::AllFiltered as a structured
error through pick_worker_by_policy_arc and select_pd_pair instead of converting
it to String. In execute_dual_dispatch, map that error to the required
workers_filtered response code while retaining server_selection_failed for other
selection failures. Add regression coverage for worker-filter exhaustion on both
prefill and decode PD legs.

---

Outside diff comments:
In `@model_gateway/src/routers/grpc/common/stages/worker_selection.rs`:
- Around line 703-721: Update assign_encode_workers and its callers to accept
the client request headers, preserving the configured worker-filter header while
overlaying each item’s synthetic content-hash routing header when constructing
SelectWorkerInfo. Keep prefill and decode behavior unchanged, and add an EPD
test proving encode selection honors the requested worker labels under dual
dispatch.
- Around line 333-401: The PD selection flow at worker_selection.rs:333-401 must
apply request worker filters before deriving target_runtime, then select
policies from the filtered prefill and decode pools; add coverage for two
runtimes where only the second has eligible labels. Apply the same root-cause
fix to the EPD flow at worker_selection.rs:495-595 by filtering encode, prefill,
and decode pools before computing their shared runtime, and add the
corresponding two-runtime label test.
🪄 Autofix

Fix all unresolved CodeRabbit comments on this PR:

  • Push a commit to this branch (recommended)
  • Create a new PR with the fixes

ℹ️ Review info
⚙️ Run configuration

Configuration used: Organization UI

Review profile: CHILL

Plan: Pro Plus

Run ID: a810be04-5d04-4ef1-b6e4-61867a8e1b9a

📥 Commits

Reviewing files that changed from the base of the PR and between caf4fb2 and caf2b91.

📒 Files selected for processing (11)
  • model_gateway/src/app_context.rs
  • model_gateway/src/config/builder.rs
  • model_gateway/src/config/types.rs
  • model_gateway/src/config/validation.rs
  • model_gateway/src/main.rs
  • model_gateway/src/policies/filter.rs
  • model_gateway/src/policies/mod.rs
  • model_gateway/src/policies/registry.rs
  • model_gateway/src/routers/grpc/common/stages/worker_selection.rs
  • model_gateway/src/routers/http/pd_router.rs
  • model_gateway/src/routers/http/router.rs

Comment thread model_gateway/src/main.rs
Comment thread model_gateway/src/routers/http/pd_router.rs
The PD path collapsed selection failures into a bare String, so filter
exhaustion lost its distinct client-facing code and surfaced as
server_selection_failed. Type the error (Unavailable vs AllFiltered):
ordinary unavailability keeps its historical wrapping, filter exhaustion
maps to 503 workers_filtered like the gRPC and HTTP-regular paths.

Also expose worker_filter_header through the Python bindings
(RouterArgs field + --worker-filter-header + PyO3 constructor,
appended at list tails to preserve positional callers).

Signed-off-by: yifeng liu <31553858+pallasathena92@users.noreply.github.com>
@github-actions github-actions Bot added the python-bindings Python bindings changes label Aug 12, 2026
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

grpc gRPC client and router changes model-gateway Model gateway crate changes python-bindings Python bindings changes

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants