refactor(runtime): extract PushRouter transport seam behind StreamingDispatch trait - #12447
Merged
Merged
Conversation
|
🎯 Code Coverage (details) 🔗 Commit SHA: 87a8ef8 | Docs | Datadog PR Page | Give us feedback! |
tanmayv25
marked this pull request as ready for review
July 30, 2026 23:38
tanmayv25
marked this pull request as draft
July 30, 2026 23:39
Contributor
WalkthroughChangesStreaming dispatch integration
Estimated code review effort: 3 (Moderate) | ~25 minutes 🚥 Pre-merge checks | ✅ 4 | ❌ 1❌ Failed checks (1 warning)
✅ Passed checks (4 passed)
✨ Finishing Touches 💡 1🛠️ Fix failing CI checks 💡
Comment |
Contributor
There was a problem hiding this comment.
Actionable comments posted: 1
🤖 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 `@lib/runtime/src/pipeline/network/egress/push_router.rs`:
- Around line 326-333: The watcher deduplication in
spawn_instance_removal_watcher must account for dispatch identity, not only
EndpointId. Enforce a single StreamingDispatch per endpoint or update
ENDPOINT_WATCHER_ACTIVE and its ownership checks to distinguish dispatch
instances, ensuring every independent dispatch receives on_instance_removed and
on_instance_added callbacks.
🪄 Autofix (Beta)
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: Path: .coderabbit.yaml
Review profile: CHILL
Plan: Enterprise
Run ID: 8fb16d73-7b31-48de-a591-56e887a90d69
📒 Files selected for processing (3)
lib/runtime/src/pipeline.rslib/runtime/src/pipeline/network/egress/addressed_router.rslib/runtime/src/pipeline/network/egress/push_router.rs
tanmayv25
force-pushed
the
refactor/router-streaming-dispatch-seam
branch
from
July 31, 2026 22:11
d21cfc7 to
d5c3bb0
Compare
tanmayv25
force-pushed
the
refactor/router-streaming-dispatch-seam
branch
2 times, most recently
from
July 31, 2026 22:25
f9851ad to
1d99755
Compare
tanmayv25
marked this pull request as ready for review
August 1, 2026 00:11
…Dispatch trait ## Summary Extract the final-hop transport in `PushRouter` behind a new `StreamingDispatch<T, U>` trait. `PushRouter` keeps instance selection, occupancy, fault detection, and migration; its `addressed` field becomes a trait object so the transport below the seam can be swapped. `AddressedPushRouter` (the request plane) is the default implementation and its behavior is unchanged. This is a behavior-preserving refactor with no new functionality. It makes the final-hop transport pluggable behind a single typed seam; no alternate transport is included in this change. Changes: - New `StreamingDispatch<T, U>` trait (`generate`, `generate_bidirectional`, `on_instance_removed` / `on_instance_added`), with `AddressedPushRouter` as the default impl delegating to its existing `AsyncEngine` path. - `PushRouter.addressed`: `Arc<AddressedPushRouter>` -> `Arc<dyn StreamingDispatch<T, U>>`. - `spawn_instance_removal_watcher` is generic and drives discovery-removal cleanup through the trait. - New `PushRouter::from_client_with_dispatch` to inject a custom dispatch. - Rename inherent `AddressedPushRouter::generate_bidirectional` -> `dispatch_bidirectional` (disambiguates from the trait method); make `AddressedRequest::into_parts` public. ## Validation - `cargo build -p dynamo-runtime` and `-p dynamo-llm` (downstream) compile. - `cargo test -p dynamo-runtime` egress suite: 67 passed, 0 failed — the request-plane path is exercised unchanged. Signed-off-by: tanmayv25 <tanmay2592@gmail.com>
tanmayv25
force-pushed
the
refactor/router-streaming-dispatch-seam
branch
from
August 1, 2026 00:20
1d99755 to
7c412cf
Compare
PeaBrane
approved these changes
Aug 3, 2026
PeaBrane
left a comment
Contributor
There was a problem hiding this comment.
✅ AI-assisted review complete. Approving with two non-blocking comments for human evaluation.
Construct a PushRouter through from_client_with_dispatch and assert a caller-supplied dispatch (not the default AddressedPushRouter) receives unary and bidirectional requests with the selected address/instance, and that discovery removal/re-addition reach its on_instance_removed / on_instance_added hooks. Existing tests only exercised the default adapter. Signed-off-by: tanmayv25 <tanmay2592@gmail.com>
tanmayv25
enabled auto-merge (squash)
August 3, 2026 23:34
hhzhang16
added a commit
that referenced
this pull request
Aug 4, 2026
dyn-3691-extract-shared-target-pid-cuda-customstorage-operation-layer * 'main' of https://github.com/ai-dynamo/dynamo: (50 commits) docs(cli): correct removed vLLM prefill-worker flag reference (#12581) docs(operator): reserve webhook Ignore for emergencies (#12563) ci(docs): make previews and checks match what actually publishes (#12339) refactor(vllm): organize custom encoder modules (#12416) feat(llm): Select reasoning output field via env var (#11464) feat(runtime): add TLS support to TCP request plane (#10921) fix: convert conditional disagg sglang warning to httperror 400 (#12578) feat(operator): add runtime feature gates (#12421) refactor(runtime): extract PushRouter transport seam behind StreamingDispatch trait (#12447) feat(replay): add deterministic canonical offline reports (#12363) build: bump ModelExpress to 0.5.0(OPS-7978) (#12455) fix(mocker): use logical KV tokens for decode timing (#12583) fix(examples): update Triton example for CUDA 13 + fix libdcgm copy (DYN-3697) (#12577) refactor(operator): implement composition-first DGD reconciliation (#12283) feat(frontend): add basetenkenizer backend (#12376) fix(profiler): configure rapid mocker without planner (#12573) docs(vllm): correct worker-role flags and document --kv-transfer-config (#12568) ci: add Kubernetes deploy test to nightly (#12090) fix(container): reuse pinned protoc in runtime image (#12535) feat(self-host): flip DYN_SELF_HOST_METADATA default to ON (gh-8749) (#11417) ... Signed-off-by: Hannah Zhang <hannahz@nvidia.com>
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Summary
Extract the final-hop transport in
PushRouterbehind a newStreamingDispatch<T, U>trait. AfterPushRouterselects a worker, thefinal hop goes through this trait object instead of a fixed
AddressedPushRouter.PushRouterkeeps instance selection, occupancy, faultdetection, and migration; only the transport below the seam can now be swapped.
AddressedPushRouter(the request plane) is the default implementation and itsbehavior is unchanged.
This is a behavior-preserving refactor with no new functionality — it makes
the final-hop transport pluggable behind a single typed seam. No alternate
transport implementation is included here.
Changes
StreamingDispatch<T, U>trait (generate,generate_bidirectional,on_instance_removed/on_instance_added), withAddressedPushRouteras thedefault impl delegating to its existing
AsyncEnginepath.PushRouter.addressed:Arc<AddressedPushRouter>->Arc<dyn StreamingDispatch<T, U>>.spawn_instance_removal_watcheris generic and drives discovery-removalcleanup through the trait.
PushRouter::from_client_with_dispatchto inject a custom dispatch(unused in this PR).
AddressedPushRouter::generate_bidirectional->dispatch_bidirectional(disambiguates from the trait method); makeAddressedRequest::into_partspublic.Validation
cargo build -p dynamo-runtimeand-p dynamo-llm(downstream) compile.cargo test -p dynamo-runtimeegress suite: 67 passed, 0 failed — therequest-plane path is exercised unchanged.