Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
15 commits
Select commit Hold shift + click to select a range
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
13 changes: 13 additions & 0 deletions components/src/dynamo/mocker/args.py
Original file line number Diff line number Diff line change
Expand Up @@ -544,6 +544,19 @@ def parse_args(argv: list[str] | None = None) -> argparse.Namespace:
"Default: 64.0 (inter-node InfiniBand). Set to 0 to disable KV transfer delay. "
"For intra-node NVLink, typical value is ~450.",
)
parser.add_argument(
"--kv-transfer-abort-timeout-ms",
type=int,
default=None,
help="Prefill-to-decode handshake timeout in milliseconds. When set, "
"the prefill worker holds KV cache after compute until decode connects to its "
"bootstrap server, up to this timeout. If decode does not arrive in time, the "
"prefill aborts the request (KV released, abort error surfaced); late-arriving "
"decodes for the same room get a clean ABORT response. Mirrors production "
"VLLM_NIXL_ABORT_REQUEST_TIMEOUT. Decode-side: must wait for local KV "
"availability before connecting to prefill. Default: None (legacy behavior, "
"prefill ACKs immediately).",
)
parser.add_argument(
"--kv-cache-dtype",
type=str,
Expand Down
14 changes: 14 additions & 0 deletions components/src/dynamo/mocker/config.py
Original file line number Diff line number Diff line change
Expand Up @@ -306,6 +306,9 @@ def build_mocker_engine_args(args: argparse.Namespace) -> MockEngineArgs:
bandwidth_g3_to_g2_gbps=getattr(args, "bandwidth_g3_to_g2_gbps", None),
bandwidth_g2_to_g4_gbps=getattr(args, "bandwidth_g2_to_g4_gbps", None),
bandwidth_g4_to_g2_gbps=getattr(args, "bandwidth_g4_to_g2_gbps", None),
kv_transfer_abort_timeout_ms=getattr(
args, "kv_transfer_abort_timeout_ms", None
),
reasoning=_parse_reasoning_config(getattr(args, "reasoning", None)),
sglang=_build_sglang_args(args),
trtllm=_build_trtllm_args(args),
Expand Down Expand Up @@ -345,6 +348,7 @@ def apply_worker_engine_args_overrides(

def build_runtime_config(
engine_args: MockEngineArgs,
engine_type: str = "vllm",
) -> tuple[int, ModelRuntimeConfig]:
rc = ModelRuntimeConfig()
# Mocker does not enforce a model context limit.
Expand All @@ -363,12 +367,22 @@ def build_runtime_config(

bootstrap_port = engine_args.bootstrap_port
if engine_args.is_prefill() and bootstrap_port is not None:
# Bootstrap rendezvous (sglang, or vLLM with an explicit --bootstrap-ports):
# advertise host/port so the (unchanged) prefill_router takes the
# Resolved/bootstrap branch. UNCHANGED behavior.
host = os.environ.get(
"DYN_HTTP_RPC_HOST", socket.gethostbyname(socket.gethostname())
)
rc.set_disaggregated_endpoint(
bootstrap_host=host, bootstrap_port=bootstrap_port
)
elif engine_args.is_prefill() and engine_type == "vllm":
# vLLM disagg WITHOUT a bootstrap address: register as a disagg/prefill
# worker with no bootstrap host/port so the unchanged prefill_router takes
# the NoBootstrapEndpoint -> output-`disaggregated_params` (NIXL) path
# (execution.rs:139,306). No router change; only which existing branch is
# taken. sglang never reaches here (it always sets bootstrap_port above).
rc.set_disaggregated_endpoint(bootstrap_host=None, bootstrap_port=None)

apply_topology_config(rc)

Expand Down
4 changes: 3 additions & 1 deletion components/src/dynamo/mocker/main.py
Original file line number Diff line number Diff line change
Expand Up @@ -196,7 +196,9 @@ async def launch_workers(args: argparse.Namespace, base_engine_args):
else:
worker_engine_args = base_engine_args

kv_cache_block_size, runtime_config = build_runtime_config(worker_engine_args)
kv_cache_block_size, runtime_config = build_runtime_config(
worker_engine_args, getattr(args, "engine_type", None) or "vllm"
)

# Create EntrypointArgs for this worker
entrypoint_args = EntrypointArgs(
Expand Down
4 changes: 3 additions & 1 deletion lib/bindings/python/rust/llm/replay.rs
Original file line number Diff line number Diff line change
Expand Up @@ -171,7 +171,7 @@ impl MockEngineArgs {
#[pymethods]
impl MockEngineArgs {
#[new]
#[pyo3(signature = (engine_type="vllm", num_gpu_blocks=None, block_size=0, max_num_seqs=Some(256), max_num_batched_tokens=Some(8192), enable_prefix_caching=true, enable_chunked_prefill=true, speedup_ratio=1.0, decode_speedup_ratio=1.0, dp_size=1, startup_time=None, worker_type="aggregated", planner_profile_data=None, aic_backend=None, aic_system=None, aic_backend_version=None, aic_tp_size=None, aic_model_path=None, aic_moe_tp_size=None, aic_moe_ep_size=None, aic_attention_dp_size=None, aic_nextn=None, aic_nextn_accept_rates=None, aic_mtp_seed=42, gpu_memory_utilization=None, mem_fraction_static=None, free_gpu_memory_fraction=None, enable_local_indexer=false, bootstrap_port=None, kv_bytes_per_token=None, kv_transfer_bandwidth=None, reasoning=None, zmq_kv_events_port=None, zmq_replay_port=None, preemption_mode="lifo", router_queue_policy=None, sglang=None, trtllm=None, num_g2_blocks=None, num_g3_blocks=None, offload_batch_size=None, bandwidth_g1_to_g2_gbps=None, bandwidth_g2_to_g1_gbps=None, bandwidth_g2_to_g3_gbps=None, bandwidth_g3_to_g2_gbps=None, enable_g4_storage=false, bandwidth_g2_to_g4_gbps=None, bandwidth_g4_to_g2_gbps=None))]
#[pyo3(signature = (engine_type="vllm", num_gpu_blocks=None, block_size=0, max_num_seqs=Some(256), max_num_batched_tokens=Some(8192), enable_prefix_caching=true, enable_chunked_prefill=true, speedup_ratio=1.0, decode_speedup_ratio=1.0, dp_size=1, startup_time=None, worker_type="aggregated", planner_profile_data=None, aic_backend=None, aic_system=None, aic_backend_version=None, aic_tp_size=None, aic_model_path=None, aic_moe_tp_size=None, aic_moe_ep_size=None, aic_attention_dp_size=None, aic_nextn=None, aic_nextn_accept_rates=None, aic_mtp_seed=42, gpu_memory_utilization=None, mem_fraction_static=None, free_gpu_memory_fraction=None, enable_local_indexer=false, bootstrap_port=None, kv_bytes_per_token=None, kv_transfer_bandwidth=None, kv_transfer_abort_timeout_ms=None, reasoning=None, zmq_kv_events_port=None, zmq_replay_port=None, preemption_mode="lifo", router_queue_policy=None, sglang=None, trtllm=None, num_g2_blocks=None, num_g3_blocks=None, offload_batch_size=None, bandwidth_g1_to_g2_gbps=None, bandwidth_g2_to_g1_gbps=None, bandwidth_g2_to_g3_gbps=None, bandwidth_g3_to_g2_gbps=None, enable_g4_storage=false, bandwidth_g2_to_g4_gbps=None, bandwidth_g4_to_g2_gbps=None))]
#[allow(clippy::too_many_arguments)]
fn new(
engine_type: &str,
Expand Down Expand Up @@ -205,6 +205,7 @@ impl MockEngineArgs {
bootstrap_port: Option<u16>,
kv_bytes_per_token: Option<usize>,
kv_transfer_bandwidth: Option<f64>,
kv_transfer_abort_timeout_ms: Option<u64>,
Comment thread
nnshah1 marked this conversation as resolved.
reasoning: Option<ReasoningConfig>,
zmq_kv_events_port: Option<u16>,
zmq_replay_port: Option<u16>,
Expand Down Expand Up @@ -275,6 +276,7 @@ impl MockEngineArgs {
.bandwidth_g3_to_g2_gbps(bandwidth_g3_to_g2_gbps)
.bandwidth_g2_to_g4_gbps(bandwidth_g2_to_g4_gbps)
.bandwidth_g4_to_g2_gbps(bandwidth_g4_to_g2_gbps)
.kv_transfer_abort_timeout_ms(kv_transfer_abort_timeout_ms)
.reasoning(reasoning.map(|config| config.inner()))
.zmq_kv_events_port(zmq_kv_events_port)
.zmq_replay_port(zmq_replay_port)
Expand Down
1 change: 1 addition & 0 deletions lib/bindings/python/src/dynamo/_core.pyi
Original file line number Diff line number Diff line change
Expand Up @@ -1893,6 +1893,7 @@ class MockEngineArgs:
bootstrap_port: Optional[int] = None,
kv_bytes_per_token: Optional[int] = None,
kv_transfer_bandwidth: Optional[float] = None,
kv_transfer_abort_timeout_ms: Optional[int] = None,
reasoning: Optional[ReasoningConfig] = None,
zmq_kv_events_port: Optional[int] = None,
zmq_replay_port: Optional[int] = None,
Expand Down
Loading
Loading