Skip to content

feat: Dynamo integration - #1514

Closed
AndyDai-nv wants to merge 2 commits into
radixark:mainfrom
AndyDai-nv:feat/dynamo-kv-router
Closed

AndyDai-nv wants to merge 2 commits into
radixark:mainfrom
AndyDai-nv:feat/dynamo-kv-router

Conversation

@AndyDai-nv

@AndyDai-nv AndyDai-nv commented Jun 29, 2026 •

Copy link
Copy Markdown

Draft — still iterating.

Adds --rollout-backend dynamo as a second option for the rollout side, wrapping NVIDIA Dynamo's dynamo.frontend + dynamo.sglang workers. SGLang remains the inference engine; only the router / dispatch layer is swapped. Default backend is unchanged.

What this adds

  • dynamo_router.start_dynamo_router — launches dynamo.frontend (plus etcd + nats sidecars when --dynamo-router-mode kv); returns (ip, port) in the same shape as the existing router launcher. Feature-gates optional flags (e.g. --router-predict-on-route) by probing --help.
  • DynamoEngine — Ray actor wrapping python -m dynamo.sglang per-engine subprocess. Mirrors SGLangEngine's public surface (init / weight-update / abort / health) so RolloutManager doesn't need a special path.
  • _generate_via_dynamo in sglang_rollout.py — issues OpenAI-style POST /v1/completions against the frontend.
  • CLI in utils/arguments.py: --rollout-backend, --dynamo-router-mode {kv,round-robin}, --dynamo-discovery-backend {etcd,file}, --dynamo-page-size, plus tunables for predict-on-route.
  • DYN_SGLANG_STREAM_INTERVAL env passthrough — dynamo_engine.py forwards this to the spawned dynamo.sglang as --stream-interval, used to throttle per-token chunk emission for non-streaming RL clients.
  • examples/dynamo_kvrouter/k8s-qwen2.5-0.5B-gsm8k-noncolocate-8gpu-dynamo.yaml — reference deployment.

@gemini-code-assist gemini-code-assist Bot left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Code Review

This pull request introduces a Dynamo-backed inference rollout backend as an alternative to SGLang, adding a new DynamoEngine Ray actor, a Dynamo frontend router launcher, and corresponding Kubernetes configuration and CLI arguments. The review feedback highlights several critical reliability and compatibility improvements: ensuring subprocesses (the engine, frontend, and etcd/NATS sidecars) are properly cleaned up on failure or exit to prevent resource leaks and port conflicts; avoiding multi-user conflicts by appending the user's UID to the temporary etcd data directory; replacing Python 3.11-specific add_note() calls with backwards-compatible exception handling; and avoiding blocking the asyncio event loop with ray.get inside an asynchronous abort function.

Important

The consumer version of Gemini Code Assist on GitHub is being sunset. Starting June 18, 2026, new organization installations will be blocked, and all code review activity will officially cease on July 17, 2026.
For more details on the timeline and next steps, please review the Help Documentation.

Comment on lines +151 to +156
logger.info("DynamoEngine[rank=%d] launching: %s", self.rank, " ".join(argv))
self.process = subprocess.Popen(argv, env=env)

# ``dynamo.sglang``'s system-status server exposes /health (not
# /health_generate which is SGLang-native). Use that for readiness.
self._wait_for_dynamo_health(timeout_s=600)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

high

If _wait_for_dynamo_health raises an exception (such as a timeout or runtime error), the newly launched dynamo.sglang worker subprocess (self.process) is left running in the background. This leaks GPU and system resources, which can cause subsequent runs to fail with Out-Of-Memory (OOM) errors. Wrap the health check in a try...except block and ensure self.shutdown() is called to clean up the process if initialization fails.

Suggested change
logger.info("DynamoEngine[rank=%d] launching: %s", self.rank, " ".join(argv))
self.process = subprocess.Popen(argv, env=env)
# ``dynamo.sglang``'s system-status server exposes /health (not
# /health_generate which is SGLang-native). Use that for readiness.
self._wait_for_dynamo_health(timeout_s=600)
logger.info("DynamoEngine[rank=%d] launching: %s", self.rank, " ".join(argv))
self.process = subprocess.Popen(argv, env=env)
# ``dynamo.sglang``'s system-status server exposes /health (not
# /health_generate which is SGLang-native). Use that for readiness.
try:
self._wait_for_dynamo_health(timeout_s=600)
except Exception:
self.shutdown()
raise
References
  1. To prevent resource leaks (e.g., counters that are not decremented), use constructs like try...finally or a with statement to ensure cleanup logic is always executed, even in the case of exceptions or early returns.

