feat: Dynamo integration - #1514
AndyDai-nv wants to merge 2 commits into
Conversation
There was a problem hiding this comment.
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.
| 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) |
There was a problem hiding this comment.
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.
| 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
- To prevent resource leaks (e.g., counters that are not decremented), use constructs like
try...finallyor awithstatement to ensure cleanup logic is always executed, even in the case of exceptions or early returns.
| _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", | ||
| ], | ||
| ) |
There was a problem hiding this comment.
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.
| _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
- When using shared temporary directories (such as
/tmp) for caching or process-specific storage, append the user's UID (e.g., usingos.getuid()) to the path to avoid permission conflicts and directory sharing issues on multi-user clusters.
| _sidecars: dict[str, _Sidecar] = {} | ||
|
|
||
|
|
||
| def _ensure_sidecar(name: str, port: int, argv: list[str]) -> _Sidecar: |
There was a problem hiding this comment.
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
- To prevent resource leaks (e.g., counters that are not decremented), use constructs like
try...finallyor awithstatement to ensure cleanup logic is always executed, even in the case of exceptions or early returns.
| logger.info("Starting Dynamo frontend: %s", " ".join(argv)) | ||
| proc = subprocess.Popen(argv, env=env) |
There was a problem hiding this comment.
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.
| 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
- To prevent resource leaks (e.g., counters that are not decremented), use constructs like
try...finallyor awithstatement to ensure cleanup logic is always executed, even in the case of exceptions or early returns.
| 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}") |
There was a problem hiding this comment.
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.
| 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}") |
| try: | ||
| resp.raise_for_status() | ||
| except requests.exceptions.HTTPError as e: | ||
| e.add_note(f"{resp.text=}") | ||
| raise |
There was a problem hiding this comment.
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.
| 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.
08f97de to
ae06806
Compare
|
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. |
Adds
--rollout-backend dynamoas a second option for the rollout side, wrapping NVIDIA Dynamo'sdynamo.frontend+dynamo.sglangworkers. SGLang remains the inference engine; only the router / dispatch layer is swapped. Default backend is unchanged.What this adds
dynamo_router.start_dynamo_router— launchesdynamo.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 wrappingpython -m dynamo.sglangper-engine subprocess. MirrorsSGLangEngine's public surface (init / weight-update / abort / health) soRolloutManagerdoesn't need a special path._generate_via_dynamoinsglang_rollout.py— issues OpenAI-stylePOST /v1/completionsagainst the frontend.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_INTERVALenv passthrough —dynamo_engine.pyforwards this to the spawneddynamo.sglangas--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.