diff --git a/projects/caliper/prometheus_metrics/queries.yaml b/projects/caliper/prometheus_metrics/queries.yaml index 67f73e525..cf5bf07e4 100644 --- a/projects/caliper/prometheus_metrics/queries.yaml +++ b/projects/caliper/prometheus_metrics/queries.yaml @@ -301,3 +301,155 @@ queries: promql: DCGM_FI_DEV_GPU_TEMP unit: celsius description: GPU temperature + + # ── Dynamo Frontend ───────────────────────────────────────────────── + + dynamo_requests_total: + category: dynamo_frontend + promql: dynamo_frontend_requests_total{namespace=~"{ns}"} + unit: count + description: Total requests handled by Dynamo frontend + + dynamo_request_rate: + category: dynamo_frontend + promql: >- + sum(rate(dynamo_frontend_requests_total + {namespace=~"{ns}"}[5m])) + unit: req/s + description: Dynamo frontend request rate + + dynamo_input_token_rate: + category: dynamo_frontend + promql: >- + sum(rate(dynamo_frontend_input_sequence_tokens_sum + {namespace=~"{ns}"}[5m])) + unit: tokens/s + description: Dynamo frontend input token throughput + + dynamo_output_token_rate: + category: dynamo_frontend + promql: >- + sum(rate(dynamo_frontend_output_sequence_tokens_sum + {namespace=~"{ns}"}[5m])) + unit: tokens/s + description: Dynamo frontend output token throughput + + dynamo_queued_requests: + category: dynamo_frontend + promql: dynamo_frontend_queued_requests{namespace=~"{ns}"} + unit: count + description: Requests queued in Dynamo frontend + + dynamo_ttft_avg: + category: dynamo_frontend + promql: >- + rate(dynamo_frontend_time_to_first_token_seconds_sum + {namespace=~"{ns}"}[5m]) + / clamp_min(rate(dynamo_frontend_time_to_first_token_seconds_count + {namespace=~"{ns}"}[5m]), 1) + unit: seconds + description: Dynamo average time to first token + + dynamo_itl_avg: + category: dynamo_frontend + promql: >- + rate(dynamo_frontend_inter_token_latency_seconds_sum + {namespace=~"{ns}"}[5m]) + / clamp_min(rate(dynamo_frontend_inter_token_latency_seconds_count + {namespace=~"{ns}"}[5m]), 1) + unit: seconds + description: Dynamo average inter-token latency + + dynamo_request_duration_avg: + category: dynamo_frontend + promql: >- + rate(dynamo_frontend_request_duration_seconds_sum + {namespace=~"{ns}"}[5m]) + / clamp_min(rate(dynamo_frontend_request_duration_seconds_count + {namespace=~"{ns}"}[5m]), 1) + unit: seconds + description: Dynamo average request duration + + dynamo_avg_input_tokens_per_request: + category: dynamo_frontend + promql: >- + rate(dynamo_frontend_input_sequence_tokens_sum + {namespace=~"{ns}"}[5m]) + / clamp_min(rate(dynamo_frontend_input_sequence_tokens_count + {namespace=~"{ns}"}[5m]), 1) + unit: tokens + description: Average input tokens per request + + dynamo_avg_output_tokens_per_request: + category: dynamo_frontend + promql: >- + rate(dynamo_frontend_output_sequence_tokens_sum + {namespace=~"{ns}"}[5m]) + / clamp_min(rate(dynamo_frontend_output_sequence_tokens_count + {namespace=~"{ns}"}[5m]), 1) + unit: tokens + description: Average output tokens per request + + # ── Dynamo Component ──────────────────────────────────────────────── + + dynamo_component_request_rate: + category: dynamo_component + promql: >- + sum by (dynamo_component) (rate(dynamo_component_requests_total + {namespace=~"{ns}"}[5m])) + unit: req/s + description: Request rate per Dynamo component + + dynamo_component_request_duration_avg: + category: dynamo_component + promql: >- + sum by (dynamo_component) ( + rate(dynamo_component_request_duration_seconds_sum + {namespace=~"{ns}"}[5m]) + / clamp_min(rate(dynamo_component_request_duration_seconds_count + {namespace=~"{ns}"}[5m]), 1) + ) + unit: seconds + description: Average request duration per Dynamo component + + dynamo_component_request_bytes_rate: + category: dynamo_component + promql: >- + sum by (dynamo_component) (rate(dynamo_component_request_bytes_total + {namespace=~"{ns}"}[5m])) + unit: bytes/s + description: Request bytes throughput per Dynamo component + + dynamo_component_response_bytes_rate: + category: dynamo_component + promql: >- + sum by (dynamo_component) (rate(dynamo_component_response_bytes_total + {namespace=~"{ns}"}[5m])) + unit: bytes/s + description: Response bytes throughput per Dynamo component + + dynamo_gpu_cache_usage_percent: + category: dynamo_component + promql: dynamo_component_gpu_cache_usage_percent{namespace=~"{ns}"} + unit: percent + description: Dynamo GPU KV cache usage percentage + + dynamo_total_blocks: + category: dynamo_component + promql: dynamo_component_total_blocks{namespace=~"{ns}"} + unit: count + description: Dynamo total KV cache blocks + + # ── Dynamo Operator ───────────────────────────────────────────────── + + dynamo_operator_reconcile_rate: + category: dynamo_operator + promql: sum(rate(dynamo_operator_reconcile_total[5m])) + unit: ops/s + description: Dynamo operator reconciliation rate + + dynamo_operator_reconcile_errors: + category: dynamo_operator + promql: sum(rate(dynamo_operator_reconcile_errors_total[5m])) + unit: errors/s + description: Dynamo operator reconciliation error rate diff --git a/projects/dynamo/orchestration/ci.py b/projects/dynamo/orchestration/ci.py new file mode 100755 index 000000000..628f06ad0 --- /dev/null +++ b/projects/dynamo/orchestration/ci.py @@ -0,0 +1,115 @@ +#!/usr/bin/env python3 +""" +Dynamo Project CI Operations +""" + +import logging +import types +from pathlib import Path + +import click + +from projects.core.agentic.config_review import trigger_config_review_for_ci +from projects.core.agentic.on_failure import agent_review_on_failure +from projects.core.ci_entrypoint.fournos_resolve import create_fournos_resolve_entrypoint +from projects.core.library import ci as ci_lib +from projects.core.library import config, env, run, vault +from projects.core.library.export import caliper_export_entrypoint +from projects.core.library.replot import caliper_replot_entrypoint +from projects.dynamo.orchestration.cleanup_phase import run as cleanup_toolbox_run +from projects.dynamo.orchestration.preflight_phase import run as preflight_toolbox_run +from projects.dynamo.orchestration.prepare_sequence import run_prepare_sequence +from projects.dynamo.orchestration.test_phase import run as test_toolbox_run + +logger = logging.getLogger(__name__) + + +def init(): + """Initialize Dynamo orchestration environment""" + env.init() + run.init() + config.init(Path(__file__).parent) + + +def list_vaults() -> list[str]: + """List all vaults (includes both mandatory and optional).""" + return vault.phase_vault_list_all() + + +@click.group(cls=ci_lib.HelpfulGroup) +@click.pass_context +@ci_lib.safe_ci_function +def main(ctx): + """Dynamo Project CI Operations for FORGE.""" + ctx.ensure_object(types.SimpleNamespace) + init() + + if ctx.invoked_subcommand == "resolve-fournos-config": + logger.info("No need to initialize the vaults for the resolve step") + return + + vault.phase_vault_init(ctx.invoked_subcommand) + + +@main.command() +@click.pass_context +@ci_lib.safe_ci_command +@agent_review_on_failure +def prepare(ctx) -> int: + """Prepare phase - Set up environment, operators, and Dynamo platform.""" + return run_prepare_sequence() + + +@main.command() +@click.pass_context +@ci_lib.safe_ci_command +@agent_review_on_failure +def preflight(ctx) -> int: + """Preflight check phase - Validate required Dynamo CRDs exist.""" + return preflight_toolbox_run() + + +@main.command() +@click.pass_context +@ci_lib.safe_ci_command +@agent_review_on_failure +def test(ctx) -> int: + """Test phase - Deploy DynamoGraphDeployment, run smoke test and benchmarks.""" + trigger_config_review_for_ci(env.BASE_ARTIFACT_DIR, async_mode=True) + return test_toolbox_run() + + +@main.command() +@click.pass_context +@ci_lib.safe_ci_command +@agent_review_on_failure +def pre_cleanup(ctx) -> int: + """Cleanup phase - Clean up Dynamo test resources.""" + from projects.dynamo.orchestration import runtime_config + + for run_spec in runtime_config.get_run_specs(): + with runtime_config.activate_run_spec(run_spec): + cleanup_toolbox_run(namespace=run_spec.namespace) + return 0 + + +@main.command() +@click.pass_context +@ci_lib.safe_ci_command +@agent_review_on_failure +def post_cleanup(ctx) -> int: + """Cleanup phase - Final cleanup of Dynamo resources.""" + from projects.dynamo.orchestration import runtime_config + + for run_spec in runtime_config.get_run_specs(): + with runtime_config.activate_run_spec(run_spec): + cleanup_toolbox_run(namespace=run_spec.namespace) + return 0 + + +main.add_command(create_fournos_resolve_entrypoint(vault_list_func=vault.phase_vault_list_all)) +main.add_command(caliper_export_entrypoint) +main.add_command(caliper_replot_entrypoint) + +if __name__ == "__main__": + main() diff --git a/projects/dynamo/orchestration/cleanup_phase.py b/projects/dynamo/orchestration/cleanup_phase.py new file mode 100644 index 000000000..fafb14e18 --- /dev/null +++ b/projects/dynamo/orchestration/cleanup_phase.py @@ -0,0 +1,32 @@ +from __future__ import annotations + +import logging + +from projects.core.dsl.utils.k8s import oc_resource_exists +from projects.dynamo.toolbox.cleanup_dynamo_resources.main import ( + run as cleanup_dynamo_resources_toolbox_run, +) + +logger = logging.getLogger(__name__) + + +def run(*, namespace: str | None = None) -> int: + """Delete Dynamo test resources from a namespace.""" + from projects.dynamo.orchestration import runtime_config + + if namespace is None: + namespace = runtime_config.get_namespace() + + if not oc_resource_exists("namespace", namespace): + logger.info("Namespace %s does not exist, nothing to clean up", namespace) + return 0 + + benchmark_job_names = runtime_config.get_benchmark_job_names() + benchmark_name = benchmark_job_names[0] if benchmark_job_names else None + + cleanup_dynamo_resources_toolbox_run( + namespace=namespace, + benchmark_job_name=benchmark_name, + ) + + return 0 diff --git a/projects/dynamo/orchestration/cli.py b/projects/dynamo/orchestration/cli.py new file mode 100755 index 000000000..decde3e85 --- /dev/null +++ b/projects/dynamo/orchestration/cli.py @@ -0,0 +1,102 @@ +#!/usr/bin/env python3 +""" +Dynamo Project CLI entrypoint +""" + +import logging +import sys +import types +from pathlib import Path + +import click + +from projects.core.library import config, env, run +from projects.core.library.cli import safe_cli_command +from projects.core.library.postprocess import postprocess_command + +logger = logging.getLogger(__name__) + + +def init(): + """Initialize Dynamo orchestration environment""" + env.init() + run.init() + config.init(Path(__file__).parent) + + +@click.group() +@click.option( + "--preset", + multiple=True, + help="Apply a preset to the configuration. Pass multiple --preset NAME to apply multiple presets.", +) +@click.pass_context +def main(ctx, preset): + """Dynamo CLI Operations.""" + ctx.ensure_object(types.SimpleNamespace) + init() + + if not preset: + return + + try: + for preset_name in preset: + logger.info(f"Applying preset: {preset_name}") + config.project.apply_preset(preset_name) + except ValueError as e: + logger.error(f"Failed to apply preset '{preset_name}': {e}") + sys.exit(1) + + +@main.command() +@click.pass_context +@safe_cli_command +def prepare(ctx): + """Prepare phase - Set up operators and Dynamo platform.""" + from projects.dynamo.orchestration.prepare_sequence import run_prepare_sequence + + exit_code = run_prepare_sequence() + sys.exit(exit_code) + + +@main.command() +@click.pass_context +@safe_cli_command +def preflight(ctx): + """Preflight check - Validate required Dynamo CRDs exist.""" + from projects.dynamo.orchestration.preflight_phase import run as preflight_run + + exit_code = preflight_run() + sys.exit(exit_code) + + +@main.command() +@click.pass_context +@safe_cli_command +def test(ctx): + """Test phase - Deploy DynamoGraphDeployment, smoke test, benchmark.""" + from projects.dynamo.orchestration.test_phase import run as test_run + + exit_code = test_run() + sys.exit(exit_code) + + +@main.command() +@click.pass_context +@safe_cli_command +def cleanup(ctx): + """Cleanup phase - Remove Dynamo test resources.""" + from projects.dynamo.orchestration import runtime_config + from projects.dynamo.orchestration.cleanup_phase import run as cleanup_run + + for run_spec in runtime_config.get_run_specs(): + with runtime_config.activate_run_spec(run_spec): + cleanup_run(namespace=run_spec.namespace) + sys.exit(0) + + +main.add_command(postprocess_command) + + +if __name__ == "__main__": + main() diff --git a/projects/dynamo/orchestration/config.d/caliper.yaml b/projects/dynamo/orchestration/config.d/caliper.yaml new file mode 100644 index 000000000..d89801810 --- /dev/null +++ b/projects/dynamo/orchestration/config.d/caliper.yaml @@ -0,0 +1,27 @@ +caliper: + postprocess: + enabled: false + export: + from: null + verbose: true + dry_run: false + backend: + mlflow: + enabled: true + secrets: + vault: + name: psap-forge-mlflow-export + mlflow_secret: mlflow-secret.yaml + config: + experiment: forge-dynamo + workspace: forge-dynamo + run_name: + description: > + Caliper export from FORGE Dynamo project + source_script: dynamo project + tags: {} + parameters: {} + metrics: {} + notifications: + enabled: true + vault: psap-forge-notifications diff --git a/projects/dynamo/orchestration/config.d/deployments.yaml b/projects/dynamo/orchestration/config.d/deployments.yaml new file mode 100644 index 000000000..37cff372a --- /dev/null +++ b/projects/dynamo/orchestration/config.d/deployments.yaml @@ -0,0 +1,71 @@ +defaults: + backend_framework: vllm + replicas: 1 + tensor_parallelism: 1 + serving_mode: aggregated + router_mode: direct + runtime_image: nvcr.io/nvidia/ai-dynamo/vllm-runtime:1.2.1 + frontend_image: nvcr.io/nvidia/ai-dynamo/dynamo-frontend:1.2.1 + hf_home: /opt/models + use_dra_resources: false + dra_resource_name: dra.llm-d.io/gpu-nic-pair + kv_transfer_config: null + kvbm_cpu_cache_gb: null + kv_block_size: "16" + frontend_cpu: "4" + router_args: [] + vllm_args: + - --gpu-memory-utilization=0.90 + - --enable-prefix-caching + - --block-size=16 + env: + DYN_STORE_KV: mem + DYN_DECODE_FALLBACK: "true" + +aggregated: + serving_mode: aggregated + +aggregated-kv-router: + serving_mode: aggregated + router_mode: kv + router_args: + - --active-decode-blocks-threshold + - None + - --active-prefill-tokens-threshold + - None + +aggregated-kvbm: + serving_mode: aggregated + router_mode: kv + kvbm_cpu_cache_gb: 100 + kv_transfer_config: '{"kv_connector":"DynamoConnector","kv_role":"kv_both"}' + +aggregated-dra: + serving_mode: aggregated + use_dra_resources: true + +aggregated-rr-4r-tp2: + serving_mode: aggregated + router_mode: round-robin + replicas: 4 + tensor_parallelism: 2 + runtime_image: nvcr.io/nvidia/ai-dynamo/vllm-runtime:1.1.1 + frontend_image: nvcr.io/nvidia/ai-dynamo/dynamo-frontend:1.1.1 + hf_home: /home/dynamo/.cache/huggingface + kv_block_size: "64" + frontend_cpu: "4" + vllm_args: + - --gpu-memory-utilization=0.90 + - --async-scheduling + - --block-size=64 + - --max-model-len=131072 + - "--hf-overrides={\"rope_scaling\":{\"rope_type\":\"yarn\",\"factor\":4.0,\"original_max_position_embeddings\":32768},\"max_position_embeddings\":131072}" + - --no-enable-log-requests + env: + DYN_STORE_KV: mem + DYN_HEALTH_CHECK_ENABLED: "false" + +disaggregated: + serving_mode: disaggregated + prefill_replicas: 1 + decode_replicas: 2 diff --git a/projects/dynamo/orchestration/config.d/metrics.yaml b/projects/dynamo/orchestration/config.d/metrics.yaml new file mode 100644 index 000000000..2a2d77488 --- /dev/null +++ b/projects/dynamo/orchestration/config.d/metrics.yaml @@ -0,0 +1,32 @@ +enabled: true +step_seconds: 15 +query_keys: + # Dynamo frontend + - dynamo_requests_total + - dynamo_request_rate + - dynamo_input_token_rate + - dynamo_output_token_rate + - dynamo_queued_requests + - dynamo_ttft_avg + - dynamo_itl_avg + - dynamo_request_duration_avg + - dynamo_avg_input_tokens_per_request + - dynamo_avg_output_tokens_per_request + # Dynamo component + - dynamo_component_request_rate + - dynamo_component_request_duration_avg + - dynamo_gpu_cache_usage_percent + - dynamo_total_blocks + # Dynamo operator + - dynamo_operator_reconcile_rate + - dynamo_operator_reconcile_errors + # Infrastructure + - gpu_utilization + - gpu_memory_used + - gpu_memory_free + - gpu_power + - cpu_usage + - memory_usage + - network_rx_bytes + - network_tx_bytes + - container_restarts diff --git a/projects/dynamo/orchestration/config.d/model_cache.yaml b/projects/dynamo/orchestration/config.d/model_cache.yaml new file mode 100644 index 000000000..faa195605 --- /dev/null +++ b/projects/dynamo/orchestration/config.d/model_cache.yaml @@ -0,0 +1,20 @@ +enabled: true +marker_filename: .forge-model-cache.json +pvc_name: dynamo-model-cache + +pvc: + name_prefix: dynamo-model + size: 15Gi + access_mode: ReadWriteOnce + storage_class_name: null + model_directory_name: model + +download: + wait_timeout_seconds: 7200 + poll_interval_seconds: 15 + pod_image_pull_policy: IfNotPresent + +hf: + downloader_image: registry.access.redhat.com/ubi9/python-311 + token_secret_name: null + token_secret_key: token diff --git a/projects/dynamo/orchestration/config.d/platform.yaml b/projects/dynamo/orchestration/config.d/platform.yaml new file mode 100644 index 000000000..2893514d2 --- /dev/null +++ b/projects/dynamo/orchestration/config.d/platform.yaml @@ -0,0 +1,54 @@ +cluster: + minimum_openshift_version: "4.17.0" + namespace: + name: forge-dynamo + prefix: dynamo + max_length: 63 + gpu_node_label_selector: nvidia.com/gpu.present=true + nfd_gpu_detection_labels: + - feature.node.kubernetes.io/pci-10de.present + - feature.node.kubernetes.io/pci-0302_10de.present + - feature.node.kubernetes.io/pci-0300_10de.present + +operators: + nfd: + display_name: Node Feature Discovery + namespace: openshift-nfd + channel: stable + source: redhat-operators + bootstrap_crd: nodefeaturediscoveries.nfd.openshift.io + gpu-operator-certified: + display_name: NVIDIA GPU Operator + namespace: nvidia-gpu-operator + channel: stable + source: certified-operators + bootstrap_crd: clusterpolicies.nvidia.com + +dynamo: + helm: + chart_path: null + chart_repo: https://helm.ngc.nvidia.com/nvidia/dynamo + chart_name: dynamo-platform + chart_version: "1.2.1" + release_name: dynamo + namespace: dynamo-system + values_override: + global: + etcd: + install: true + nats: + install: true + operator: + wait_timeout_seconds: 300 + required_crds: + - dynamographdeployments.nvidia.com + - dynamocomponentdeployments.nvidia.com + +smoke: + pod_name: dynamo-smoke + client_image: curlimages/curl:8.11.1 + endpoint_path: /v1/completions + request_timeout_seconds: 120 + +artifacts: + capture_namespace_events: true diff --git a/projects/dynamo/orchestration/config.d/project.yaml b/projects/dynamo/orchestration/config.d/project.yaml new file mode 100644 index 000000000..92523283b --- /dev/null +++ b/projects/dynamo/orchestration/config.d/project.yaml @@ -0,0 +1,3 @@ +name: dynamo +args: + - smoke diff --git a/projects/dynamo/orchestration/config.d/runtime.yaml b/projects/dynamo/orchestration/config.d/runtime.yaml new file mode 100644 index 000000000..2aa9f4fdd --- /dev/null +++ b/projects/dynamo/orchestration/config.d/runtime.yaml @@ -0,0 +1,8 @@ +selected_preset: null +model_name: Qwen/Qwen3-0.6B +deployment_profile: aggregated +smoke_request_key: default +benchmark_tool: null +benchmark_key: null +job_name: null +namespace_override: null diff --git a/projects/dynamo/orchestration/config.d/workloads.yaml b/projects/dynamo/orchestration/config.d/workloads.yaml new file mode 100644 index 000000000..aa744d951 --- /dev/null +++ b/projects/dynamo/orchestration/config.d/workloads.yaml @@ -0,0 +1,81 @@ +smoke_requests: + default: + prompt: San Francisco is a + max_tokens: 50 + temperature: 0.7 + +guidellm_benchmarks: + short: + job_name: guidellm-benchmark + image: ghcr.io/vllm-project/guidellm:v0.5.4 + pvc_size: 1Gi + timeout_seconds: 900 + outputs: json + args: + backend_type: openai_http + rate: 1 + rate_type: concurrent + max_seconds: 120 + sample_requests: 20 + data: prompt_tokens=256,output_tokens=128 + + concurrent-1k-1k: + job_name: guidellm-benchmark + image: ghcr.io/vllm-project/guidellm:v0.5.4 + pvc_size: 1Gi + timeout_seconds: 3600 + args: + backend_type: openai_http + rate_type: concurrent + rate: [300, 200, 100, 50, 1] + data: prompt_tokens=1000,output_tokens=1000 + max_seconds: 600 + + multi-turn: + job_name: guidellm-benchmark + image: ghcr.io/vllm-project/guidellm:v0.5.4 + pvc_size: 1Gi + timeout_seconds: 3600 + args: + backend_type: openai_http + rate_type: concurrent + rate: [32, 64, 128, 256, 512] + data: prompt_tokens=128,output_tokens=128,turns=5,prefix_tokens=10000,prefix_count={2*rate} + max_requests: "{10*rate}" + +aiperf_benchmarks: + mooncake-trace-2k: + timeout_seconds: 7200 + dataset_url: https://raw.githubusercontent.com/kvcache-ai/Mooncake/refs/heads/main/FAST25-release/traces/conversation_trace.jsonl + dataset_type: mooncake_trace + dataset_cap: 2000 + endpoint_type: chat + endpoint_path: /v1/chat/completions + streaming: true + fixed_schedule: true + fixed_schedule_auto_offset: true + synthesis_max_isl: 131072 + + mooncake-trace-full: + timeout_seconds: 14400 + dataset_url: https://raw.githubusercontent.com/kvcache-ai/Mooncake/refs/heads/main/FAST25-release/traces/conversation_trace.jsonl + dataset_type: mooncake_trace + dataset_cap: null + endpoint_type: chat + endpoint_path: /v1/chat/completions + streaming: true + fixed_schedule: true + fixed_schedule_auto_offset: true + synthesis_max_isl: 131072 + + mooncake-toolagent: + timeout_seconds: 7200 + dataset_url: https://raw.githubusercontent.com/kvcache-ai/Mooncake/refs/heads/main/FAST25-release/traces/tool_agent_trace.jsonl + dataset_type: mooncake_trace + dataset_cap: 2000 + endpoint_type: chat + endpoint_path: /v1/chat/completions + streaming: true + fixed_schedule: true + fixed_schedule_auto_offset: true + synthesis_max_isl: 131072 diff --git a/projects/dynamo/orchestration/manifests/dynamo-graph-deployment.yaml b/projects/dynamo/orchestration/manifests/dynamo-graph-deployment.yaml new file mode 100644 index 000000000..5ed492a26 --- /dev/null +++ b/projects/dynamo/orchestration/manifests/dynamo-graph-deployment.yaml @@ -0,0 +1,14 @@ +apiVersion: nvidia.com/v1alpha1 +kind: DynamoGraphDeployment +metadata: + name: "" + namespace: "" + labels: + app.kubernetes.io/managed-by: forge + forge.openshift.io/project: dynamo +spec: + backendFramework: vllm + pvcs: + - name: model-cache + create: false + services: {} diff --git a/projects/dynamo/orchestration/preflight_phase.py b/projects/dynamo/orchestration/preflight_phase.py new file mode 100644 index 000000000..b7014b7c1 --- /dev/null +++ b/projects/dynamo/orchestration/preflight_phase.py @@ -0,0 +1,36 @@ +from __future__ import annotations + +import logging + +from projects.core.dsl.utils.k8s import oc_resource_exists +from projects.dynamo.orchestration import runtime_config + +logger = logging.getLogger(__name__) + + +def run() -> int: + """Validate that required Dynamo CRDs exist in the cluster.""" + logger.info("Starting Dynamo preflight checks") + + dynamo_config = runtime_config.get_dynamo_config() + required_crds = dynamo_config["required_crds"] + + missing_crds = [] + + for crd_name in required_crds: + logger.info(f"Checking for CRD: {crd_name}") + if not oc_resource_exists("crd", crd_name): + missing_crds.append(crd_name) + logger.error(f"Required CRD not found: {crd_name}") + else: + logger.info(f"CRD found: {crd_name}") + + if missing_crds: + logger.error( + f"Preflight check failed - missing {len(missing_crds)} required CRDs: " + f"{', '.join(missing_crds)}" + ) + return 1 + + logger.info("Preflight checks completed successfully - all required Dynamo CRDs are available") + return 0 diff --git a/projects/dynamo/orchestration/prepare_phase.py b/projects/dynamo/orchestration/prepare_phase.py new file mode 100644 index 000000000..6b45f93ac --- /dev/null +++ b/projects/dynamo/orchestration/prepare_phase.py @@ -0,0 +1,204 @@ +from __future__ import annotations + +import json +import logging +from typing import Any + +from projects.cluster.toolbox.cluster_deploy_operator import main as cluster_deploy_operator +from projects.cluster.toolbox.wait_for_crds import main as wait_for_crds_command +from projects.core.dsl.utils.k8s import oc, oc_get_json +from projects.core.library import env +from projects.core.orchestration.utils.k8s import ensure_namespace +from projects.dynamo.orchestration import runtime_config +from projects.dynamo.orchestration.cleanup_phase import run as cleanup_toolbox_run +from projects.dynamo.toolbox.capture_dynamo_state.main import ( + run as capture_dynamo_state_toolbox_run, +) +from projects.dynamo.toolbox.deploy_dynamo_platform.main import ( + run as deploy_dynamo_platform_toolbox_run, +) +from projects.gpu_operator.toolbox.bootstrap_gpu_clusterpolicy import ( + main as bootstrap_gpu_clusterpolicy, +) +from projects.gpu_operator.toolbox.bootstrap_nfd_instance import main as bootstrap_nfd_instance +from projects.kserve.toolbox.prepare_hf_model_cache.main import ( + run as prepare_hf_model_cache_toolbox_run, +) + +logger = logging.getLogger(__name__) + + +def operator_spec_by_package(platform: dict[str, Any], package: str) -> dict[str, Any]: + operators = platform["operators"] + if isinstance(operators, dict): + if package in operators: + return {"package": package, **operators[package]} + raise KeyError(f"Unknown operator package in dynamo platform config: {package}") + + for operator_spec in operators: + if operator_spec["package"] == package: + return operator_spec + raise KeyError(f"Unknown operator package in dynamo platform config: {package}") + + +def verify_oc_access() -> None: + oc("whoami") + + +def verify_cluster_version() -> None: + platform = runtime_config.get_platform_config() + version_info = oc("version", "-o", "json") + payload = json.loads(version_info.stdout) + + openshift_version = ( + payload.get("openshiftVersion") + or payload.get("releaseClientVersion") + or payload.get("clientVersion", {}).get("gitVersion") + or payload.get("serverVersion", {}).get("gitVersion") + or payload.get("serverVersion", {}).get("platform") + ) + if not openshift_version: + raise RuntimeError("Could not determine OpenShift version from `oc version -o json`") + + minimum = platform["cluster"]["minimum_openshift_version"] + if runtime_config.version_tuple(openshift_version) < runtime_config.version_tuple(minimum): + raise RuntimeError( + f"Cluster version {openshift_version} is older than the dynamo minimum {minimum}" + ) + + +def ensure_operator_subscription(operator_spec: dict[str, str]) -> dict[str, object]: + return cluster_deploy_operator.run( + package_name=operator_spec["package"], + target_namespace=operator_spec["namespace"], + source_name=operator_spec["source"], + channel=operator_spec["channel"], + source_namespace=operator_spec.get("source_namespace", "openshift-marketplace"), + display_name=operator_spec.get("display_name", operator_spec["package"]), + artifact_dirname_suffix=f"_{operator_spec['package']}", + ) + + +def prepare_nfd() -> None: + platform = runtime_config.get_platform_config() + operator_spec = operator_spec_by_package(platform, "nfd") + ensure_operator_subscription(operator_spec) + wait_for_crds_command.run( + crd_names=[operator_spec["bootstrap_crd"]], + display_name="NFD bootstrap CRD", + ) + bootstrap_nfd_instance.run() + + +def prepare_gpu_operator() -> None: + platform = runtime_config.get_platform_config() + operator_spec = operator_spec_by_package(platform, "gpu-operator-certified") + ensure_operator_subscription(operator_spec) + wait_for_crds_command.run( + crd_names=[operator_spec["bootstrap_crd"]], + display_name="GPU Operator bootstrap CRD", + ) + bootstrap_gpu_clusterpolicy.run() + + +def deploy_dynamo_platform() -> None: + """Deploy the Dynamo platform via Helm (operator + etcd + NATS).""" + dynamo_config = runtime_config.get_dynamo_config() + helm_config = dynamo_config["helm"] + + deploy_dynamo_platform_toolbox_run( + chart_repo=helm_config.get("chart_repo"), + chart_name=helm_config["chart_name"], + chart_version=helm_config["chart_version"], + chart_path=helm_config.get("chart_path"), + release_name=helm_config["release_name"], + namespace=helm_config["namespace"], + values_override=helm_config.get("values_override", {}), + wait_timeout_seconds=dynamo_config["operator"]["wait_timeout_seconds"], + ) + + +def wait_for_dynamo_crds() -> None: + """Wait for Dynamo CRDs to be registered after Helm install.""" + dynamo_config = runtime_config.get_dynamo_config() + wait_for_crds_command.run( + crd_names=dynamo_config["required_crds"], + display_name="Dynamo CRDs", + ) + + +def ensure_test_namespace() -> None: + namespace = runtime_config.get_namespace() + ensure_namespace( + namespace, + labels={ + "app.kubernetes.io/managed-by": "forge", + "forge.openshift.io/project": "dynamo", + }, + ) + + +def cleanup_previous_run() -> None: + namespace = runtime_config.get_namespace() + cleanup_toolbox_run(namespace=namespace) + + +def prepare_model_cache() -> None: + model_cache = runtime_config.get_model_cache_config() + + if not model_cache.get("enabled", False): + logger.info("Model cache disabled") + return + + model_uri = runtime_config.get_model_uri() + + if model_uri.startswith(("pvc://", "pvc+hf://")): + logger.info("Skipping cache for PVC-based model: %s", model_uri) + return + + namespace = runtime_config.get_namespace() + model_slug = runtime_config.get_model_slug() + + common_args = { + "namespace": namespace, + "namespace_is_managed": runtime_config.get_namespace_is_managed(), + "model_key": model_slug, + "model_uri": model_uri, + "pvc_size": model_cache["pvc"]["size"], + "access_mode": model_cache["pvc"]["access_mode"], + "storage_class_name": model_cache["pvc"].get("storage_class_name"), + "pvc_name_prefix": model_cache["pvc"]["name_prefix"], + "model_directory_name": model_cache["pvc"]["model_directory_name"], + } + + if model_uri.startswith("hf://"): + prepare_hf_model_cache_toolbox_run( + **common_args, + downloader_image=model_cache["hf"]["downloader_image"], + hf_token_file_path=None, + ) + else: + raise ValueError(f"Unsupported model URI scheme: {model_uri}") + + +def verify_gpu_nodes() -> None: + platform = runtime_config.get_platform_config() + selector = platform["cluster"]["gpu_node_label_selector"] + data = oc_get_json("nodes", selector=selector, ignore_not_found=True) + items = data.get("items", []) if data else [] + if not items: + raise RuntimeError( + f"No GPU nodes found with selector {selector}. Dynamo requires GPUs." + ) + + +def capture_prepare_state() -> None: + artifact_dir = env.ARTIFACT_DIR + namespace = runtime_config.get_namespace() + dynamo_config = runtime_config.get_dynamo_config() + + capture_dynamo_state_toolbox_run( + artifact_dir=artifact_dir, + namespace=namespace, + dynamo_namespace=dynamo_config["helm"]["namespace"], + ) diff --git a/projects/dynamo/orchestration/prepare_sequence.py b/projects/dynamo/orchestration/prepare_sequence.py new file mode 100644 index 000000000..c011bfc43 --- /dev/null +++ b/projects/dynamo/orchestration/prepare_sequence.py @@ -0,0 +1,24 @@ +from __future__ import annotations + +from projects.core.library import env +from projects.dynamo.orchestration import prepare_phase, runtime_config + + +def run_prepare_sequence() -> int: + prepare_phase.verify_oc_access() + prepare_phase.verify_cluster_version() + prepare_phase.prepare_nfd() + prepare_phase.prepare_gpu_operator() + prepare_phase.deploy_dynamo_platform() + prepare_phase.wait_for_dynamo_crds() + + for run_spec in runtime_config.get_run_specs(): + with runtime_config.activate_run_spec(run_spec): + with env.NextArtifactDir(f"prepare_{run_spec.artifact_dirname}"): + prepare_phase.ensure_test_namespace() + prepare_phase.cleanup_previous_run() + prepare_phase.prepare_model_cache() + prepare_phase.verify_gpu_nodes() + prepare_phase.capture_prepare_state() + + return 0 diff --git a/projects/dynamo/orchestration/presets.d/presets.yaml b/projects/dynamo/orchestration/presets.d/presets.yaml new file mode 100644 index 000000000..1462db06f --- /dev/null +++ b/projects/dynamo/orchestration/presets.d/presets.yaml @@ -0,0 +1,60 @@ +__multiple: true + +smoke: + runtime: + model_name: Qwen/Qwen3-0.6B + deployment_profile: aggregated + benchmark_tool: guidellm + benchmark_key: short + +disagg-smoke: + runtime: + model_name: Qwen/Qwen3-0.6B + deployment_profile: disaggregated + benchmark_tool: guidellm + benchmark_key: short + +aiperf-mooncake: + runtime: + model_name: Qwen/Qwen3-32B + deployment_profile: aggregated + benchmark_tool: aiperf + benchmark_key: mooncake-trace-2k + +aiperf-mooncake-toolagent: + runtime: + model_name: Qwen/Qwen3-32B + deployment_profile: aggregated + benchmark_tool: aiperf + benchmark_key: mooncake-toolagent + +qwen32b-4r-mooncake: + runtime: + model_name: Qwen/Qwen3-32B-FP8 + deployment_profile: aggregated-rr-4r-tp2 + benchmark_tool: aiperf + benchmark_key: mooncake-trace-2k + +qwen32b-4r-mooncake-100: + runtime: + model_name: Qwen/Qwen3-32B-FP8 + deployment_profile: aggregated-rr-4r-tp2 + benchmark_tool: aiperf + benchmark_key: mooncake-trace-2k + workloads: + aiperf_benchmarks: + mooncake-trace-2k: + dataset_cap: 100 + +full: + runtime: + model_name: + - meta-llama/Llama-3-70B + - Qwen/Qwen3-32B + deployment_profile: + - aggregated + - disaggregated + benchmark_tool: guidellm + benchmark_key: + - concurrent-1k-1k + - multi-turn diff --git a/projects/dynamo/orchestration/render_graph_deployment.py b/projects/dynamo/orchestration/render_graph_deployment.py new file mode 100644 index 000000000..f97a3ba52 --- /dev/null +++ b/projects/dynamo/orchestration/render_graph_deployment.py @@ -0,0 +1,217 @@ +from __future__ import annotations + +import copy +from pathlib import Path +from typing import Any + +import yaml + + +def _load_yaml(path: Path) -> Any: + with path.open(encoding="utf-8") as handle: + return yaml.safe_load(handle) + + +def render_graph_deployment( + *, + config_dir: str | Path, + namespace: str, + model_name: str, + model_slug: str, + deployment_profile: dict[str, Any], + model_cache: dict[str, Any], + dynamo_config: dict[str, Any], +) -> dict[str, Any]: + """Render a DynamoGraphDeployment manifest from config inputs.""" + template_path = Path(config_dir) / "manifests" / "dynamo-graph-deployment.yaml" + manifest = _load_yaml(template_path) + + manifest["metadata"]["name"] = f"dynamo-{model_slug}" + manifest["metadata"]["namespace"] = namespace + manifest["metadata"].setdefault("labels", {}) + manifest["metadata"]["labels"].update({ + "app.kubernetes.io/managed-by": "forge", + "forge.openshift.io/project": "dynamo", + }) + + p = deployment_profile + backend = p.get("backend_framework", "vllm") + manifest["spec"]["backendFramework"] = backend + + pvc_name = model_cache.get("pvc_name", "model-cache") + manifest["spec"]["pvcs"] = [{"name": pvc_name, "create": False}] + + vllm_args = _build_vllm_args(p.get("vllm_args", [])) + if p.get("kv_transfer_config"): + vllm_args.extend(["--kv-transfer-config", f"'{p['kv_transfer_config']}'"]) + + hf_home = p.get("hf_home", "/opt/models") + base_env = [ + {"name": "SERVED_MODEL_NAME", "value": model_name}, + {"name": "MODEL_PATH", "value": model_name}, + {"name": "HF_HOME", "value": hf_home}, + ] + if p.get("kvbm_cpu_cache_gb"): + base_env.append({"name": "DYN_KVBM_CPU_CACHE_GB", "value": str(p["kvbm_cpu_cache_gb"])}) + for key, value in p.get("env", {}).items(): + base_env.append({"name": key, "value": str(value)}) + + ctx = _BuildCtx( + profile=p, backend=backend, vllm_args=vllm_args, base_env=base_env, + hf_home=hf_home, model_name=model_name, pvc_name=pvc_name, + runtime_image=p["runtime_image"], frontend_image=p["frontend_image"], + tp=p.get("tensor_parallelism", 1), + kv_block_size=p.get("kv_block_size", "16"), + router_mode=p.get("router_mode", "direct"), + use_dra=p.get("use_dra_resources", False), + dra_name=p.get("dra_resource_name", "dra.llm-d.io/gpu-nic-pair"), + ) + + serving_mode = p.get("serving_mode", "aggregated") + services = _build_routing_service(ctx) + if serving_mode == "disaggregated": + services["VllmPrefillWorker"] = _build_worker(ctx, mode="prefill", + replicas=p.get("prefill_replicas", 1)) + services["VllmDecodeWorker"] = _build_worker(ctx, mode="decode", + replicas=p.get("decode_replicas", 2)) + else: + services["VllmWorker"] = _build_worker(ctx, mode="aggregated", + replicas=p.get("replicas", 1)) + + manifest["spec"]["services"] = services + return manifest + + +class _BuildCtx: + """Carries resolved config through the builder functions.""" + __slots__ = ( + "profile", "backend", "vllm_args", "base_env", "hf_home", "model_name", + "pvc_name", "runtime_image", "frontend_image", "tp", "kv_block_size", + "router_mode", "use_dra", "dra_name", + ) + + def __init__(self, **kwargs): + for k, v in kwargs.items(): + setattr(self, k, v) + + +def _build_worker(ctx: _BuildCtx, *, mode: str, replicas: int) -> dict[str, Any]: + disagg_flag = "" + if mode in ("prefill", "decode"): + disagg_flag = f" --disaggregation-mode {mode}" + + worker_cmd = ( + f"python3 -m dynamo.{ctx.backend} " + f"--model $MODEL_PATH --served-model-name $SERVED_MODEL_NAME " + f"--tensor-parallel-size {ctx.tp} --data-parallel-size 1 " + + " ".join(ctx.vllm_args) + + " --kv-events-config '{\"enable_kv_cache_events\":true}'" + + f" --block-size {ctx.kv_block_size}" + + disagg_flag + ) + + sidecar_mode = "direct" if ctx.router_mode == "kv" else ctx.router_mode + + return { + "componentType": "worker", + "volumeMounts": [{"name": ctx.pvc_name, "mountPoint": ctx.hf_home}], + "sharedMemory": {"size": "2Gi"}, + "frontendSidecar": { + "image": ctx.runtime_image, + "args": ["-m", "dynamo.frontend", "--router-mode", sidecar_mode], + }, + "extraPodSpec": { + "mainContainer": { + "env": copy.deepcopy(ctx.base_env), + "args": [worker_cmd], + "command": ["/bin/sh", "-c"], + "image": ctx.runtime_image, + "workingDir": f"/workspace/examples/backends/{ctx.backend}", + }, + }, + "replicas": replicas, + "resources": _build_gpu_resources(ctx.tp, use_dra=ctx.use_dra, dra_name=ctx.dra_name), + } + + +def _build_routing_service(ctx: _BuildCtx) -> dict[str, Any]: + if ctx.router_mode == "kv": + router_args = ctx.profile.get("router_args", []) + args = ["--router-mode", "kv", "--router-reset-states"] + router_args + cpu = ctx.profile.get("frontend_cpu", "4") + return { + "Frontend": { + "componentType": "frontend", + "replicas": 1, + "extraPodSpec": { + "mainContainer": { + "image": ctx.runtime_image, + "workingDir": "/workspace", + "command": ["python3", "-m", "dynamo.frontend"], + "args": args, + "env": [{"name": "HF_HOME", "value": ctx.hf_home}], + }, + }, + "resources": {"requests": {"cpu": cpu}, "limits": {"cpu": cpu}}, + }, + } + return {"Epp": _build_epp_service(ctx)} + + +def _build_epp_service(ctx: _BuildCtx) -> dict[str, Any]: + return { + "componentType": "epp", + "replicas": 1, + "extraPodSpec": { + "mainContainer": { + "image": ctx.frontend_image, + "env": [ + {"name": "DYN_KV_CACHE_BLOCK_SIZE", "value": ctx.kv_block_size}, + {"name": "DYN_MODEL_NAME", "value": ctx.model_name}, + {"name": "DYN_DECODE_FALLBACK", "value": "true"}, + ], + }, + }, + "eppConfig": { + "config": { + "plugins": [ + {"type": "disagg-profile-handler"}, + {"name": "decode-filter", "type": "label-filter", "parameters": { + "label": "nvidia.com/dynamo-sub-component-type", + "validValues": ["decode"], "allowsNoLabel": True, + }}, + {"name": "picker", "type": "max-score-picker"}, + {"name": "dyn-decode", "type": "dyn-decode-scorer"}, + ], + "schedulingProfiles": [{ + "name": "decode", + "plugins": [ + {"pluginRef": "decode-filter", "weight": 1}, + {"pluginRef": "dyn-decode", "weight": 1}, + {"pluginRef": "picker", "weight": 1}, + ], + }], + }, + }, + } + + +def _build_gpu_resources(tp: int, *, use_dra: bool, dra_name: str) -> dict[str, Any]: + v = str(tp) + if use_dra: + return {"limits": {"custom": {dra_name: v}}, "requests": {"custom": {dra_name: v}}} + return {"limits": {"gpu": v}, "requests": {"gpu": v}} + + +def _build_vllm_args(vllm_args: dict[str, Any] | list[str]) -> list[str]: + if isinstance(vllm_args, list): + return [str(arg) for arg in vllm_args] + rendered = [] + for key, value in vllm_args.items(): + cli_key = key.replace("_", "-") + if isinstance(value, bool): + if value: + rendered.append(f"--{cli_key}") + continue + rendered.append(f"--{cli_key}={value}") + return rendered diff --git a/projects/dynamo/orchestration/runtime_config.py b/projects/dynamo/orchestration/runtime_config.py new file mode 100644 index 000000000..64476d923 --- /dev/null +++ b/projects/dynamo/orchestration/runtime_config.py @@ -0,0 +1,353 @@ +from __future__ import annotations + +import ast +import copy +import logging +import re +from contextlib import contextmanager +from dataclasses import dataclass +from itertools import product +from pathlib import Path +from typing import Any + +import yaml + +from projects.core.dsl.utils import slugify_identifier, truncate_k8s_name +from projects.core.library import config, env, run + +logger = logging.getLogger(__name__) +RUNTIME_DIR = Path(__file__).resolve().parent +PROJECT_DIR = RUNTIME_DIR.parent +ORCHESTRATION_DIR = PROJECT_DIR / "orchestration" +CONFIG_DIR = ORCHESTRATION_DIR + + +@dataclass(frozen=True) +class RunSpec: + model_name: str + model_slug: str + deployment_profile_name: str + deployment_profile_slug: str + namespace: str + artifact_dirname: str + namespace_is_managed: bool + + +def init() -> Path: + if not logging.getLogger().handlers: + logging.basicConfig(level=logging.INFO, format="%(levelname)s: %(message)s") + + env.init() + run.init() + if config.project is None: + config.init(CONFIG_DIR) + ensure_artifact_directories(env.ARTIFACT_DIR) + return env.ARTIFACT_DIR + + +def ensure_artifact_directories(artifact_dir: Path) -> None: + for relative in ("src", "artifacts", "artifacts/results"): + (artifact_dir / relative).mkdir(parents=True, exist_ok=True) + + +def _get_runtime_value(key: str) -> Any: + return config.project.get_config(f"runtime.{key}", None) + + +def _normalize_string_or_list(value: Any, field_name: str) -> list[str]: + if value in (None, ""): + return [] + + if isinstance(value, list): + return [str(item).strip() for item in value if str(item).strip()] + + if not isinstance(value, str): + raise ValueError(f"{field_name} must be a string or list of strings, got {type(value)}") + + value = value.strip() + if value.startswith("[") and value.endswith("]"): + try: + parsed = ast.literal_eval(value) + if isinstance(parsed, list): + return [str(item).strip() for item in parsed if str(item).strip()] + except (ValueError, SyntaxError): + items = [item.strip() for item in value[1:-1].split(",")] + return [item for item in items if item] + + return [value] if value else [] + + +def _deep_merge(base: Any, override: Any) -> Any: + if not isinstance(base, dict) or not isinstance(override, dict): + return copy.deepcopy(override) + + merged = copy.deepcopy(base) + for key, value in override.items(): + if key in merged and isinstance(merged[key], dict) and isinstance(value, dict): + merged[key] = _deep_merge(merged[key], value) + else: + merged[key] = copy.deepcopy(value) + return merged + + +def _load_yaml(path: Path) -> Any: + with path.open(encoding="utf-8") as handle: + return yaml.safe_load(handle) + + +def get_config_dir() -> Path: + return ORCHESTRATION_DIR + + +def get_job_name() -> str: + job_name = _get_runtime_value("job_name") + if job_name: + return job_name + + preset_name = config.project.get_config("runtime.selected_preset") + return f"local-{preset_name}" + + +def get_platform_config() -> dict[str, Any]: + return normalize_platform_config(copy.deepcopy(config.project.get_config("platform"))) + + +def get_dynamo_config() -> dict[str, Any]: + """Get the Dynamo-specific platform configuration (Helm, operator, infra).""" + platform = get_platform_config() + return copy.deepcopy(platform["dynamo"]) + + +def _derive_base_namespace() -> str: + platform_data = get_platform_config() + namespace_override = _get_runtime_value("namespace_override") + namespace_config = platform_data["cluster"]["namespace"] + default_namespace = namespace_config.get("name") + + if namespace_override: + return namespace_override + if default_namespace: + return default_namespace + + return derive_namespace( + get_job_name(), + namespace_config["prefix"], + namespace_config["max_length"], + ) + + +def get_namespace() -> str: + return _derive_base_namespace() + + +_namespace_is_managed_override: bool | None = None + + +def get_namespace_is_managed() -> bool: + if _namespace_is_managed_override is not None: + return _namespace_is_managed_override + + namespace_override = _get_runtime_value("namespace_override") + platform_data = get_platform_config() + default_namespace = platform_data["cluster"]["namespace"].get("name") + + return namespace_override is None and default_namespace is None + + +def get_model_name() -> str: + model_names = _normalize_string_or_list(_get_runtime_value("model_name"), "runtime.model_name") + if len(model_names) != 1: + raise ValueError( + f"Expected exactly one runtime.model_name in the active dynamo run, got {model_names}" + ) + return model_names[0] + + +def get_model_slug(model_name: str | None = None) -> str: + return slugify_identifier(model_name or get_model_name(), max_length=32) + + +def get_model_uri(model_name: str | None = None) -> str: + name = model_name or get_model_name() + if name.startswith(("hf://", "oci://", "pvc://", "pvc+hf://")): + return name + return f"hf://{name}" + + +def get_served_model_name(model_name: str | None = None) -> str: + return model_name or get_model_name() + + +def get_model_cache_config() -> dict[str, Any]: + return copy.deepcopy(config.project.get_config("model_cache")) + + +def get_benchmark_tool() -> str | None: + return _get_runtime_value("benchmark_tool") + + +def get_benchmark_key() -> str | None: + return _get_runtime_value("benchmark_key") + + +def get_benchmark_keys() -> list[str]: + return _normalize_string_or_list(_get_runtime_value("benchmark_key"), "runtime.benchmark_key") + + +def get_benchmark_job_names() -> list[str]: + """Job names to clean up, based on benchmark_tool + benchmark_key.""" + tool = get_benchmark_tool() + key = get_benchmark_key() + if not tool or not key: + return [] + + if tool == "guidellm": + bench = config.project.get_config(f"workloads.guidellm_benchmarks.{key}", {}) + name = bench.get("job_name", "guidellm-benchmark") + return [name] if name else [] + elif tool == "aiperf": + from projects.core.dsl.utils import slugify_identifier + return [f"aiperf-{slugify_identifier(key, max_length=32)}"] + return [] + + +def get_deployment_profile_name() -> str: + deployment_profiles = _normalize_string_or_list( + _get_runtime_value("deployment_profile"), + "runtime.deployment_profile", + ) + if len(deployment_profiles) != 1: + raise ValueError( + "Expected exactly one runtime.deployment_profile in the active dynamo run, " + f"got {deployment_profiles}" + ) + return deployment_profiles[0] + + +def get_deployment_profile() -> dict[str, Any]: + profile_name = get_deployment_profile_name() + deployment_config = copy.deepcopy(config.project.get_config("deployments")) + profile = copy.deepcopy(deployment_config[profile_name]) + defaults = copy.deepcopy(deployment_config.get("defaults", {})) + + return _deep_merge(defaults, profile) + + +def get_smoke_request() -> dict[str, Any]: + smoke_request_key = config.project.get_config("runtime.smoke_request_key") + return copy.deepcopy(config.project.get_config(f"workloads.smoke_requests.{smoke_request_key}")) + + +def get_run_specs() -> list[RunSpec]: + model_names = _normalize_string_or_list(_get_runtime_value("model_name"), "runtime.model_name") + profile_names = _normalize_string_or_list( + _get_runtime_value("deployment_profile"), + "runtime.deployment_profile", + ) + + if not model_names: + raise ValueError("runtime.model_name must be set to a model name or list of model names") + if not profile_names: + raise ValueError( + "runtime.deployment_profile must be set to a deployment profile or list of profiles" + ) + + combinations = list(product(model_names, profile_names)) + base_namespace = _derive_base_namespace() + namespace_max_length = get_platform_config()["cluster"]["namespace"]["max_length"] + namespace_is_managed = get_namespace_is_managed() + run_specs: list[RunSpec] = [] + + for model_name, profile_name in combinations: + model_slug = get_model_slug(model_name) + profile_slug = slugify_identifier(profile_name, max_length=24) + if len(combinations) == 1: + namespace = base_namespace + artifact_dirname = "dynamo_run" + else: + namespace = truncate_k8s_name( + f"{base_namespace}-{model_slug}-{profile_slug}", + max_length=namespace_max_length, + ) + artifact_dirname = f"dynamo_run_{profile_slug}_{model_slug}" + + run_specs.append( + RunSpec( + model_name=model_name, + model_slug=model_slug, + deployment_profile_name=profile_name, + deployment_profile_slug=profile_slug, + namespace=namespace, + artifact_dirname=artifact_dirname, + namespace_is_managed=namespace_is_managed, + ) + ) + + return run_specs + + +@contextmanager +def activate_run_spec(run_spec: RunSpec): + global _namespace_is_managed_override + + saved = { + "model_name": _get_runtime_value("model_name"), + "deployment_profile": _get_runtime_value("deployment_profile"), + "namespace_override": _get_runtime_value("namespace_override"), + } + prev_managed_override = _namespace_is_managed_override + + config.project.set_config("runtime.model_name", run_spec.model_name) + config.project.set_config("runtime.deployment_profile", run_spec.deployment_profile_name) + config.project.set_config("runtime.namespace_override", run_spec.namespace) + _namespace_is_managed_override = run_spec.namespace_is_managed + try: + yield + finally: + config.project.set_config("runtime.model_name", saved["model_name"]) + config.project.set_config("runtime.deployment_profile", saved["deployment_profile"]) + config.project.set_config("runtime.namespace_override", saved["namespace_override"]) + _namespace_is_managed_override = prev_managed_override + + +def normalize_platform_config(platform_data: dict[str, Any]) -> dict[str, Any]: + cluster = platform_data["cluster"] + if "namespace" not in cluster: + cluster["namespace"] = { + "name": cluster.pop("namespace_name", None), + "prefix": cluster.pop("namespace_prefix"), + "max_length": cluster.pop("namespace_max_length"), + } + + operators = platform_data["operators"] + if isinstance(operators, list): + platform_data["operators"] = { + operator_spec["package"]: { + key: value for key, value in operator_spec.items() if key != "package" + } + for operator_spec in operators + } + + return platform_data + + +def derive_namespace(job_name: str, prefix: str, max_length: int) -> str: + slug = re.sub(r"[^a-z0-9-]+", "-", job_name.lower()) + slug = re.sub(r"-{2,}", "-", slug).strip("-") + if not slug: + slug = "run" + + if slug.startswith(f"{prefix}-"): + namespace = slug + else: + namespace = f"{prefix}-{slug}" + + namespace = namespace[:max_length].rstrip("-") + if not namespace: + raise ValueError(f"Could not derive a valid namespace from job name: {job_name}") + return namespace + + +def version_tuple(value: str) -> tuple[int, ...]: + numbers = re.findall(r"\d+", value) + return tuple(int(number) for number in numbers[:3]) diff --git a/projects/dynamo/orchestration/test_phase.py b/projects/dynamo/orchestration/test_phase.py new file mode 100644 index 000000000..a4a2cb1bc --- /dev/null +++ b/projects/dynamo/orchestration/test_phase.py @@ -0,0 +1,447 @@ +from __future__ import annotations + +import logging +import sys +import time +from pathlib import Path +from typing import Any + +from projects.core.dsl import shell +from projects.core.dsl.utils import slugify_identifier +from projects.core.dsl.utils.k8s import oc, oc_get_json +from projects.core.library import env +from projects.core.library.postprocess import run_and_postprocess, write_test_labels +from projects.core.library.run import SignalInterrupt +from projects.core.orchestration.utils.k8s import ensure_namespace +from projects.dynamo.orchestration.prepare_phase import prepare_model_cache +from projects.dynamo.orchestration.render_graph_deployment import render_graph_deployment +from projects.dynamo.orchestration.utils import write_yaml +from projects.guidellm.toolbox.run_guidellm_benchmark import build_guidellm_args +from projects.guidellm.toolbox.run_guidellm_benchmark import main as run_guidellm_benchmark_command +from projects.guidellm.toolbox.run_smoke_request import main as run_smoke_request_command + +logger = logging.getLogger(__name__) + + +def create_test_labels() -> None: + from projects.dynamo.orchestration import runtime_config + + model_name = runtime_config.get_model_name() + deployment_profile = runtime_config.get_deployment_profile_name() + benchmark_tool = runtime_config.get_benchmark_tool() + benchmark_key = runtime_config.get_benchmark_key() + + labels = { + "model_name": model_name, + "deployment_profile": deployment_profile, + "framework": "dynamo", + } + + if benchmark_tool: + labels["benchmark_tool"] = benchmark_tool + if benchmark_key: + labels["benchmark_key"] = benchmark_key + + write_test_labels(env.ARTIFACT_DIR, labels) + logger.info("Created test labels: %s", labels) + + +def run() -> int: + return run_and_postprocess(do_test) + + +def run_finalizers( + endpoint_url: str | None, + primary_exc: tuple[type[BaseException], BaseException, Any] | None, + finalizer_exc: tuple[type[BaseException], BaseException, Any] | None, +) -> tuple[type[BaseException], BaseException, Any] | None: + def _run_finalizer(description: str, callback, **kwargs): + try: + callback(**kwargs) + except Exception: + if primary_exc is None: + logger.exception("Finalizer failed: %s", description) + return finalizer_exc or sys.exc_info() + logger.exception("Ignoring %s failure after primary test failure", description) + return finalizer_exc + + from projects.dynamo.orchestration import runtime_config + + namespace = runtime_config.get_namespace() + platform = runtime_config.get_platform_config() + capture_namespace_events = platform["artifacts"]["capture_namespace_events"] + + finalizer_exc = _run_finalizer( + "capture dynamo state", + capture_dynamo_state, + ) + finalizer_exc = _run_finalizer( + "write endpoint URL", + write_endpoint_url, + artifact_dir=env.ARTIFACT_DIR, + endpoint_url=endpoint_url, + ) + finalizer_exc = _run_finalizer( + "capture namespace events", + capture_namespace_events_after_test, + artifact_dir=env.ARTIFACT_DIR, + namespace=namespace, + capture_namespace_events=capture_namespace_events, + ) + finalizer_exc = _run_finalizer( + "cleanup runtime resources", + cleanup_test_resources, + ) + + return primary_exc, finalizer_exc + + +def do_test() -> int: + from projects.dynamo.orchestration import runtime_config + + namespace = runtime_config.get_namespace() + + ensure_namespace( + namespace, + labels={ + "app.kubernetes.io/managed-by": "forge", + "forge.openshift.io/project": "dynamo", + }, + ) + + endpoint_url: str | None = None + primary_exc: tuple[type[BaseException], BaseException, Any] | None = None + finalizer_exc: tuple[type[BaseException], BaseException, Any] | None = None + + with env.NextArtifactDir("dynamo_test"): + try: + create_test_labels() + + endpoint_url = deploy_graph_deployment() + + if not endpoint_url: + raise ValueError("Failed to discover endpoint URL from DynamoGraphDeployment") + run_smoke_request(endpoint_url=endpoint_url) + + run_benchmark(endpoint_url=endpoint_url) + except Exception: + primary_exc = sys.exc_info() + except SignalInterrupt: + primary_exc = sys.exc_info() + finally: + do_finalizers = True + if primary_exc and isinstance(primary_exc[1], SignalInterrupt): + logging.warning("Caught a SignalInterrupt, skipping the finalizers") + do_finalizers = False + + if do_finalizers: + primary_exc, finalizer_exc = run_finalizers( + endpoint_url, primary_exc, finalizer_exc + ) + + if primary_exc is not None: + raise primary_exc[1].with_traceback(primary_exc[2]) + + if finalizer_exc is not None: + raise finalizer_exc[1].with_traceback(finalizer_exc[2]) + + return 0 + + +def deploy_graph_deployment() -> str: + """Deploy DynamoGraphDeployment and return the endpoint URL.""" + logger.info("Starting DynamoGraphDeployment deployment") + + from projects.dynamo.orchestration import runtime_config + + namespace = runtime_config.get_namespace() + + _prepare_model_cache() + + manifest_path = _build_graph_deployment_manifest() + + logger.info("Applying DynamoGraphDeployment from manifest: %s", manifest_path) + oc("apply", "-f", str(manifest_path), "-n", namespace, check=True) + + _wait_for_dynamo_ready(namespace=namespace) + + endpoint_url = _discover_endpoint(namespace=namespace) + + logger.info("DynamoGraphDeployment deployed, endpoint: %s", endpoint_url) + return endpoint_url + + +def _prepare_model_cache() -> None: + from projects.dynamo.orchestration import runtime_config + + model_name = runtime_config.get_model_name() + logger.info("Preparing model cache for model: %s", model_name) + prepare_model_cache() + + +def _build_graph_deployment_manifest() -> Path: + from projects.dynamo.orchestration import runtime_config + + config_dir = runtime_config.get_config_dir() + namespace = runtime_config.get_namespace() + model_name = runtime_config.get_model_name() + model_slug = runtime_config.get_model_slug(model_name) + deployment_profile = runtime_config.get_deployment_profile() + model_cache = runtime_config.get_model_cache_config() + dynamo_config = runtime_config.get_dynamo_config() + + manifest = render_graph_deployment( + config_dir=config_dir, + namespace=namespace, + model_name=model_name, + model_slug=model_slug, + deployment_profile=deployment_profile, + model_cache=model_cache, + dynamo_config=dynamo_config, + ) + + artifacts_dir = env.ARTIFACT_DIR / "artifacts" + artifacts_dir.mkdir(parents=True, exist_ok=True) + manifest_path = artifacts_dir / "dynamo-graph-deployment.yaml" + write_yaml(manifest_path, manifest) + + logger.info("Built DynamoGraphDeployment manifest: %s", manifest_path) + return manifest_path + + +def _wait_for_dynamo_ready( + *, + namespace: str, + timeout_seconds: int = 600, + poll_interval: int = 10, +) -> None: + """Wait for all Dynamo pods to be Running.""" + logger.info("Waiting for Dynamo pods to be ready in %s (timeout=%ds)", namespace, timeout_seconds) + deadline = time.time() + timeout_seconds + + while time.time() < deadline: + pods_data = oc_get_json("pods", namespace=namespace, ignore_not_found=True) + if not pods_data: + time.sleep(poll_interval) + continue + + items = pods_data.get("items", []) + if not items: + logger.info("No pods found yet, waiting...") + time.sleep(poll_interval) + continue + + all_ready = True + for pod in items: + phase = pod.get("status", {}).get("phase", "Unknown") + name = pod.get("metadata", {}).get("name", "unknown") + if phase not in ("Running", "Succeeded"): + logger.info("Pod %s is %s", name, phase) + all_ready = False + + if all_ready and items: + logger.info("All %d Dynamo pods are ready", len(items)) + return + + time.sleep(poll_interval) + + raise TimeoutError( + f"Dynamo pods did not become ready within {timeout_seconds}s in namespace {namespace}" + ) + + +def _discover_endpoint(*, namespace: str) -> str: + """Discover the inference endpoint URL from the Dynamo deployment. + + Checks in order: + 1. Frontend service (standalone router, port 8000) + 2. Gateway service (inference-gateway, port 80) + 3. Any service on port 8000 + """ + # Check for standalone frontend service first + svc_data = oc_get_json("services", namespace=namespace, ignore_not_found=True) + if svc_data: + items = svc_data.get("items", []) + # Prefer frontend service + for svc in items: + svc_name = svc.get("metadata", {}).get("name", "") + if "frontend" in svc_name: + for port in svc.get("spec", {}).get("ports", []): + if port.get("port") == 8000: + return f"http://{svc_name}.{namespace}.svc.cluster.local:8000" + + # Then gateway + for svc in items: + svc_name = svc.get("metadata", {}).get("name", "") + if "inference-gateway" in svc_name: + for port in svc.get("spec", {}).get("ports", []): + if port.get("port") in (80, 8080): + return f"http://{svc_name}.{namespace}.svc.cluster.local:{port['port']}" + + # Fallback: any service on 8000 + for svc in items: + for port in svc.get("spec", {}).get("ports", []): + if port.get("port") == 8000: + svc_name = svc["metadata"]["name"] + return f"http://{svc_name}.{namespace}.svc.cluster.local:8000" + + raise RuntimeError(f"Could not discover Dynamo endpoint in namespace {namespace}") + + +def run_smoke_request(*, endpoint_url: str) -> dict[str, object]: + from projects.dynamo.orchestration import runtime_config + + namespace = runtime_config.get_namespace() + platform = runtime_config.get_platform_config() + smoke = platform["smoke"] + smoke_request = runtime_config.get_smoke_request() + + return run_smoke_request_command.run( + namespace=namespace, + endpoint_url=endpoint_url, + pod_name=smoke["pod_name"], + client_image=smoke["client_image"], + endpoint_path=smoke["endpoint_path"], + request_timeout_seconds=smoke["request_timeout_seconds"], + served_model_name=runtime_config.get_served_model_name(), + prompt=smoke_request["prompt"], + max_tokens=smoke_request["max_tokens"], + temperature=smoke_request["temperature"], + ) + + +def run_benchmark(*, endpoint_url: str) -> None: + """Dispatch to guidellm or aiperf based on runtime.benchmark_tool.""" + from projects.core.library import config + + tool = config.project.get_config("runtime.benchmark_tool", None) + key = config.project.get_config("runtime.benchmark_key", None) + + if not tool or not key: + logger.info("No benchmark configured (benchmark_tool=%s, benchmark_key=%s), skipping", tool, key) + return + + if tool == "guidellm": + bench_config = config.project.get_config(f"workloads.guidellm_benchmarks.{key}", None) + if bench_config is None: + raise ValueError( + f"benchmark_key '{key}' not found in workloads.guidellm_benchmarks. " + f"Available: {list(config.project.get_config('workloads.guidellm_benchmarks', {}).keys())}" + ) + _run_guidellm_benchmark(endpoint_url=endpoint_url) + elif tool == "aiperf": + bench_config = config.project.get_config(f"workloads.aiperf_benchmarks.{key}", None) + if bench_config is None: + raise ValueError( + f"benchmark_key '{key}' not found in workloads.aiperf_benchmarks. " + f"Available: {list(config.project.get_config('workloads.aiperf_benchmarks', {}).keys())}" + ) + _run_aiperf_benchmark(endpoint_url=endpoint_url) + else: + raise ValueError(f"Unknown benchmark_tool '{tool}'. Must be 'guidellm' or 'aiperf'.") + + +def _run_guidellm_benchmark(*, endpoint_url: str) -> None: + from projects.core.library import config + from projects.dynamo.orchestration import runtime_config + + namespace = runtime_config.get_namespace() + key = config.project.get_config("runtime.benchmark_key") + benchmark = config.project.get_config(f"workloads.guidellm_benchmarks.{key}") + + guidellm_args = build_guidellm_args(benchmark) + if not any(arg.startswith("--processor=") for arg in guidellm_args): + guidellm_args.append(f"--processor={runtime_config.get_model_name()}") + + artifact_name = f"benchmark_{slugify_identifier(key, max_length=48)}" + with env.NextArtifactDir(artifact_name): + run_guidellm_benchmark_command.run( + endpoint_url=endpoint_url, + name=benchmark.get("job_name", "guidellm-benchmark"), + namespace=namespace, + image=benchmark.get("image", "ghcr.io/vllm-project/guidellm:v0.5.4"), + timeout=benchmark.get("timeout_seconds", 3600), + pvc_size=benchmark.get("pvc_size", "1Gi"), + guidellm_args=guidellm_args, + ) + + +def _run_aiperf_benchmark(*, endpoint_url: str) -> None: + from projects.core.library import config + from projects.dynamo.orchestration import runtime_config + from projects.dynamo.toolbox.run_aiperf_benchmark.main import run as run_aiperf + + key = config.project.get_config("runtime.benchmark_key") + aiperf_config = config.project.get_config(f"workloads.aiperf_benchmarks.{key}") + model_name = runtime_config.get_model_name() + + artifact_name = f"aiperf_{slugify_identifier(key, max_length=48)}" + with env.NextArtifactDir(artifact_name): + run_aiperf( + endpoint_url=endpoint_url, + model_name=model_name, + name=f"aiperf-{slugify_identifier(key, max_length=32)}", + namespace=runtime_config.get_namespace(), + dataset_url=aiperf_config["dataset_url"], + dataset_type=aiperf_config["dataset_type"], + dataset_cap=aiperf_config.get("dataset_cap"), + endpoint_type=aiperf_config.get("endpoint_type", "chat"), + endpoint_path=aiperf_config.get("endpoint_path", "/v1/chat/completions"), + streaming=aiperf_config.get("streaming", True), + fixed_schedule=aiperf_config.get("fixed_schedule", True), + fixed_schedule_auto_offset=aiperf_config.get("fixed_schedule_auto_offset", True), + synthesis_max_isl=aiperf_config.get("synthesis_max_isl"), + ) + + +def capture_dynamo_state() -> None: + from projects.dynamo.orchestration import runtime_config + from projects.dynamo.toolbox.capture_dynamo_state.main import run as capture_state + + namespace = runtime_config.get_namespace() + dynamo_config = runtime_config.get_dynamo_config() + + capture_state( + artifact_dir=env.ARTIFACT_DIR, + namespace=namespace, + dynamo_namespace=dynamo_config["helm"]["namespace"], + ) + + +def write_endpoint_url(*, artifact_dir: Path, endpoint_url: str | None) -> None: + if not endpoint_url: + return + + endpoint_file = artifact_dir / "artifacts" / "endpoint.url" + endpoint_file.parent.mkdir(parents=True, exist_ok=True) + endpoint_file.write_text(f"{endpoint_url}\n", encoding="utf-8") + + +def cleanup_test_resources() -> None: + from projects.dynamo.orchestration import runtime_config + from projects.dynamo.toolbox.cleanup_dynamo_resources.main import run as cleanup + + namespace = runtime_config.get_namespace() + benchmark_job_names = runtime_config.get_benchmark_job_names() or [None] + + for benchmark_job_name in benchmark_job_names: + cleanup( + namespace=namespace, + benchmark_job_name=benchmark_job_name, + ) + + +def capture_namespace_events_after_test( + *, + artifact_dir: Path, + namespace: str, + capture_namespace_events: bool, +) -> None: + if not capture_namespace_events: + return + + shell.run( + f"oc get events -n {namespace} --sort-by=.metadata.creationTimestamp", + check=False, + stdout_dest=artifact_dir / "artifacts" / "namespace.events.txt", + ) diff --git a/projects/dynamo/orchestration/utils.py b/projects/dynamo/orchestration/utils.py new file mode 100644 index 000000000..fcd471253 --- /dev/null +++ b/projects/dynamo/orchestration/utils.py @@ -0,0 +1,11 @@ +from __future__ import annotations + +from pathlib import Path + +import yaml + + +def write_yaml(path: Path, data: dict) -> None: + path.parent.mkdir(parents=True, exist_ok=True) + with path.open("w", encoding="utf-8") as fh: + yaml.dump(data, fh, default_flow_style=False, sort_keys=False) diff --git a/projects/dynamo/tests/test_render_graph_deployment.py b/projects/dynamo/tests/test_render_graph_deployment.py new file mode 100644 index 000000000..21de7e733 --- /dev/null +++ b/projects/dynamo/tests/test_render_graph_deployment.py @@ -0,0 +1,179 @@ +"""Tests for DynamoGraphDeployment manifest rendering.""" + +from __future__ import annotations + +from pathlib import Path + +import pytest +import yaml + +from projects.dynamo.orchestration.render_graph_deployment import render_graph_deployment + + +@pytest.fixture +def config_dir(tmp_path: Path) -> Path: + """Create a temporary config directory with the manifest template.""" + manifests_dir = tmp_path / "manifests" + manifests_dir.mkdir() + + template = { + "apiVersion": "nvidia.com/v1alpha1", + "kind": "DynamoGraphDeployment", + "metadata": { + "name": "", + "namespace": "", + "labels": { + "app.kubernetes.io/managed-by": "forge", + "forge.openshift.io/project": "dynamo", + }, + }, + "spec": { + "backendFramework": "vllm", + "pvcs": [{"name": "model-cache", "create": False}], + "services": {}, + }, + } + + with (manifests_dir / "dynamo-graph-deployment.yaml").open("w") as fh: + yaml.dump(template, fh) + + return tmp_path + + +@pytest.fixture +def base_profile() -> dict: + return { + "backend_framework": "vllm", + "replicas": 2, + "tensor_parallelism": 1, + "serving_mode": "aggregated", + "runtime_image": "nvcr.io/nvidia/ai-dynamo/vllm-runtime:1.2.1", + "frontend_image": "nvcr.io/nvidia/ai-dynamo/dynamo-frontend:1.2.1", + "vllm_args": ["--gpu-memory-utilization=0.90", "--enable-prefix-caching"], + "env": {"DYN_STORE_KV": "mem"}, + } + + +@pytest.fixture +def model_cache() -> dict: + return {"enabled": True, "pvc": {"name_prefix": "dynamo-model"}} + + +@pytest.fixture +def dynamo_config() -> dict: + return {"helm": {"namespace": "dynamo-system"}} + + +def test_aggregated_manifest_has_correct_structure( + config_dir, base_profile, model_cache, dynamo_config +): + manifest = render_graph_deployment( + config_dir=config_dir, + namespace="test-ns", + model_name="Qwen/Qwen3-0.6B", + model_slug="qwen-qwen3-0-6b", + deployment_profile=base_profile, + model_cache=model_cache, + dynamo_config=dynamo_config, + ) + + assert manifest["apiVersion"] == "nvidia.com/v1alpha1" + assert manifest["kind"] == "DynamoGraphDeployment" + assert manifest["metadata"]["name"] == "dynamo-qwen-qwen3-0-6b" + assert manifest["metadata"]["namespace"] == "test-ns" + assert manifest["spec"]["backendFramework"] == "vllm" + + services = manifest["spec"]["services"] + assert "Epp" in services + assert "VllmWorker" in services + assert services["VllmWorker"]["replicas"] == 2 + assert services["VllmWorker"]["componentType"] == "worker" + + +def test_disaggregated_manifest_has_prefill_and_decode( + config_dir, base_profile, model_cache, dynamo_config +): + base_profile["serving_mode"] = "disaggregated" + base_profile["prefill_replicas"] = 1 + base_profile["decode_replicas"] = 3 + + manifest = render_graph_deployment( + config_dir=config_dir, + namespace="test-ns", + model_name="Qwen/Qwen3-0.6B", + model_slug="qwen-qwen3-0-6b", + deployment_profile=base_profile, + model_cache=model_cache, + dynamo_config=dynamo_config, + ) + + services = manifest["spec"]["services"] + assert "Epp" in services + assert "VllmPrefillWorker" in services + assert "VllmDecodeWorker" in services + assert "VllmWorker" not in services + + assert services["VllmPrefillWorker"]["replicas"] == 1 + assert services["VllmDecodeWorker"]["replicas"] == 3 + + +def test_epp_service_has_correct_config( + config_dir, base_profile, model_cache, dynamo_config +): + manifest = render_graph_deployment( + config_dir=config_dir, + namespace="test-ns", + model_name="Qwen/Qwen3-0.6B", + model_slug="qwen-qwen3-0-6b", + deployment_profile=base_profile, + model_cache=model_cache, + dynamo_config=dynamo_config, + ) + + epp = manifest["spec"]["services"]["Epp"] + assert epp["componentType"] == "epp" + assert epp["replicas"] == 1 + assert "eppConfig" in epp + + env_vars = { + e["name"]: e["value"] + for e in epp["extraPodSpec"]["mainContainer"]["env"] + } + assert env_vars["DYN_MODEL_NAME"] == "Qwen/Qwen3-0.6B" + assert env_vars["DYN_DECODE_FALLBACK"] == "true" + + +def test_gpu_resources_match_tensor_parallelism( + config_dir, base_profile, model_cache, dynamo_config +): + base_profile["tensor_parallelism"] = 4 + + manifest = render_graph_deployment( + config_dir=config_dir, + namespace="test-ns", + model_name="meta-llama/Llama-3-70B", + model_slug="meta-llama-llama-3-70b", + deployment_profile=base_profile, + model_cache=model_cache, + dynamo_config=dynamo_config, + ) + + worker = manifest["spec"]["services"]["VllmWorker"] + assert worker["resources"]["limits"]["gpu"] == "4" + assert worker["resources"]["requests"]["gpu"] == "4" + + +def test_forge_labels_present(config_dir, base_profile, model_cache, dynamo_config): + manifest = render_graph_deployment( + config_dir=config_dir, + namespace="test-ns", + model_name="Qwen/Qwen3-0.6B", + model_slug="qwen-qwen3-0-6b", + deployment_profile=base_profile, + model_cache=model_cache, + dynamo_config=dynamo_config, + ) + + labels = manifest["metadata"]["labels"] + assert labels["app.kubernetes.io/managed-by"] == "forge" + assert labels["forge.openshift.io/project"] == "dynamo" diff --git a/projects/dynamo/tests/test_runtime_config.py b/projects/dynamo/tests/test_runtime_config.py new file mode 100644 index 000000000..1eb0b2c1f --- /dev/null +++ b/projects/dynamo/tests/test_runtime_config.py @@ -0,0 +1,73 @@ +"""Tests for Dynamo runtime configuration.""" + +from __future__ import annotations + +from projects.dynamo.orchestration.runtime_config import ( + _deep_merge, + _normalize_string_or_list, + derive_namespace, + version_tuple, +) + + +class TestNormalizeStringOrList: + def test_none_returns_empty(self): + assert _normalize_string_or_list(None, "test") == [] + + def test_empty_string_returns_empty(self): + assert _normalize_string_or_list("", "test") == [] + + def test_single_string_returns_list(self): + assert _normalize_string_or_list("aggregated", "test") == ["aggregated"] + + def test_list_passthrough(self): + assert _normalize_string_or_list(["a", "b"], "test") == ["a", "b"] + + def test_bracket_string_parsed(self): + result = _normalize_string_or_list("[aggregated, disaggregated]", "test") + assert result == ["aggregated", "disaggregated"] + + def test_quoted_bracket_string_parsed(self): + result = _normalize_string_or_list("['aggregated', 'disaggregated']", "test") + assert result == ["aggregated", "disaggregated"] + + +class TestDeepMerge: + def test_override_scalar(self): + assert _deep_merge({"a": 1}, {"a": 2}) == {"a": 2} + + def test_merge_nested_dicts(self): + base = {"a": {"b": 1, "c": 2}} + override = {"a": {"c": 3, "d": 4}} + assert _deep_merge(base, override) == {"a": {"b": 1, "c": 3, "d": 4}} + + def test_override_replaces_non_dict(self): + assert _deep_merge({"a": [1, 2]}, {"a": [3]}) == {"a": [3]} + + +class TestDeriveNamespace: + def test_basic(self): + assert derive_namespace("my-job", "dynamo", 63) == "dynamo-my-job" + + def test_already_prefixed(self): + assert derive_namespace("dynamo-test", "dynamo", 63) == "dynamo-test" + + def test_truncation(self): + ns = derive_namespace("very-long-job-name-that-exceeds", "dynamo", 20) + assert len(ns) <= 20 + + def test_special_chars_slugified(self): + ns = derive_namespace("My Job/Test", "dynamo", 63) + assert "/" not in ns + assert " " not in ns + + +class TestVersionTuple: + def test_semver(self): + assert version_tuple("4.17.3") == (4, 17, 3) + + def test_with_prefix(self): + assert version_tuple("v1.2.1") == (1, 2, 1) + + def test_openshift_version(self): + assert version_tuple("4.19.9-0.nightly") == (4, 19, 9) diff --git a/projects/dynamo/toolbox/__init__.py b/projects/dynamo/toolbox/__init__.py new file mode 100644 index 000000000..e69de29bb diff --git a/projects/dynamo/toolbox/capture_dynamo_state/main.py b/projects/dynamo/toolbox/capture_dynamo_state/main.py new file mode 100644 index 000000000..5efa85d6b --- /dev/null +++ b/projects/dynamo/toolbox/capture_dynamo_state/main.py @@ -0,0 +1,80 @@ +from __future__ import annotations + +import logging +from pathlib import Path + +from projects.core.dsl import entrypoint, execute_tasks, task +from projects.core.dsl.utils.k8s import oc + +logger = logging.getLogger(__name__) + + +@entrypoint +def run( + *, + artifact_dir: Path, + namespace: str, + dynamo_namespace: str = "dynamo-system", + capture_namespace_events: bool = True, +) -> int: + execute_tasks(locals()) + return 0 + + +@task +def setup_artifacts_directory(args, ctx): + ctx.artifacts_dir = args.artifact_dir / "artifacts" + ctx.artifacts_dir.mkdir(parents=True, exist_ok=True) + return f"Artifacts directory prepared: {ctx.artifacts_dir}" + + +@task +def capture_dynamo_operator_state(args, ctx): + """Capture Dynamo operator pod status and logs.""" + dest = ctx.artifacts_dir / "dynamo-operator-pods.txt" + oc("get", "pods", "-n", args.dynamo_namespace, "-l", "app=dynamo-operator", + "-o", "wide", check=False, stdout_dest=str(dest)) + + log_dest = ctx.artifacts_dir / "dynamo-operator-logs.txt" + oc("logs", "-n", args.dynamo_namespace, "-l", "app=dynamo-operator", + "--tail=200", check=False, stdout_dest=str(log_dest)) + + +@task +def capture_dynamo_crds(args, ctx): + """Capture DynamoGraphDeployment and DynamoComponentDeployment resources.""" + dest = ctx.artifacts_dir / "dynamographdeployments.yaml" + oc("get", "dynamographdeployments", "-n", args.namespace, + "-o", "yaml", check=False, stdout_dest=str(dest)) + + dest2 = ctx.artifacts_dir / "dynamocomponentdeployments.yaml" + oc("get", "dynamocomponentdeployments", "-n", args.namespace, + "-o", "yaml", check=False, stdout_dest=str(dest2)) + + +@task +def capture_infrastructure_state(args, ctx): + """Capture etcd and NATS pod status.""" + for component in ["etcd", "nats"]: + dest = ctx.artifacts_dir / f"{component}-pods.txt" + oc("get", "pods", "-n", args.dynamo_namespace, "-l", f"app={component}", + "-o", "wide", check=False, stdout_dest=str(dest)) + + +@task +def capture_worker_pods(args, ctx): + """Capture Dynamo worker and frontend pod status in the test namespace.""" + dest = ctx.artifacts_dir / "dynamo-pods.txt" + oc("get", "pods", "-n", args.namespace, "-o", "wide", + check=False, stdout_dest=str(dest)) + + +@task +def capture_namespace_events(args, ctx): + """Capture namespace events if enabled.""" + if not args.capture_namespace_events: + return + + dest = ctx.artifacts_dir / "namespace-events.txt" + oc("get", "events", "-n", args.namespace, "--sort-by=.metadata.creationTimestamp", + check=False, stdout_dest=str(dest)) diff --git a/projects/dynamo/toolbox/check_gateway_health/main.py b/projects/dynamo/toolbox/check_gateway_health/main.py new file mode 100644 index 000000000..c42dfebbd --- /dev/null +++ b/projects/dynamo/toolbox/check_gateway_health/main.py @@ -0,0 +1,602 @@ +""" +Dynamo Gateway Stack Health Check + +Automated 6-step diagnostic that walks the inference data path bottom-up: + Gateway proxy → Gateway programmed → HTTPRoute accepted → + EPP running → ext_proc connectivity → EPP→Worker routing + +Stops at the first failure and prints the fix. Designed for the three +gateway controllers we encounter: agentgateway, kgateway, Istio. +""" + +from __future__ import annotations + +import json +import logging +import re +import sys +from pathlib import Path + +from projects.core.dsl import entrypoint, execute_tasks, task +from projects.core.dsl.utils.k8s import oc, oc_get_json + +logger = logging.getLogger(__name__) + +PASS = "\033[32mPASS\033[0m" +FAIL = "\033[31mFAIL\033[0m" +WARN = "\033[33mWARN\033[0m" +INFO = "\033[34mINFO\033[0m" +BOLD = "\033[1m" +RESET = "\033[0m" + + +def _print_step(num: int, title: str) -> None: + print(f"\n{'='*60}") + print(f" Step {num}: {title}") + print(f"{'='*60}") + + +def _print_result(status: str, message: str) -> None: + print(f" [{status}] {message}") + + +def _print_fix(title: str, command: str) -> None: + print(f" {BOLD}Fix — {title}{RESET}") + for line in command.strip().split("\n"): + print(f" $ {line}") + + +def _get_pods_by_label(namespace: str, label: str) -> list[dict]: + data = oc_get_json("pods", namespace=namespace, selector=label, ignore_not_found=True) + if not data: + return [] + return data.get("items", []) + + +def _get_conditions(obj: dict) -> dict[str, dict]: + conditions = {} + for c in obj.get("status", {}).get("conditions", []): + conditions[c["type"]] = c + return conditions + + +@entrypoint +def run( + *, + namespace: str, + gateway_name: str = "inference-gateway", + httproute_name: str | None = None, + model_name: str = "Qwen/Qwen3-0.6B", + artifact_dir: Path | None = None, + skip_smoke: bool = False, +) -> int: + execute_tasks(locals()) + return 0 + + +@task +def step_1_gateway_proxy_pod(args, ctx): + """Check if the gateway proxy pod is running.""" + _print_step(1, "Gateway proxy pod running?") + + label = f"gateway.networking.k8s.io/gateway-name={args.gateway_name}" + pods = _get_pods_by_label(args.namespace, label) + + if not pods: + _print_result(FAIL, "No gateway proxy pods found") + _print_result(INFO, f"Label selector: {label}") + + # Check if GatewayClass exists + result = oc("get", "gatewayclass", "--no-headers", check=False) + if result.returncode != 0 or not result.stdout.strip(): + _print_fix( + "No GatewayClass — install a gateway controller", + "# agentgateway (recommended for Dynamo):\n" + "helm upgrade -i agentgateway-crds oci://cr.agentgateway.dev/charts/agentgateway-crds \\\n" + " --create-namespace --namespace agentgateway-system --version v1.0.0\n" + "helm upgrade -i agentgateway oci://cr.agentgateway.dev/charts/agentgateway \\\n" + " --namespace agentgateway-system --version v1.0.0 \\\n" + " --set inferenceExtension.enabled=true --wait", + ) + else: + _print_result(INFO, f"GatewayClasses found:\n{result.stdout.strip()}") + _print_fix( + "Gateway exists but no pods — check Gateway resource", + f"oc get gateway {args.gateway_name} -n {args.namespace} -o yaml", + ) + ctx.healthy = False + return + + pod = pods[0] + pod_name = pod["metadata"]["name"] + phase = pod["status"].get("phase", "Unknown") + containers = pod["status"].get("containerStatuses", []) + ready = all(c.get("ready", False) for c in containers) + restarts = sum(c.get("restartCount", 0) for c in containers) + + if phase == "Running" and ready: + _print_result(PASS, f"Pod {pod_name} is Running and Ready (restarts={restarts})") + ctx.healthy = True + ctx.gw_pod_name = pod_name + return + + if phase == "Running" and not ready: + _print_result(WARN, f"Pod {pod_name} Running but NOT Ready (restarts={restarts})") + if restarts > 3: + _print_result(FAIL, "CrashLoopBackOff — checking logs for root cause") + _diagnose_gateway_crash(args.namespace, pod_name) + ctx.healthy = False + return + + # Check events for SCC issues + events = oc( + "get", "events", "-n", args.namespace, + "--field-selector", f"reason=FailedCreate", + "--sort-by=.metadata.creationTimestamp", + "-o", "json", check=False, + ) + if events.returncode == 0: + event_data = json.loads(events.stdout) + for event in event_data.get("items", []): + msg = event.get("message", "") + if "forbidden" in msg.lower() and args.gateway_name in event.get("involvedObject", {}).get("name", ""): + if "NET_BIND_SERVICE" in msg: + _print_result(FAIL, "SCC blocks NET_BIND_SERVICE capability") + _print_fix( + "Allow NET_BIND_SERVICE in dynamo-frontend-scc", + f"oc patch scc dynamo-frontend-scc --type merge \\\n" + f" -p '{{\"allowedCapabilities\":[\"NET_BIND_SERVICE\"]}}'\n" + f"oc adm policy add-scc-to-user dynamo-frontend-scc \\\n" + f" -z {args.gateway_name} -n {args.namespace}\n" + f"oc rollout restart deployment {args.gateway_name} -n {args.namespace}", + ) + elif "runAsUser" in msg or "fsGroup" in msg: + uid_match = re.search(r"Invalid value: (?:\[)?(\d+)", msg) + uid = uid_match.group(1) if uid_match else "?" + _print_result(FAIL, f"SCC blocks runAsUser/fsGroup {uid}") + _print_fix( + "Grant SCC to gateway service account", + f"oc adm policy add-scc-to-user dynamo-frontend-scc \\\n" + f" -z {args.gateway_name} -n {args.namespace}\n" + f"oc rollout restart deployment {args.gateway_name} -n {args.namespace}", + ) + else: + _print_result(FAIL, f"SCC rejection: {msg[:120]}...") + ctx.healthy = False + return + + _print_result(FAIL, f"Pod {pod_name} in phase {phase}") + ctx.healthy = False + + +def _diagnose_gateway_crash(namespace: str, pod_name: str) -> None: + logs = oc("logs", pod_name, "-n", namespace, "--previous", "--tail=50", check=False) + if logs.returncode != 0: + logs = oc("logs", pod_name, "-n", namespace, "--tail=50", check=False) + if logs.returncode != 0: + _print_result(INFO, "Could not retrieve logs") + return + + log_text = logs.stdout or "" + + if "too many open files" in log_text.lower() or "socket" in log_text.lower() and "failed" in log_text.lower(): + _print_result(FAIL, "File descriptor exhaustion — too many clusters loaded") + _print_fix( + "Switch from kgateway to agentgateway", + "# kgateway loads ALL cluster endpoints and exhausts fd limits on busy clusters.\n" + "# Install agentgateway instead:\n" + "helm upgrade -i agentgateway oci://cr.agentgateway.dev/charts/agentgateway \\\n" + " --namespace agentgateway-system --version v1.0.0 \\\n" + " --set inferenceExtension.enabled=true --wait\n" + "# Then recreate Gateway with gatewayClassName: agentgateway", + ) + return + + if "x509" in log_text or "certificate" in log_text.lower() and "unknown authority" in log_text.lower(): + _print_result(FAIL, "Istio CA certificate mismatch") + _print_fix( + "Compare CA fingerprints (requires team coordination to fix)", + "# What pods trust:\n" + "oc get cm istio-ca-root-cert -n istio-system \\\n" + " -o jsonpath='{.data.root-cert\\.pem}' | openssl x509 -noout -fingerprint\n" + "# What istiod signs with:\n" + "oc get secret istio-ca-secret -n istio-system \\\n" + " -o jsonpath='{.data.ca-cert\\.pem}' | base64 -d | openssl x509 -noout -fingerprint\n" + "# If different → CA was rotated. Delete istio-ca-secret + restart istiod.", + ) + return + + # Generic crash + error_lines = [l for l in log_text.split("\n") if any(k in l.lower() for k in ["error", "fatal", "critical", "panic", "abort"])] + if error_lines: + _print_result(INFO, "Error lines from logs:") + for line in error_lines[:5]: + print(f" {line.strip()[:120]}") + else: + _print_result(INFO, "No obvious error pattern — check full logs:") + print(f" $ oc logs {pod_name} -n {namespace} --previous") + + +@task +def step_2_gateway_programmed(args, ctx): + """Check if the Gateway resource is programmed.""" + if not getattr(ctx, "healthy", True): + _print_step(2, "Gateway programmed?") + _print_result(WARN, "Skipped — gateway pod not healthy (Step 1)") + return + + _print_step(2, "Gateway programmed?") + + gw_data = oc_get_json( + f"gateway/{args.gateway_name}", namespace=args.namespace, ignore_not_found=True + ) + + if not gw_data: + _print_result(FAIL, f"Gateway '{args.gateway_name}' not found in namespace {args.namespace}") + _print_fix( + "Create a Gateway", + f"cat < int: + execute_tasks(locals()) + return 0 + + +@task +def delete_graph_deployments(args, ctx): + """Delete all DynamoGraphDeployments in the namespace.""" + oc("delete", "dynamographdeployments", "--all", "-n", args.namespace, + "--ignore-not-found", check=False) + logger.info("Deleted DynamoGraphDeployments in %s", args.namespace) + + +@task +def delete_component_deployments(args, ctx): + """Delete all DynamoComponentDeployments in the namespace.""" + oc("delete", "dynamocomponentdeployments", "--all", "-n", args.namespace, + "--ignore-not-found", check=False) + logger.info("Deleted DynamoComponentDeployments in %s", args.namespace) + + +@task +def delete_benchmark_resources(args, ctx): + """Delete benchmark job and associated resources.""" + if not args.benchmark_job_name: + return + + oc("delete", "job", args.benchmark_job_name, "-n", args.namespace, + "--ignore-not-found", check=False) + oc("delete", "pod", "-n", args.namespace, "-l", + f"job-name={args.benchmark_job_name}", "--ignore-not-found", check=False) + logger.info("Deleted benchmark resources for %s", args.benchmark_job_name) + + +@task +def delete_forge_labeled_resources(args, ctx): + """Delete all Forge-managed resources in the namespace.""" + for kind in ["pods", "services", "configmaps", "jobs"]: + oc("delete", kind, "-n", args.namespace, + "-l", "forge.openshift.io/project=dynamo", + "--ignore-not-found", check=False) + logger.info("Deleted Forge-labeled resources in %s", args.namespace) diff --git a/projects/dynamo/toolbox/deploy_dynamo_platform/main.py b/projects/dynamo/toolbox/deploy_dynamo_platform/main.py new file mode 100644 index 000000000..7c13c35c7 --- /dev/null +++ b/projects/dynamo/toolbox/deploy_dynamo_platform/main.py @@ -0,0 +1,98 @@ +from __future__ import annotations + +import logging +from pathlib import Path +from typing import Any + +from projects.core.dsl import entrypoint, execute_tasks, shell, task +from projects.core.dsl.utils.k8s import oc + +logger = logging.getLogger(__name__) + + +@entrypoint +def run( + *, + chart_repo: str | None = None, + chart_name: str = "dynamo-platform", + chart_version: str = "1.2.1", + chart_path: str | None = None, + release_name: str = "dynamo", + namespace: str = "dynamo-system", + values_override: dict[str, Any] | None = None, + wait_timeout_seconds: int = 300, + artifact_dir: Path | None = None, +) -> int: + execute_tasks(locals()) + return 0 + + +@task +def ensure_helm_repo(args, ctx): + """Add the Dynamo Helm repository if using a remote chart.""" + if args.chart_path: + logger.info("Using local chart path: %s, skipping repo add", args.chart_path) + ctx.chart_ref = args.chart_path + return + + if not args.chart_repo: + raise ValueError("Either chart_repo or chart_path must be provided") + + shell.run( + f"helm repo add dynamo {args.chart_repo} --force-update", + check=True, + ) + shell.run("helm repo update dynamo", check=True) + ctx.chart_ref = f"dynamo/{args.chart_name}" + + +@task +def create_namespace(args, ctx): + """Ensure the Dynamo platform namespace exists.""" + oc("create", "namespace", args.namespace, check=False) + logger.info("Namespace %s ensured", args.namespace) + + +@task +def install_or_upgrade_chart(args, ctx): + """Helm install/upgrade the Dynamo platform chart.""" + cmd_parts = [ + "helm", "upgrade", "--install", + args.release_name, + ctx.chart_ref, + f"--namespace={args.namespace}", + f"--timeout={args.wait_timeout_seconds}s", + "--wait", + f"--version={args.chart_version}", + ] + + if args.values_override: + import tempfile + values_file = Path(tempfile.mkdtemp()) / "values-override.yaml" + import yaml + with values_file.open("w") as fh: + yaml.dump(args.values_override, fh) + cmd_parts.extend(["-f", str(values_file)]) + + cmd = " ".join(cmd_parts) + logger.info("Running: %s", cmd) + + result = shell.run(cmd, check=True) + + if args.artifact_dir: + output_file = args.artifact_dir / "artifacts" / "helm-install-output.txt" + output_file.parent.mkdir(parents=True, exist_ok=True) + output_file.write_text(result.stdout or "", encoding="utf-8") + + return f"Helm chart {ctx.chart_ref} installed as {args.release_name}" + + +@task +def wait_for_operator(args, ctx): + """Wait for the Dynamo operator deployment to be available.""" + shell.run( + f"oc rollout status deployment/dynamo-operator " + f"-n {args.namespace} --timeout={args.wait_timeout_seconds}s", + check=True, + ) + logger.info("Dynamo operator is ready in namespace %s", args.namespace) diff --git a/projects/dynamo/toolbox/run_aiperf_benchmark/main.py b/projects/dynamo/toolbox/run_aiperf_benchmark/main.py new file mode 100644 index 000000000..5adb365e0 --- /dev/null +++ b/projects/dynamo/toolbox/run_aiperf_benchmark/main.py @@ -0,0 +1,501 @@ +""" +Run an aiperf benchmark as a K8s Job against a Dynamo endpoint. + +Deploys aiperf in a pod (pip install at runtime), writes results to a PVC, +polls for completion, extracts the JSON summary to the local artifact dir. +Follows the same Job+PVC pattern as the guidellm toolbox. +""" + +from __future__ import annotations + +import json +import logging +import os +import subprocess +from pathlib import Path + +from projects.core.dsl import entrypoint, execute_tasks, retry, task +from projects.core.dsl.utils import write_json, write_text +from projects.core.dsl.utils.k8s import oc, oc_apply, oc_get_json + +logger = logging.getLogger(__name__) + +AIPERF_VERSION = "0.7.0" +AIPERF_IMAGE = "python:3.12-slim" + + +@entrypoint +def run( + *, + endpoint_url: str, + model_name: str, + name: str = "aiperf-benchmark", + namespace: str = "", + pvc_name: str = "forge-dynamo-results", + pvc_size: str = "1Gi", + timeout: int = 7200, + artifact_dir: Path | None = None, + # Dataset + dataset_url: str = "https://raw.githubusercontent.com/kvcache-ai/Mooncake/refs/heads/main/FAST25-release/traces/conversation_trace.jsonl", + dataset_type: str = "mooncake_trace", + dataset_cap: int | None = 2000, + # Endpoint + endpoint_type: str = "chat", + endpoint_path: str = "/v1/chat/completions", + streaming: bool = True, + # Schedule + fixed_schedule: bool = True, + fixed_schedule_auto_offset: bool = True, + # Limits + synthesis_max_isl: int | None = 131072, + # aiperf options + tokenizer: str | None = None, +) -> int: + execute_tasks(locals()) + return 0 + + +@task +def validate_parameters(args, ctx): + """Validate and resolve namespace.""" + if not args.namespace: + result = oc("project", "-q", check=False) + if result.returncode == 0: + ctx.namespace = result.stdout.strip() + else: + raise RuntimeError("Could not auto-detect namespace") + else: + ctx.namespace = args.namespace + + ctx.job_name = args.name + ctx.pvc_name = args.pvc_name + ctx.results_subpath = f"aiperf-{ctx.job_name}" + return f"Namespace: {ctx.namespace}, Job: {ctx.job_name}" + + +@task +def cleanup_previous(args, ctx): + """Delete previous aiperf job if exists.""" + _best_effort_delete("aiperf job", "delete", "job", ctx.job_name, + "-n", ctx.namespace, "--ignore-not-found=true") + _best_effort_delete("aiperf copy pod", "delete", "pod", f"{ctx.job_name}-copy", + "-n", ctx.namespace, "--ignore-not-found=true") + + +@task +def ensure_results_pvc(args, ctx): + """Ensure results PVC exists.""" + existing = oc("get", "pvc", ctx.pvc_name, "-n", ctx.namespace, + "--ignore-not-found", "-oname", check=False) + if existing.stdout.strip(): + logger.info("Results PVC %s already exists", ctx.pvc_name) + return + + oc_apply( + args.artifact_dir / "src" / "aiperf-pvc.yaml", + { + "apiVersion": "v1", + "kind": "PersistentVolumeClaim", + "metadata": { + "name": ctx.pvc_name, + "namespace": ctx.namespace, + "labels": { + "app.kubernetes.io/managed-by": "forge", + "forge.openshift.io/project": "dynamo", + }, + }, + "spec": { + "accessModes": ["ReadWriteOnce"], + "resources": {"requests": {"storage": args.pvc_size}}, + }, + }, + ) + logger.info("Created results PVC %s", ctx.pvc_name) + + +@task +def create_aiperf_job(args, ctx): + """Render and apply the aiperf benchmark Job.""" + (args.artifact_dir / "src").mkdir(parents=True, exist_ok=True) + + script = _build_aiperf_script(args, ctx) + tokenizer = args.tokenizer or args.model_name + + job_manifest = { + "apiVersion": "batch/v1", + "kind": "Job", + "metadata": { + "name": ctx.job_name, + "namespace": ctx.namespace, + "labels": { + "app.kubernetes.io/managed-by": "forge", + "forge.openshift.io/project": "dynamo", + }, + }, + "spec": { + "backoffLimit": 0, + "template": { + "spec": { + "restartPolicy": "Never", + "volumes": [ + {"name": "results", "persistentVolumeClaim": {"claimName": ctx.pvc_name}}, + ], + "containers": [{ + "name": "aiperf", + "image": AIPERF_IMAGE, + "command": ["/bin/bash", "-c", script], + "volumeMounts": [ + {"name": "results", "mountPath": "/results"}, + ], + "env": [ + {"name": "AIPERF_HTTP_SSL_VERIFY", "value": "false"}, + ], + "resources": { + "requests": {"cpu": "2", "memory": "4Gi"}, + "limits": {"cpu": "4", "memory": "8Gi"}, + }, + }], + }, + }, + }, + } + + oc_apply(args.artifact_dir / "src" / "aiperf-job.yaml", job_manifest) + logger.info("Created aiperf job %s", ctx.job_name) + + +def _build_aiperf_script(args, ctx) -> str: + """Build the shell script that runs inside the Job pod.""" + tokenizer = args.tokenizer or args.model_name + aiperf_args = [ + f"--model {args.model_name}", + f"--url {args.endpoint_url}", + f"--endpoint-type {args.endpoint_type}", + f"--endpoint {args.endpoint_path}", + "--input-file /tmp/dataset.jsonl", + f"--custom-dataset-type {args.dataset_type}", + f"--tokenizer {tokenizer}", + f"--artifact-dir /results/{ctx.results_subpath}", + "--ui none", + ] + if args.streaming: + aiperf_args.append("--streaming") + if args.fixed_schedule: + aiperf_args.append("--fixed-schedule") + if args.fixed_schedule_auto_offset: + aiperf_args.append("--fixed-schedule-auto-offset") + if args.synthesis_max_isl is not None: + aiperf_args.append(f"--synthesis-max-isl {args.synthesis_max_isl}") + + cap_code = "" + if args.dataset_cap: + cap_code = ( + f"lines = open('/tmp/dataset_full.jsonl').readlines()[:{args.dataset_cap}]\n" + "open('/tmp/dataset.jsonl', 'w').writelines(lines)\n" + f"print(f'Capped to {{len(lines)}} entries')" + ) + dl_target = "/tmp/dataset_full.jsonl" + else: + dl_target = "/tmp/dataset.jsonl" + cap_code = "" + + return f"""set -e +export HOME=/tmp PIP_CACHE_DIR=/tmp/pip-cache +pip install -q --user aiperf=={AIPERF_VERSION} 2>&1 | tail -3 +export PATH="/tmp/.local/bin:$PATH" +python3 -c " +import urllib.request +urllib.request.urlretrieve('{args.dataset_url}', '{dl_target}') +{cap_code} +print('Dataset ready') +" +mkdir -p /results/{ctx.results_subpath} +aiperf profile {' '.join(aiperf_args)} +echo "aiperf completed" +""" + + +@retry(attempts=360, delay=10, backoff=1.0) +@task +def wait_for_completion(args, ctx): + """Poll until aiperf job completes.""" + active = oc("get", "job", ctx.job_name, "-n", ctx.namespace, + "-o", "jsonpath={.status.active}", check=False) + if active.returncode == 0 and active.stdout.strip() == "1": + logger.info("Job %s still running...", ctx.job_name) + return False + + succeeded = oc("get", "job", ctx.job_name, "-n", ctx.namespace, + "-o", "jsonpath={.status.succeeded}", check=False) + failed = oc("get", "job", ctx.job_name, "-n", ctx.namespace, + "-o", "jsonpath={.status.failed}", check=False) + + if succeeded.returncode == 0 and succeeded.stdout.strip() == "1": + return f"aiperf job {ctx.job_name} completed" + + if failed.returncode == 0 and failed.stdout.strip() == "1": + _capture_job_state(args.artifact_dir, ctx.namespace, ctx.job_name) + raise RuntimeError(f"aiperf job {ctx.job_name} failed — check artifacts for logs") + + return False + + +@task +def capture_job_state(args, ctx): + """Capture job logs and pod state.""" + _capture_job_state(args.artifact_dir, ctx.namespace, ctx.job_name) + + +@task +def create_copy_pod(args, ctx): + """Create a pod to read results from PVC.""" + pod_data = oc_get_json("pods", namespace=ctx.namespace, + selector=f"job-name={ctx.job_name}", ignore_not_found=True) + node_name = None + if pod_data and pod_data.get("items"): + node_name = pod_data["items"][0].get("spec", {}).get("nodeName") + + copy_pod = { + "apiVersion": "v1", + "kind": "Pod", + "metadata": { + "name": f"{ctx.job_name}-copy", + "namespace": ctx.namespace, + }, + "spec": { + "restartPolicy": "Never", + "volumes": [ + {"name": "results", "persistentVolumeClaim": {"claimName": ctx.pvc_name}}, + ], + "containers": [{ + "name": "copy", + "image": "busybox", + "command": ["sleep", "3600"], + "volumeMounts": [ + {"name": "results", "mountPath": "/results"}, + ], + }], + }, + } + if node_name: + copy_pod["spec"]["nodeName"] = node_name + + oc_apply(args.artifact_dir / "src" / "aiperf-copy-pod.yaml", copy_pod) + + +@retry(attempts=24, delay=5, backoff=1.0) +@task +def wait_copy_pod_ready(args, ctx): + """Wait for copy pod to be ready.""" + payload = oc_get_json("pod", name=f"{ctx.job_name}-copy", namespace=ctx.namespace) + conditions = payload.get("status", {}).get("conditions", []) + if any(c.get("type") == "Ready" and c.get("status") == "True" for c in conditions): + return f"Copy pod ready" + return False + + +@task +def extract_results(args, ctx): + """Extract ALL aiperf output files from PVC via copy pod.""" + results_dir = args.artifact_dir / "artifacts" / "results" + results_dir.mkdir(parents=True, exist_ok=True) + + copy_pod = f"{ctx.job_name}-copy" + remote_base = f"/results/{ctx.results_subpath}" + + # List all files on PVC + listing = oc("exec", "-n", ctx.namespace, copy_pod, + "--", "find", remote_base, "-type", "f", check=False, log_stdout=False) + if listing.returncode != 0: + logger.warning("Could not list aiperf results at %s", remote_base) + return + + remote_files = [l.strip() for l in (listing.stdout or "").split("\n") if l.strip()] + logger.info("Found %d result files on PVC", len(remote_files)) + + for remote_path in remote_files: + rel = remote_path.replace(remote_base + "/", "", 1) + local_path = results_dir / rel + local_path.parent.mkdir(parents=True, exist_ok=True) + + result = oc("exec", "-n", ctx.namespace, copy_pod, + "--", "cat", remote_path, check=False, log_stdout=False) + if result.returncode == 0 and result.stdout: + write_text(local_path, result.stdout) + logger.info(" extracted: %s (%d bytes)", rel, len(result.stdout)) + else: + logger.warning(" failed: %s", rel) + + # Parse summary + summary_path = results_dir / "profile_export_aiperf.json" + if summary_path.exists(): + try: + full = json.loads(summary_path.read_text()) + ctx.aiperf_results = full + summary = { + "request_count": _m(full, "request_count"), + "error_count": _m(full, "error_request_count"), + "duration_s": round(_m(full, "benchmark_duration") or 0, 1), + "throughput_rps": round(_m(full, "request_throughput") or 0, 2), + "output_tps": round(_m(full, "output_token_throughput") or 0, 2), + "total_tps": round(_m(full, "total_token_throughput") or 0, 2), + "ttft_avg_ms": round(_m(full, "time_to_first_token") or 0, 2), + "ttft_p95_ms": round(_m(full, "time_to_first_token", "p95") or 0, 2), + "itl_avg_ms": round(_m(full, "inter_token_latency") or 0, 2), + "itl_p95_ms": round(_m(full, "inter_token_latency", "p95") or 0, 2), + "latency_avg_ms": round(_m(full, "request_latency") or 0, 2), + "latency_p95_ms": round(_m(full, "request_latency", "p95") or 0, 2), + } + write_json(results_dir / "aiperf_summary.json", summary) + logger.info("=== aiperf Results ===") + for k, v in summary.items(): + logger.info(" %s: %s", k, v) + except (json.JSONDecodeError, TypeError) as e: + logger.warning("Failed to parse aiperf results: %s", e) + + +@task +def build_mlflow_metadata(args, ctx): + """Build MLflow-compatible metadata file for caliper-export.""" + results = getattr(ctx, "aiperf_results", None) + if not results: + return + + metadata = { + "parameters": { + "benchmark_tool": "aiperf", + "model": args.model_name, + "endpoint_url": args.endpoint_url, + "endpoint_type": args.endpoint_type, + "dataset_type": args.dataset_type, + "dataset_cap": str(args.dataset_cap or "full"), + "streaming": str(args.streaming), + "synthesis_max_isl": str(args.synthesis_max_isl or ""), + }, + "metrics": { + "throughput/request_throughput": _m(results, "request_throughput") or 0, + "tokens/output_token_throughput": _m(results, "output_token_throughput") or 0, + "tokens/total_token_throughput": _m(results, "total_token_throughput") or 0, + "latency/request_latency_ms": _m(results, "request_latency") or 0, + "latency/request_latency_p95_ms": _m(results, "request_latency", "p95") or 0, + "ttft/time_to_first_token_ms": _m(results, "time_to_first_token") or 0, + "ttft/time_to_first_token_p95_ms": _m(results, "time_to_first_token", "p95") or 0, + "itl/inter_token_latency_ms": _m(results, "inter_token_latency") or 0, + "itl/inter_token_latency_p95_ms": _m(results, "inter_token_latency", "p95") or 0, + "request_count": _m(results, "request_count") or 0, + "error_request_count": _m(results, "error_request_count") or 0, + "total_output_tokens": _m(results, "total_output_tokens") or 0, + }, + "tags": { + "benchmark_tool": "aiperf", + "framework": "dynamo", + }, + "description": ( + f"aiperf benchmark: {args.model_name}, " + f"dataset={args.dataset_type} cap={args.dataset_cap}, " + f"endpoint={args.endpoint_url}" + ), + } + + meta_path = args.artifact_dir / "artifacts" / "mlflow_run_metadata.json" + meta_path.parent.mkdir(parents=True, exist_ok=True) + write_json(meta_path, metadata) + logger.info("MLflow metadata written to %s", meta_path) + + +@task +def export_to_mlflow(args, ctx): + """Push results + metadata to MLflow via Caliper export (if configured).""" + meta_path = args.artifact_dir / "artifacts" / "mlflow_run_metadata.json" + results_dir = args.artifact_dir / "artifacts" / "results" + + if not meta_path.exists() or not results_dir.exists(): + logger.info("No metadata or results to export, skipping MLflow push") + return + + try: + from projects.caliper.engine.file_export.runner import run_file_export + except ImportError: + logger.info("Caliper not available, skipping MLflow export") + return + + # Read MLflow connection from env (set by vault or manually) + tracking_uri = os.environ.get("MLFLOW_TRACKING_URI") + if not tracking_uri: + logger.info("MLFLOW_TRACKING_URI not set, skipping MLflow export") + return + + run_metadata = json.loads(meta_path.read_text()) + insecure = os.environ.get("MLFLOW_TRACKING_INSECURE_TLS", "").lower() == "true" + connection = {} + if os.environ.get("MLFLOW_TRACKING_USERNAME"): + connection["username"] = os.environ["MLFLOW_TRACKING_USERNAME"] + connection["password"] = os.environ.get("MLFLOW_TRACKING_PASSWORD", "") + if insecure: + connection["insecure_tls"] = True + + experiment = os.environ.get("MLFLOW_EXPERIMENT", "forge-dynamo") + run_name = os.environ.get("MLFLOW_RUN_NAME", f"forge-{args.name}") + + logger.info("Exporting to MLflow: %s experiment=%s", tracking_uri, experiment) + results = run_file_export( + source=results_dir, + backends=["mlflow"], + dry_run=False, + mlflow_tracking_uri=tracking_uri, + mlflow_experiment=experiment, + mlflow_run_id=None, + mlflow_run_name=run_name, + mlflow_insecure_tls=insecure, + mlflow_connection=connection if connection else None, + mlflow_run_metadata=run_metadata, + verbose=True, + ) + for r in results: + logger.info("MLflow export: %s — %s", r.status, r.detail) + if r.metadata: + for k, v in r.metadata.items(): + logger.info(" %s: %s", k, v) + + +@task +def cleanup_job_resources(args, ctx): + """Delete job and copy pod, keep PVC for future runs.""" + _best_effort_delete("copy pod", "delete", "pod", f"{ctx.job_name}-copy", + "-n", ctx.namespace, "--ignore-not-found=true") + _best_effort_delete("aiperf job", "delete", "job", ctx.job_name, + "-n", ctx.namespace, "--ignore-not-found=true") + + +def _best_effort_delete(description: str, *oc_args: str) -> None: + try: + oc(*oc_args, check=False, timeout_seconds=60) + except subprocess.TimeoutExpired: + logger.warning("Timed out deleting %s", description) + + +def _capture_job_state(artifact_dir: Path, namespace: str, job_name: str) -> None: + artifacts = artifact_dir / "artifacts" + artifacts.mkdir(parents=True, exist_ok=True) + + result = oc("logs", f"job/{job_name}", "-n", namespace, check=False, log_stdout=False) + if result.returncode == 0 and result.stdout: + write_text(artifacts / "aiperf_job.logs", result.stdout) + + result = oc("get", "job", job_name, "-n", namespace, "-oyaml", check=False, log_stdout=False) + if result.returncode == 0 and result.stdout: + write_text(artifacts / "aiperf_job.yaml", result.stdout) + + +def _m(data: dict, key: str, field: str = "avg"): + v = data.get(key, {}) + if isinstance(v, dict): + r = v.get(field) + try: + return float(r) if r is not None else None + except (TypeError, ValueError): + return None + try: + return float(v) + except (TypeError, ValueError): + return None