Comment on lines +90 to +102
_ensure_sidecar(
"etcd",
2379,
[
"etcd",
"--listen-client-urls", "http://0.0.0.0:2379",
"--advertise-client-urls", "http://127.0.0.1:2379",
"--listen-peer-urls", "http://0.0.0.0:2380",
"--initial-advertise-peer-urls", "http://127.0.0.1:2380",
"--initial-cluster", "default=http://127.0.0.1:2380",
"--data-dir", "/tmp/dynamo-etcd",
],
)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

high

Using a hardcoded shared temporary directory /tmp/dynamo-etcd for etcd data can lead to permission conflicts and data corruption on multi-user clusters if multiple users run jobs on the same node. As per the general guidelines, append the user's UID (using os.getuid()) to the path to ensure uniqueness and avoid conflicts.

Suggested change
_ensure_sidecar(
"etcd",
2379,
[
"etcd",
"--listen-client-urls", "http://0.0.0.0:2379",
"--advertise-client-urls", "http://127.0.0.1:2379",
"--listen-peer-urls", "http://0.0.0.0:2380",
"--initial-advertise-peer-urls", "http://127.0.0.1:2380",
"--initial-cluster", "default=http://127.0.0.1:2380",
"--data-dir", "/tmp/dynamo-etcd",
],
)
_ensure_sidecar(
"etcd",
2379,
[
"etcd",
"--listen-client-urls", "http://0.0.0.0:2379",
"--advertise-client-urls", "http://127.0.0.1:2379",
"--listen-peer-urls", "http://0.0.0.0:2380",
"--initial-advertise-peer-urls", "http://127.0.0.1:2380",
"--initial-cluster", "default=http://127.0.0.1:2380",
"--data-dir", f"/tmp/dynamo-etcd-{os.getuid()}" if hasattr(os, "getuid") else "/tmp/dynamo-etcd",
],
)
References
  1. When using shared temporary directories (such as /tmp) for caching or process-specific storage, append the user's UID (e.g., using os.getuid()) to the path to avoid permission conflicts and directory sharing issues on multi-user clusters.

Comment on lines +55 to +58
_sidecars: dict[str, _Sidecar] = {}


def _ensure_sidecar(name: str, port: int, argv: list[str]) -> _Sidecar:

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

high

The sidecar processes (etcd and nats-server) started via subprocess.Popen are not daemonized or cleaned up when the main Python process exits. On shared clusters or Kubernetes nodes, these orphaned processes will continue running and permanently block ports 2379 and 4222, preventing future jobs from running. Register an atexit handler to terminate these sidecars on exit.

import atexit

_sidecars: dict[str, _Sidecar] = {}

def _cleanup_sidecars():
    for sc in _sidecars.values():
        if sc.process is not None:
            try:
                sc.process.terminate()
                sc.process.wait(timeout=5)
            except Exception:
                try:
                    sc.process.kill()
                except Exception:
                    pass

atexit.register(_cleanup_sidecars)


def _ensure_sidecar(name: str, port: int, argv: list[str]) -> _Sidecar:
References
  1. To prevent resource leaks (e.g., counters that are not decremented), use constructs like try...finally or a with statement to ensure cleanup logic is always executed, even in the case of exceptions or early returns.

Comment on lines +193 to +194
logger.info("Starting Dynamo frontend: %s", " ".join(argv))
proc = subprocess.Popen(argv, env=env)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

high

The Dynamo frontend process is launched via subprocess.Popen but is never terminated when the main training process exits or crashes. This leaks the frontend process and blocks the router port on the host. Register an atexit handler to clean up the frontend process.

Suggested change
logger.info("Starting Dynamo frontend: %s", " ".join(argv))
proc = subprocess.Popen(argv, env=env)
logger.info("Starting Dynamo frontend: %s", " ".join(argv))
proc = subprocess.Popen(argv, env=env)
import atexit
atexit.register(lambda: proc.terminate() if proc.poll() is None else None)
References
  1. To prevent resource leaks (e.g., counters that are not decremented), use constructs like try...finally or a with statement to ensure cleanup logic is always executed, even in the case of exceptions or early returns.

