Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
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
10 changes: 10 additions & 0 deletions components/src/dynamo/mocker/args.py
Original file line number Diff line number Diff line change
Expand Up @@ -285,6 +285,16 @@ def parse_args(argv: list[str] | None = None) -> argparse.Namespace:
default=None,
help="AIC system name (e.g., 'h200_sxm'). Used with --aic-perf-model.",
)
parser.add_argument(
"--aic-backend",
type=str,
default=None,
choices=["vllm", "sglang", "trtllm"],
help="AIC backend name used for perf database lookups. When unset, "
"falls back to --engine-type. Set this to decouple the AIC perf model "
"from the simulated engine type (e.g. simulate with vllm while using "
"trtllm AIC data).",
)
parser.add_argument(
"--aic-backend-version",
type=str,
Expand Down
6 changes: 5 additions & 1 deletion components/src/dynamo/mocker/config.py
Original file line number Diff line number Diff line change
Expand Up @@ -58,7 +58,11 @@ def build_mocker_engine_args(args: argparse.Namespace) -> MockEngineArgs:
aic_moe_ep_size = None
aic_attention_dp_size = None
if getattr(args, "aic_perf_model", False):
aic_backend = getattr(args, "engine_type", None) or "vllm"
aic_backend = (
getattr(args, "aic_backend", None)
or getattr(args, "engine_type", None)
or "vllm"
)
aic_system = getattr(args, "aic_system", None)
aic_backend_version = getattr(args, "aic_backend_version", None)
aic_tp_size = getattr(args, "aic_tp_size", None)
Expand Down
17 changes: 17 additions & 0 deletions components/src/dynamo/mocker/tests/unit/test_config.py
Original file line number Diff line number Diff line change
Expand Up @@ -53,6 +53,7 @@ def make_args(**overrides):
"sglang_schedule_conservativeness": None,
"aic_perf_model": False,
"aic_system": None,
"aic_backend": None,
"aic_backend_version": None,
"aic_tp_size": None,
"aic_moe_tp_size": None,
Expand Down Expand Up @@ -239,6 +240,22 @@ def test_build_mocker_engine_args_preserves_cli_mapped_fields(tmp_path):
assert "has_perf_model" not in payload


def test_aic_backend_override_decouples_from_engine_type():
args = make_args(
engine_type="vllm",
aic_perf_model=True,
aic_system="h200_sxm",
aic_backend="trtllm",
aic_tp_size=4,
)

engine_args = CONFIG.build_mocker_engine_args(args)
payload = json.loads(engine_args.dump_json())

assert payload["engine_type"] == "vllm"
assert payload["aic_backend"] == "trtllm"


def test_mock_engine_args_from_json_ignores_legacy_has_perf_model_field():
payload = {
"engine_type": "vllm",
Expand Down
234 changes: 97 additions & 137 deletions components/src/dynamo/profiler/interpolation.py
Original file line number Diff line number Diff line change
Expand Up @@ -13,7 +13,14 @@
# See the License for the specific language governing permissions and
# limitations under the License.

"""Interpolation curve generation for planner pre-deployment sweeping."""
"""Real-GPU interpolation curve generation for planner thorough-mode sweeps.

Rapid-mode interpolation is no longer generated by the profiler: the planner
runs AIConfigurator in-process at bootstrap (``planner/monitoring/aic_interpolation.py``)
and the mocker pulls AIC perf data at runtime via ``--aic-perf-model`` flags
injected by :func:`dynamo.profiler.utils.dgd_generation.generate_mocker_config`.
This module only handles the thorough path (real GPUs → NPZ on disk).
"""

import logging
import os
Expand All @@ -30,19 +37,12 @@
)
from dynamo.profiler.utils.defaults import EngineType
from dynamo.profiler.utils.dgdr_v1beta1_types import DynamoGraphDeploymentRequestSpec
from dynamo.profiler.utils.estimate_perf import AIConfiguratorPerfEstimator
from dynamo.profiler.utils.profile_common import (
ProfilerOperationalConfig,
inject_tolerations_into_dgd,
)
from dynamo.profiler.utils.profile_decode import (
profile_decode,
profile_decode_aiconfigurator,
)
from dynamo.profiler.utils.profile_prefill import (
profile_prefill,
profile_prefill_aiconfigurator,
)
from dynamo.profiler.utils.profile_decode import profile_decode
from dynamo.profiler.utils.profile_prefill import profile_prefill

logger = logging.getLogger(__name__)

Expand All @@ -53,19 +53,17 @@ async def run_interpolation(
disagg_config: dict,
best_prefill_config: PickedParallelConfig,
best_decode_config: PickedParallelConfig,
model: str,
system: str,
backend: str,
isl: int,
osl: int,
sweep_max_context_length: int,
deployment_clients: list[DynamoDeploymentClient],
job_tolerations: list | None = None,
) -> None:
"""Generate interpolation curves for the planner based on sweep mode.
"""Generate real-GPU interpolation curves for thorough-mode deployments.

Takes the output disagg DGD config and uses ``convert_config`` to strip
it down to standalone prefill / decode engines for profiling.
it down to standalone prefill / decode engines for profiling. Rapid mode
short-circuits here because its interpolation is now handled by the
planner (AIC in-process) and the mocker (``--aic-perf-model`` at runtime).
"""
planner_cfg = (
dgdr.features.planner if (dgdr.features and dgdr.features.planner) else None
Expand All @@ -74,9 +72,11 @@ async def run_interpolation(
if planner_cfg and planner_cfg.pre_deployment_sweeping_mode:
sweep_mode = planner_cfg.pre_deployment_sweeping_mode

if sweep_mode == PlannerPreDeploymentSweepMode.None_:
if sweep_mode != PlannerPreDeploymentSweepMode.Thorough:
logger.info(
"Planner pre-deployment sweeping is disabled — skipping interpolation."
"Skipping real-GPU interpolation for sweep_mode=%s; rapid-mode "
"consumers (planner, mocker) use AIC at runtime.",
sweep_mode,
)
return

Expand All @@ -97,58 +97,42 @@ async def run_interpolation(
with open(prefill_config_fn, "w") as f:
yaml.dump(prefill_config, f)

if sweep_mode == PlannerPreDeploymentSweepMode.Rapid:
logger.info("Using AIC simulation for prefill interpolation.")
estimator = AIConfiguratorPerfEstimator(
hf_id=model,
system=system.lower(),
backend=backend,
)
profile_prefill_aiconfigurator(
work_dir,
best_prefill_gpus,
sweep_max_context_length,
ops.prefill_interpolation_granularity,
estimator,
tp_size=best_prefill_config.tp_size,
)
elif sweep_mode == PlannerPreDeploymentSweepMode.Thorough:
logger.info("Using real GPUs for prefill interpolation.")
frontend_port = config_modifier.get_port(prefill_config)
client = DynamoDeploymentClient(
namespace=ops.k8s_namespace,
base_log_dir=work_dir,
model_name=model_name,
frontend_port=frontend_port,
deployment_name=prefill_config["metadata"]["name"],
)
deployment_clients.append(client)
await client.create_deployment(prefill_config_fn)
logger.info("Waiting for prefill interpolation deployment...")
try:
await client.wait_for_deployment_ready(timeout=ops.deployment_timeout)
except TimeoutError:
logger.error("Prefill interpolation deployment timed out, skipping.")
await client.delete_deployment()
deployment_clients.remove(client)
return

await client.get_deployment_logs()
base_url = client.get_service_url()

profile_prefill(
work_dir,
model_name,
model_path,
base_url,
best_prefill_gpus,
sweep_max_context_length,
ops.prefill_interpolation_granularity,
attention_dp_size=best_prefill_config.dp,
)

logger.info("Using real GPUs for prefill interpolation.")
frontend_port = config_modifier.get_port(prefill_config)
client = DynamoDeploymentClient(
namespace=ops.k8s_namespace,
base_log_dir=work_dir,
model_name=model_name,
frontend_port=frontend_port,
deployment_name=prefill_config["metadata"]["name"],
)
deployment_clients.append(client)
await client.create_deployment(prefill_config_fn)
logger.info("Waiting for prefill interpolation deployment...")
try:
await client.wait_for_deployment_ready(timeout=ops.deployment_timeout)
except TimeoutError:
logger.error("Prefill interpolation deployment timed out, skipping.")
await client.delete_deployment()
deployment_clients.remove(client)
return

await client.get_deployment_logs()
base_url = client.get_service_url()

profile_prefill(
work_dir,
model_name,
model_path,
base_url,
best_prefill_gpus,
sweep_max_context_length,
ops.prefill_interpolation_granularity,
attention_dp_size=best_prefill_config.dp,
)

await client.delete_deployment()
deployment_clients.remove(client)

# --- Decode interpolation ---
decode_config = config_modifier.convert_config(disagg_config, EngineType.DECODE)
Expand All @@ -161,74 +145,50 @@ async def run_interpolation(
with open(decode_config_fn, "w") as f:
yaml.dump(decode_config, f)

if sweep_mode == PlannerPreDeploymentSweepMode.Rapid:
logger.info("Using AIC simulation for decode interpolation.")
estimator = AIConfiguratorPerfEstimator(
hf_id=model,
system=system.lower(),
backend=backend,
)
attention_dp_size = best_decode_config.dp
max_kv_tokens = estimator.get_max_kv_tokens(
isl,
osl,
tp_size=best_decode_config.tp_size,
)
profile_decode_aiconfigurator(
work_dir,
best_decode_gpus,
max_kv_tokens,
sweep_max_context_length,
ops.decode_interpolation_granularity,
estimator,
attention_dp_size,
tp_size=best_decode_config.tp_size,
)
elif sweep_mode == PlannerPreDeploymentSweepMode.Thorough:
logger.info("Using real GPUs for decode interpolation.")
frontend_port = config_modifier.get_port(decode_config)
client = DynamoDeploymentClient(
namespace=ops.k8s_namespace,
base_log_dir=work_dir,
model_name=model_name,
frontend_port=frontend_port,
deployment_name=decode_config["metadata"]["name"],
)
deployment_clients.append(client)
await client.create_deployment(decode_config_fn)
logger.info("Waiting for decode interpolation deployment...")
try:
await client.wait_for_deployment_ready(timeout=ops.deployment_timeout)
except TimeoutError:
logger.error("Decode interpolation deployment timed out, skipping.")
await client.delete_deployment()
deployment_clients.remove(client)
return

await client.get_deployment_logs()

attention_dp_size = best_decode_config.dp
decode_cfg = Config.model_validate(decode_config)
decode_service_name = get_service_name_by_type(
decode_cfg, backend, SubComponentType.DECODE
).lower()
max_kv_tokens = config_modifier.get_kv_cache_size_from_dynamo_log(
f"{work_dir}/{client.deployment_name}/{decode_service_name}/0.log",
attention_dp_size=attention_dp_size,
)
base_url = client.get_service_url()

profile_decode(
work_dir,
model_name,
model_path,
base_url,
best_decode_gpus,
max_kv_tokens,
sweep_max_context_length,
ops.decode_interpolation_granularity,
attention_dp_size,
)

logger.info("Using real GPUs for decode interpolation.")
frontend_port = config_modifier.get_port(decode_config)
client = DynamoDeploymentClient(
namespace=ops.k8s_namespace,
base_log_dir=work_dir,
model_name=model_name,
frontend_port=frontend_port,
deployment_name=decode_config["metadata"]["name"],
)
deployment_clients.append(client)
await client.create_deployment(decode_config_fn)
logger.info("Waiting for decode interpolation deployment...")
try:
await client.wait_for_deployment_ready(timeout=ops.deployment_timeout)
except TimeoutError:
logger.error("Decode interpolation deployment timed out, skipping.")
await client.delete_deployment()
deployment_clients.remove(client)
return

await client.get_deployment_logs()

attention_dp_size = best_decode_config.dp
decode_cfg = Config.model_validate(decode_config)
decode_service_name = get_service_name_by_type(
decode_cfg, backend, SubComponentType.DECODE
).lower()
max_kv_tokens = config_modifier.get_kv_cache_size_from_dynamo_log(
f"{work_dir}/{client.deployment_name}/{decode_service_name}/0.log",
attention_dp_size=attention_dp_size,
)
base_url = client.get_service_url()

profile_decode(
work_dir,
model_name,
model_path,
base_url,
best_decode_gpus,
max_kv_tokens,
sweep_max_context_length,
ops.decode_interpolation_granularity,
attention_dp_size,
)

await client.delete_deployment()
deployment_clients.remove(client)
4 changes: 0 additions & 4 deletions components/src/dynamo/profiler/profile_sla.py
Original file line number Diff line number Diff line change
Expand Up @@ -414,11 +414,7 @@ async def run_profile(
dgd_config,
best_prefill_config,
best_decode_config,
model,
system,
resolved_backend,
isl,
osl,
sweep_max_context_length,
deployment_clients,
job_tolerations=job_tolerations,
Expand Down
Loading
Loading