Comment on lines +509 to +512
try:
ray.get([engine.abort_request.remote(abort_all=True) for engine in engines])
except Exception as e:
logger.warning(f"Dynamo abort_request failed: {e}")

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

high

Calling ray.get inside an async def function blocks the event loop, freezing all other concurrent async tasks (such as in-flight generation requests or HTTP servers) until the remote abort_request calls complete. Since ObjectRef is awaitable in Ray, use await asyncio.gather to asynchronously wait for the abort requests without blocking the event loop.

Suggested change
try:
ray.get([engine.abort_request.remote(abort_all=True) for engine in engines])
except Exception as e:
logger.warning(f"Dynamo abort_request failed: {e}")
try:
await asyncio.gather(*[engine.abort_request.remote(abort_all=True) for engine in engines])
except Exception as e:
logger.warning(f"Dynamo abort_request failed: {e}")

Comment on lines +337 to +341
try:
resp.raise_for_status()
except requests.exceptions.HTTPError as e:
e.add_note(f"{resp.text=}")
raise

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

medium

e.add_note() was introduced in Python 3.11 (PEP 678). Since many machine learning and deep learning environments still run on Python 3.10 or older, calling add_note will raise an AttributeError and mask the original HTTP error. Use a backwards-compatible way to include the response text in the exception, such as raising a new RuntimeError from the original exception.

Suggested change
try:
resp.raise_for_status()
except requests.exceptions.HTTPError as e:
e.add_note(f"{resp.text=}")
raise
try:
resp.raise_for_status()
except requests.exceptions.HTTPError as e:
raise RuntimeError(f"HTTP error calling route {route}: {resp.text}") from e

Drop-in Dynamo frontend + dynamo.sglang workers behind --rollout-backend
dynamo, reusing Miles' router/(ip,port) and Ray-actor contracts:

- dynamo_router.start_dynamo_router: launch dynamo.frontend (+etcd/nats
  sidecars for KV mode) in place of sgl_router
- DynamoEngine: Ray actor wrapping python -m dynamo.sglang, reusing
  SGLang server-arg derivation and control-plane RPCs
- sglang_rollout: dispatch to _generate_via_dynamo (OpenAI /v1/completions)
- arguments: --rollout-backend / --dynamo-* knobs
- docker/Dockerfile.dynamo-kvrouter overlay image
…ynamo KV router

- examples/dynamo_kvrouter/k8s-qwen2.5-0.5B-gsm8k-noncolocate-8gpu-dynamo.yaml:
  reference deployment running Qwen2.5-0.5B + GSM8K on 8 GPUs non-colocate
  (4 train TP=2 DP=2 + 4 rollout 2xTP=2) with --rollout-backend dynamo.
- dynamo_engine.py: forward DYN_SGLANG_STREAM_INTERVAL to dynamo.sglang
  as --stream-interval. Default 50 mirrors SGLang's
  DEFAULT_FORCE_STREAM_INTERVAL and throttles per-token chunk emission
  for non-streaming RL clients (~13% wall-clock win on 0.5B+GSM8K;
  partial-rollout cadence preserved).
- Drop docker/Dockerfile.dynamo-kvrouter; it hardcoded local build paths
  and didn't generalize. Build steps belong in a per-user runbook, not
  the repo.
@AndyDai-nv
AndyDai-nv force-pushed the feat/dynamo-kv-router branch from 08f97de to ae06806 Compare July 1, 2026 21:21
@AndyDai-nv
AndyDai-nv marked this pull request as ready for review July 1, 2026 21:30
@AndyDai-nv

Copy link
Copy Markdown
Author

Closing this pre-refactor Dynamo integration in favor of a fresh plan based on current main's attached external-engine provider (#2514) and the standalone SGLang sidecar. We will retain native SGLang control/weight updates and use Dynamo for request routing; the old dynamo.sglang-owned engine path is no longer the target. Branch retained for reference.

@AndyDai-nv AndyDai-nv closed this Sep 29, 2026
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant