sskj/scripts/common/adaptive_bench_lib.sh
yy-fighting 0fdd3749c7 Add P800 sglang TP/DP matrix experiment
- New experiments/p800/dsv4_p800_sglang_tp_dp_matrix: TP8/DP1, TP4/DP2,
  TP2/DP4 matrix with smoke results; TP2/DP4 documents the weight-loading
  OOM root cause (274 GiB INT8 weights sharded only across TP group).
- Launch args drop --ep-size/--chunked-prefill-size/--max-prefill-tokens/
  --max-running-requests; experts fall back to TP sharding.
- Move dsv4_p800_256k_4k_probe under experiments/p800/.
- scripts/common: jq-free parsing fixes in adaptive_bench_lib.sh and
  parse_backend.py.
- .gitignore: cover raw_outputs under nested platform experiment layout.
2026-07-16 06:49:21 +00:00

802 lines
27 KiB
Bash
Executable File

#!/usr/bin/env bash
# Shared adaptive-concurrency search loop for serving benchmarks.
ADAPTIVE_COMMON_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)"
ADAPTIVE_HELPER="${ADAPTIVE_COMMON_DIR}/adaptive_concurrency.py"
adaptive_timestamp() {
date --iso-8601=seconds
}
adaptive_json_number_or_null() {
local value="${1:-}"
if [[ -n "$value" ]]; then
printf '%s' "$value"
else
printf 'null'
fi
}
adaptive_warmup_request_count() {
local concurrency="$1"
if (( BENCH_WARMUP_MAX_REQUESTS > 0 && concurrency > BENCH_WARMUP_MAX_REQUESTS )); then
printf '%s' "$BENCH_WARMUP_MAX_REQUESTS"
else
printf '%s' "$concurrency"
fi
}
export -f adaptive_warmup_request_count
adaptive_validate_config() {
local name
for name in SEARCH_START_CONCURRENCY SEARCH_MAX_CONCURRENCY \
NUM_PROMPTS_MULTIPLIER PLATEAU_PATIENCE MAX_POINT_RETRIES \
BENCH_WARMUP_MAX_REQUESTS; do
if ! [[ "${!name}" =~ ^[0-9]+$ ]]; then
log "ERROR: ${name} must be a non-negative integer, got ${!name}"
return 1
fi
done
if (( SEARCH_START_CONCURRENCY < 1 )); then
log "ERROR: SEARCH_START_CONCURRENCY must be at least 1"
return 1
fi
if (( SEARCH_MAX_CONCURRENCY < SEARCH_START_CONCURRENCY )); then
log "ERROR: SEARCH_MAX_CONCURRENCY must be >= SEARCH_START_CONCURRENCY"
return 1
fi
if [[ -n "${SEARCH_ADDEND:-}" ]]; then
if ! [[ "${SEARCH_ADDEND}" =~ ^[0-9]+$ ]] || (( SEARCH_ADDEND < 1 )); then
log "ERROR: SEARCH_ADDEND must be at least 1"
return 1
fi
else
if ! [[ "${SEARCH_MULTIPLIER:-}" =~ ^[0-9]+$ ]] || (( SEARCH_MULTIPLIER < 2 )); then
log "ERROR: SEARCH_MULTIPLIER must be at least 2"
return 1
fi
fi
if (( NUM_PROMPTS_MULTIPLIER < 1 || PLATEAU_PATIENCE < 1 )); then
log "ERROR: NUM_PROMPTS_MULTIPLIER and PLATEAU_PATIENCE must be at least 1"
return 1
fi
if ! [[ "$TPS_MIN_GAIN_PCT" =~ ^[0-9]+([.][0-9]+)?$ ]]; then
log "ERROR: TPS_MIN_GAIN_PCT must be a non-negative number"
return 1
fi
if ! [[ "$RANDOM_RANGE_RATIO" =~ ^[0-9]+([.][0-9]+)?$ ]] || \
! awk -v ratio="$RANDOM_RANGE_RATIO" 'BEGIN { exit !(ratio >= 0 && ratio <= 1) }'; then
log "ERROR: RANDOM_RANGE_RATIO must be between 0.0 and 1.0"
return 1
fi
if ! [[ "${TTFT_SLO_MS:-4000}" =~ ^[0-9]+$ ]] || (( TTFT_SLO_MS < 1 )); then
log "ERROR: TTFT_SLO_MS must be a positive integer, got ${TTFT_SLO_MS}"
return 1
fi
if [[ "${ENABLE_TTFT_SLO_STOP:-1}" != "0" && "${ENABLE_TTFT_SLO_STOP:-1}" != "1" ]]; then
log "ERROR: ENABLE_TTFT_SLO_STOP must be 0 or 1, got ${ENABLE_TTFT_SLO_STOP}"
return 1
fi
}
adaptive_start_gpu_monitor() {
local csv_path="$1"
mkdir -p "$(dirname "$csv_path")"
nvidia-smi \
--query-gpu=timestamp,index,memory.used,memory.total,utilization.gpu \
--format=csv \
-l "$GPU_MEM_SAMPLE_INTERVAL_S" \
> "$csv_path" 2>/dev/null &
echo $!
}
adaptive_stop_gpu_monitor() {
local pid="$1"
if kill -0 "$pid" 2>/dev/null; then
kill "$pid" 2>/dev/null || true
sleep 1
kill -9 "$pid" 2>/dev/null || true
fi
}
adaptive_warmup() {
local log_path="$1"
log "running synthetic warmup"
timeout "$SCENARIO_TIMEOUT_S" bash -c \
'engine_run_bench 1024 128 1 1 /dev/null' \
>> "$log_path" 2>&1
}
adaptive_restart_server() {
local tp="$1"
local dp="$2"
log "restarting ${ENGINE} server tp=${tp} dp=${dp}"
engine_stop_server "$tp" "$dp"
sleep "$SERVER_RESTART_COOLDOWN_S"
engine_start_server "$tp" "$dp"
adaptive_warmup "${ADAPTIVE_LOG_DIR}/${ENGINE}_tp${tp}_dp${dp}.warmup.log"
}
adaptive_append_point_from_metrics() {
local metrics_file="$1"
local tp="$2"
local dp="$3"
local mark="$4"
local isl="$5"
local osl="$6"
local concurrency="$7"
local num_prompts="$8"
local attempt="$9"
local gain_pct="${10:-}"
local plateau_streak="${11:-0}"
local raw_file="${12}"
local detail_log="${13}"
local gain_json
gain_json="$(adaptive_json_number_or_null "$gain_pct")"
"$PYTHON" - \
"$metrics_file" \
"$(adaptive_timestamp)" \
"$ENGINE" \
"$tp" "$dp" "$mark" "$isl" "$osl" "$concurrency" "$num_prompts" \
"$(adaptive_warmup_request_count "$concurrency")" \
"$attempt" "$gain_json" "$plateau_streak" "$raw_file" "$detail_log" \
>> "$ADAPTIVE_POINTS_JSONL" <<'PY'
import json, sys
(path, timestamp, engine, tp, dp, mark, isl, osl, concurrency,
num_prompts, warmup_requests, attempt, gain_json, plateau_streak,
raw_file, detail_log) = sys.argv[1:17]
with open(path, "r", encoding="utf-8") as f:
data = json.load(f)
data.update({
"timestamp": timestamp,
"engine": engine,
"tp": int(tp),
"dp": int(dp),
"mark": mark,
"isl": int(isl),
"osl": int(osl),
"concurrency": int(concurrency),
"num_prompts": int(num_prompts),
"warmup_requests": int(warmup_requests),
"attempt": int(attempt),
"gain_pct": json.loads(gain_json),
"plateau_streak": int(plateau_streak),
"error_type": "",
"raw_file": raw_file,
"detail_log": detail_log,
})
print(json.dumps(data, separators=(",", ":"), ensure_ascii=False))
PY
}
adaptive_append_failed_point() {
local tp="$1"
local dp="$2"
local mark="$3"
local isl="$4"
local osl="$5"
local concurrency="$6"
local num_prompts="$7"
local attempt="$8"
local status="$9"
local error_type="${10}"
local raw_file="${11}"
local detail_log="${12}"
"$PYTHON" - \
"$(adaptive_timestamp)" \
"$ENGINE" \
"$tp" "$dp" "$mark" "$isl" "$osl" "$concurrency" "$num_prompts" \
"$(adaptive_warmup_request_count "$concurrency")" \
"$attempt" "$status" "$error_type" "$raw_file" "$detail_log" \
>> "$ADAPTIVE_POINTS_JSONL" <<'PY'
import json, sys
(timestamp, engine, tp, dp, mark, isl, osl, concurrency, num_prompts,
warmup_requests, attempt, status, error_type, raw_file,
detail_log) = sys.argv[1:16]
print(json.dumps({
"timestamp": timestamp,
"engine": engine,
"tp": int(tp),
"dp": int(dp),
"mark": mark,
"isl": int(isl),
"osl": int(osl),
"concurrency": int(concurrency),
"num_prompts": int(num_prompts),
"warmup_requests": int(warmup_requests),
"attempt": int(attempt),
"status": status,
"completed": 0,
"failed": int(num_prompts),
"error_type": error_type,
"validation_errors": [],
"gain_pct": None,
"plateau_streak": 0,
"raw_file": raw_file,
"detail_log": detail_log,
}, separators=(",", ":"), ensure_ascii=False))
PY
}
adaptive_append_shape_summary() {
local tp="$1"
local dp="$2"
local mark="$3"
local isl="$4"
local osl="$5"
local status="$6"
local stop_reason="$7"
local tested_points="$8"
local max_successful_concurrency="$9"
local saturation_concurrency="${10:-}"
local stop_probe_concurrency="${11:-}"
local best_tps_concurrency="${12:-}"
local best_total_tps="${13:-}"
local last_total_tps="${14:-}"
local saturation_json stop_probe_json best_c_json best_tps_json last_tps_json
saturation_json="$(adaptive_json_number_or_null "$saturation_concurrency")"
stop_probe_json="$(adaptive_json_number_or_null "$stop_probe_concurrency")"
best_c_json="$(adaptive_json_number_or_null "$best_tps_concurrency")"
best_tps_json="$(adaptive_json_number_or_null "$best_total_tps")"
last_tps_json="$(adaptive_json_number_or_null "$last_total_tps")"
"$PYTHON" - \
"$(adaptive_timestamp)" \
"$ENGINE" \
"$tp" "$dp" "$mark" "$isl" "$osl" "$status" "$stop_reason" \
"$tested_points" "$SEARCH_MAX_CONCURRENCY" "$max_successful_concurrency" \
"$saturation_json" "$stop_probe_json" "$best_c_json" "$best_tps_json" "$last_tps_json" \
>> "$ADAPTIVE_SHAPES_JSONL" <<'PY'
import json, sys
(timestamp, engine, tp, dp, mark, isl, osl, status, stop_reason,
tested_points, search_cap, max_successful_concurrency,
saturation_json, stop_probe_json, best_c_json, best_tps_json,
last_tps_json) = sys.argv[1:18]
print(json.dumps({
"timestamp": timestamp,
"engine": engine,
"tp": int(tp),
"dp": int(dp),
"mark": mark,
"isl": int(isl),
"osl": int(osl),
"status": status,
"stop_reason": stop_reason,
"tested_points": int(tested_points),
"search_cap": int(search_cap),
"max_successful_concurrency": int(max_successful_concurrency),
"saturation_concurrency": json.loads(saturation_json),
"stop_probe_concurrency": json.loads(stop_probe_json),
"best_tps_concurrency": json.loads(best_c_json),
"best_total_tps": json.loads(best_tps_json),
"last_total_tps": json.loads(last_tps_json),
}, separators=(",", ":"), ensure_ascii=False))
PY
}
adaptive_run_point() {
local tp="$1"
local dp="$2"
local isl="$3"
local osl="$4"
local concurrency="$5"
local num_prompts="$6"
local config_root="$7"
local attempt bench_rc parse_rc gpu_pid
POINT_RAW_FILE=""
POINT_DETAIL_LOG=""
POINT_METRICS_FILE=""
POINT_ATTEMPT=0
POINT_ERROR_TYPE=""
for (( attempt = 1; attempt <= MAX_POINT_RETRIES + 1; attempt++ )); do
POINT_ATTEMPT="$attempt"
POINT_RAW_FILE="${config_root}/raw_outputs/${ENGINE}_adaptive_c${concurrency}_i${isl}_o${osl}_a${attempt}.jsonl"
POINT_DETAIL_LOG="${config_root}/logs/${ENGINE}_c${concurrency}_i${isl}_o${osl}_a${attempt}.log"
POINT_METRICS_FILE="${config_root}/metrics/${ENGINE}_c${concurrency}_i${isl}_o${osl}_a${attempt}.json"
local gpu_csv="${config_root}/gpu_logs/gpu_c${concurrency}_i${isl}_o${osl}_a${attempt}.csv"
log "probe ${ENGINE} tp=${tp} dp=${dp} isl=${isl} osl=${osl} c=${concurrency} warmup=$(adaptive_warmup_request_count "$concurrency") n=${num_prompts} attempt=${attempt}"
gpu_pid="$(adaptive_start_gpu_monitor "$gpu_csv")"
bench_rc=0
timeout "$SCENARIO_TIMEOUT_S" bash -c \
'engine_run_bench "$@"' _ \
"$isl" "$osl" "$concurrency" "$num_prompts" "$POINT_RAW_FILE" \
> "$POINT_DETAIL_LOG" 2>&1 || bench_rc=$?
adaptive_stop_gpu_monitor "$gpu_pid"
if (( bench_rc == 0 )); then
parse_rc=0
"$PYTHON" "$ADAPTIVE_HELPER" parse-result \
--input "$POINT_RAW_FILE" \
--output "$POINT_METRICS_FILE" \
--expected-prompts "$num_prompts" \
--isl "$isl" \
--osl "$osl" \
--input-tolerance-pct "$INPUT_LENGTH_TOLERANCE_PCT" \
--output-tolerance-pct "$OUTPUT_LENGTH_TOLERANCE_PCT" \
>> "$POINT_DETAIL_LOG" 2>&1 || parse_rc=$?
if (( parse_rc == 0 )); then
return 0
fi
POINT_ERROR_TYPE="workload validation failed"
return 21
fi
if engine_detect_oom "$POINT_DETAIL_LOG" "$tp" "$dp"; then
POINT_ERROR_TYPE="detected out of memory"
return 20
fi
POINT_ERROR_TYPE="benchmark rc=${bench_rc}"
if (( attempt <= MAX_POINT_RETRIES )); then
log "probe failed without OOM; restarting server before retry"
if ! adaptive_restart_server "$tp" "$dp"; then
POINT_ERROR_TYPE="benchmark failed and server restart failed"
return 22
fi
continue
fi
return 22
done
}
adaptive_run_shape() {
local tp="$1"
local dp="$2"
local mark="$3"
local isl="$4"
local osl="$5"
local config_root="$6"
local concurrency="$SEARCH_START_CONCURRENCY"
local previous_tps=""
local current_tps=""
local best_tps=""
local best_concurrency=""
local max_successful_concurrency=0
local last_total_tps=""
local plateau_streak=0
local plateau_start_concurrency=""
local saturation_concurrency=""
local stop_probe_concurrency=""
local stop_reason="SEARCH_CAP_REACHED"
local shape_status="COMPLETED"
local tested_points=0
local gain_pct=""
local meaningful_gain=1
local point_rc=0
ADAPTIVE_RESTART_NEEDED=0
ADAPTIVE_FATAL_WORKLOAD=0
while (( concurrency <= SEARCH_MAX_CONCURRENCY )); do
local num_prompts=$((concurrency * NUM_PROMPTS_MULTIPLIER))
point_rc=0
adaptive_run_point "$tp" "$dp" "$isl" "$osl" "$concurrency" "$num_prompts" "$config_root" || point_rc=$?
tested_points=$((tested_points + 1))
stop_probe_concurrency="$concurrency"
if (( point_rc == 0 )); then
current_tps="$("$PYTHON" -c 'import json, sys; print(json.load(open(sys.argv[1]))["total_tps"])' "$POINT_METRICS_FILE")"
gain_pct=""
meaningful_gain=1
if [[ -n "$previous_tps" ]]; then
IFS=$'\t' read -r gain_pct meaningful_gain < <(
"$PYTHON" "$ADAPTIVE_HELPER" gain \
--previous "$previous_tps" \
--current "$current_tps" \
--threshold-pct "$TPS_MIN_GAIN_PCT"
)
if (( meaningful_gain == 1 )); then
plateau_streak=0
plateau_start_concurrency=""
else
plateau_streak=$((plateau_streak + 1))
if [[ -z "$plateau_start_concurrency" ]]; then
plateau_start_concurrency="$concurrency"
fi
fi
fi
if [[ -z "$best_tps" ]] || awk -v current="$current_tps" -v best="$best_tps" 'BEGIN { exit !(current > best) }'; then
best_tps="$current_tps"
best_concurrency="$concurrency"
fi
max_successful_concurrency="$concurrency"
last_total_tps="$current_tps"
# TTFT SLO check: stop searching if TTFT P95 exceeds the SLO threshold
local ttft_p95_ms=""
if [[ "${ENABLE_TTFT_SLO_STOP:-1}" == "1" ]]; then
ttft_p95_ms="$("$PYTHON" -c 'import json, sys; print(json.load(open(sys.argv[1]))["ttft_p95_ms"])' "$POINT_METRICS_FILE")"
if awk -v ttft="$ttft_p95_ms" -v slo="${TTFT_SLO_MS:-4000}" 'BEGIN { exit !(ttft > slo) }'; then
adaptive_append_point_from_metrics \
"$POINT_METRICS_FILE" "$tp" "$dp" "$mark" "$isl" "$osl" \
"$concurrency" "$num_prompts" "$POINT_ATTEMPT" "$gain_pct" \
"$plateau_streak" "$POINT_RAW_FILE" "$POINT_DETAIL_LOG"
stop_reason="TTFT_SLO_EXCEEDED"
shape_status="TTFT_BOUNDARY"
log "probe result c=${concurrency} total_tps=${current_tps} ttft_p95=${ttft_p95_ms}ms > SLO=${TTFT_SLO_MS}ms, stopping shape"
break
fi
fi
adaptive_append_point_from_metrics \
"$POINT_METRICS_FILE" "$tp" "$dp" "$mark" "$isl" "$osl" \
"$concurrency" "$num_prompts" "$POINT_ATTEMPT" "$gain_pct" \
"$plateau_streak" "$POINT_RAW_FILE" "$POINT_DETAIL_LOG"
log "probe result c=${concurrency} total_tps=${current_tps} gain_pct=${gain_pct:-first} plateau_streak=${plateau_streak}/${PLATEAU_PATIENCE}"
if (( plateau_streak >= PLATEAU_PATIENCE )); then
saturation_concurrency="$plateau_start_concurrency"
stop_reason="TPS_PLATEAU"
break
fi
previous_tps="$current_tps"
elif (( point_rc == 20 )); then
adaptive_append_failed_point \
"$tp" "$dp" "$mark" "$isl" "$osl" "$concurrency" "$num_prompts" \
"$POINT_ATTEMPT" "OOM" "$POINT_ERROR_TYPE" "$POINT_RAW_FILE" "$POINT_DETAIL_LOG"
shape_status="OOM_BOUNDARY"
stop_reason="OOM"
ADAPTIVE_RESTART_NEEDED=1
break
elif (( point_rc == 21 )); then
adaptive_append_point_from_metrics \
"$POINT_METRICS_FILE" "$tp" "$dp" "$mark" "$isl" "$osl" \
"$concurrency" "$num_prompts" "$POINT_ATTEMPT" "" 0 \
"$POINT_RAW_FILE" "$POINT_DETAIL_LOG"
shape_status="INVALID_WORKLOAD"
stop_reason="INVALID_WORKLOAD"
ADAPTIVE_FATAL_WORKLOAD=1
break
else
adaptive_append_failed_point \
"$tp" "$dp" "$mark" "$isl" "$osl" "$concurrency" "$num_prompts" \
"$POINT_ATTEMPT" "FAILED" "$POINT_ERROR_TYPE" "$POINT_RAW_FILE" "$POINT_DETAIL_LOG"
shape_status="FAILED"
stop_reason="BENCHMARK_FAILED"
ADAPTIVE_RESTART_NEEDED=1
break
fi
if (( concurrency == SEARCH_MAX_CONCURRENCY )); then
stop_reason="SEARCH_CAP_REACHED"
break
fi
local next_concurrency
if [[ -n "${SEARCH_ADDEND:-}" ]]; then
next_concurrency=$((concurrency + SEARCH_ADDEND))
else
next_concurrency=$((concurrency * SEARCH_MULTIPLIER))
fi
if (( next_concurrency > SEARCH_MAX_CONCURRENCY )); then
next_concurrency="$SEARCH_MAX_CONCURRENCY"
fi
concurrency="$next_concurrency"
done
adaptive_append_shape_summary \
"$tp" "$dp" "$mark" "$isl" "$osl" "$shape_status" "$stop_reason" \
"$tested_points" "$max_successful_concurrency" "$saturation_concurrency" \
"$stop_probe_concurrency" "$best_concurrency" "$best_tps" "$last_total_tps"
log "shape done tp=${tp} dp=${dp} isl=${isl} osl=${osl} stop=${stop_reason} saturation_c=${saturation_concurrency:-none} best_c=${best_concurrency:-none} best_total_tps=${best_tps:-none}"
}
adaptive_matches_filter() {
local value="$1"
local filter="${2:-}"
[[ -z "$filter" || " $filter " == *" $value "* ]]
}
adaptive_print_sequence() {
local concurrency="$SEARCH_START_CONCURRENCY"
local -a sequence=()
while (( concurrency <= SEARCH_MAX_CONCURRENCY )); do
sequence+=("$concurrency")
if (( concurrency == SEARCH_MAX_CONCURRENCY )); then
break
fi
local next_concurrency
if [[ -n "${SEARCH_ADDEND:-}" ]]; then
next_concurrency=$((concurrency + SEARCH_ADDEND))
else
next_concurrency=$((concurrency * SEARCH_MULTIPLIER))
fi
if (( next_concurrency > SEARCH_MAX_CONCURRENCY )); then
next_concurrency="$SEARCH_MAX_CONCURRENCY"
fi
concurrency="$next_concurrency"
done
printf '%s ' "${sequence[@]}"
}
adaptive_cleanup_active_server() {
if [[ -n "${ADAPTIVE_ACTIVE_TP:-}" && -n "${ADAPTIVE_ACTIVE_DP:-}" ]]; then
engine_stop_server "$ADAPTIVE_ACTIVE_TP" "$ADAPTIVE_ACTIVE_DP" || true
ADAPTIVE_ACTIVE_TP=""
ADAPTIVE_ACTIVE_DP=""
fi
}
adaptive_shape_already_tested() {
local tp="$1"
local dp="$2"
local mark="$3"
local isl="$4"
local osl="$5"
local resume_shapes_jsonl="$6"
if [[ ! -f "$resume_shapes_jsonl" ]]; then
return 1
fi
local count
count="$("$PYTHON" - "$tp" "$dp" "$mark" "$isl" "$osl" "$resume_shapes_jsonl" 2>/dev/null <<'PY'
import json, sys
tp, dp, mark, isl, osl = int(sys.argv[1]), int(sys.argv[2]), sys.argv[3], int(sys.argv[4]), int(sys.argv[5])
count = 0
with open(sys.argv[6], "r", encoding="utf-8") as f:
for line in f:
line = line.strip()
if not line:
continue
try:
d = json.loads(line)
except Exception:
continue
if d.get("tp") == tp and d.get("dp") == dp and d.get("mark") == mark \
and d.get("isl") == isl and d.get("osl") == osl:
count += 1
print(count)
PY
)"
if [[ -z "$count" ]]; then
return 1
fi
if (( count > 0 )); then
return 0
fi
return 1
}
adaptive_main() {
local resume_mode=0
local resume_shapes_jsonl=""
if [[ -n "${RESUME_RUN_ID:-}" ]]; then
ADAPTIVE_RUN_ROOT="${RESULT_BASE}/${RESUME_RUN_ID}"
if [[ ! -d "$ADAPTIVE_RUN_ROOT" ]]; then
log "ERROR: RESUME_RUN_ID=${RESUME_RUN_ID} does not exist under ${RESULT_BASE}"
return 1
fi
resume_mode=1
RUN_ID="$RESUME_RUN_ID"
resume_shapes_jsonl="${ADAPTIVE_RUN_ROOT}/adaptive_shapes.jsonl"
if [[ ! -f "$resume_shapes_jsonl" ]]; then
log "WARNING: no adaptive_shapes.jsonl found in ${ADAPTIVE_RUN_ROOT}; will resume with empty history"
fi
log "RESUME mode: continuing run_id=${RESUME_RUN_ID}"
else
RUN_ID="${RUN_ID:-adaptive_$(date '+%Y%m%d-%H%M%S')}"
ADAPTIVE_RUN_ROOT="${RESULT_BASE}/${RUN_ID}"
fi
RESULT_BASE="${RESULT_BASE:-${SCRIPT_DIR}/adaptive_results}"
MATRIX_FILE="${MATRIX_FILE:-${SCRIPT_DIR}/matrix.json}"
MATRIX_MODE="${MATRIX_MODE:-Y}"
ADAPTIVE_LOG_DIR="${ADAPTIVE_RUN_ROOT}/logs"
ADAPTIVE_POINTS_JSONL="${ADAPTIVE_RUN_ROOT}/adaptive_points.jsonl"
ADAPTIVE_SHAPES_JSONL="${ADAPTIVE_RUN_ROOT}/adaptive_shapes.jsonl"
local shapes_tsv="${ADAPTIVE_RUN_ROOT}/shapes.tsv"
mkdir -p "$ADAPTIVE_LOG_DIR"
if (( resume_mode == 0 )); then
: > "$ADAPTIVE_POINTS_JSONL"
: > "$ADAPTIVE_SHAPES_JSONL"
fi
log_init "${ADAPTIVE_LOG_DIR}/orchestrator.log"
adaptive_validate_config
"$PYTHON" "$ADAPTIVE_HELPER" shapes \
--matrix "$MATRIX_FILE" \
--mode "$MATRIX_MODE" \
> "$shapes_tsv"
local tokenize_prompt_json=true
if [[ "$BENCH_DATASET_NAME" == "random" ]]; then
local required_dataset_rows=$((SEARCH_MAX_CONCURRENCY * NUM_PROMPTS_MULTIPLIER))
local dataset_capacity=0
local dataset_error=""
if [[ ! -f "$DATASET_PATH" ]]; then
dataset_error="DATASET_PATH does not exist: ${DATASET_PATH}"
elif ! dataset_capacity="$(
"$PYTHON" "$ADAPTIVE_HELPER" dataset-capacity --path "$DATASET_PATH" 2>/dev/null
)"; then
dataset_error="DATASET_PATH is not a valid ShareGPT JSON file: ${DATASET_PATH}"
fi
if [[ -n "$dataset_error" ]]; then
if [[ "$DRY_RUN" == "1" ]]; then
log "WARNING: ${dataset_error}"
else
log "ERROR: ${dataset_error}"
return 1
fi
fi
if [[ -z "$dataset_error" ]] && (( dataset_capacity < required_dataset_rows )); then
if [[ "$DRY_RUN" == "1" ]]; then
log "WARNING: DATASET_PATH has ${dataset_capacity} usable conversations, but a real run may request ${required_dataset_rows}."
else
log "ERROR: DATASET_PATH has ${dataset_capacity} usable conversations, but adaptive search may request ${required_dataset_rows}. Set DATASET_PATH to a larger ShareGPT file or explicitly use BENCH_DATASET_NAME=random-ids."
return 1
fi
fi
tokenize_prompt_json=false
elif [[ "$BENCH_DATASET_NAME" != random* ]]; then
log "ERROR: BENCH_DATASET_NAME must be random or start with random"
return 1
fi
log "experiment=${EXPERIMENT_NAME} engine=${ENGINE} run_id=${RUN_ID}"
log "search_start=${SEARCH_START_CONCURRENCY} search_cap=${SEARCH_MAX_CONCURRENCY} multiplier=${SEARCH_MULTIPLIER} min_gain_pct=${TPS_MIN_GAIN_PCT} patience=${PLATEAU_PATIENCE} prompts_multiplier=${NUM_PROMPTS_MULTIPLIER} warmup_cap=${BENCH_WARMUP_MAX_REQUESTS}"
log "dataset=${BENCH_DATASET_NAME} random_range_ratio=${RANDOM_RANGE_RATIO} tokenize_prompt=${tokenize_prompt_json} matrix=${MATRIX_FILE} mode=${MATRIX_MODE}"
"$PYTHON" - \
"$EXPERIMENT_NAME" \
"$ENGINE" \
"$RUN_ID" \
"$MODEL_PATH" \
"$HARDWARE" \
"$MATRIX_FILE" \
"$BENCH_DATASET_NAME" \
"$tokenize_prompt_json" \
"$RANDOM_RANGE_RATIO" \
"$SEARCH_START_CONCURRENCY" \
"$SEARCH_MAX_CONCURRENCY" \
"$SEARCH_MULTIPLIER" \
"$NUM_PROMPTS_MULTIPLIER" \
"$TPS_MIN_GAIN_PCT" \
"$PLATEAU_PATIENCE" \
"$BENCH_WARMUP_MAX_REQUESTS" \
"${TTFT_SLO_MS:-4000}" \
"${ENABLE_TTFT_SLO_STOP:-1}" \
> "${ADAPTIVE_RUN_ROOT}/run_manifest.json" <<'PY'
import json, sys
(experiment, engine, run_id, model, hardware, matrix, dataset,
tokenize_prompt, random_range_ratio, search_start, search_cap,
search_multiplier, prompts_multiplier, min_gain_pct, plateau_patience,
warmup_max_requests, ttft_slo_ms, enable_ttft_slo_stop) = sys.argv[1:19]
print(json.dumps({
"experiment": experiment,
"engine": engine,
"run_id": run_id,
"model": model,
"hardware": hardware,
"matrix": matrix,
"dataset": dataset,
"tokenize_prompt": tokenize_prompt == "true",
"random_range_ratio": float(random_range_ratio),
"search": {
"start_concurrency": int(search_start),
"max_concurrency": int(search_cap),
"multiplier": int(search_multiplier),
"num_prompts_multiplier": int(prompts_multiplier),
"min_tps_gain_pct": float(min_gain_pct),
"plateau_patience": int(plateau_patience),
"warmup_max_requests": int(warmup_max_requests),
"ttft_slo_ms": float(ttft_slo_ms),
"enable_ttft_slo_stop": int(enable_ttft_slo_stop),
},
}, indent=2, ensure_ascii=False))
PY
ADAPTIVE_ACTIVE_TP=""
ADAPTIVE_ACTIVE_DP=""
trap adaptive_cleanup_active_server EXIT
trap 'adaptive_cleanup_active_server; exit 130' INT
trap 'adaptive_cleanup_active_server; exit 143' TERM
local abort_all=0
local cfg tp dp config_root server_args shape_count mark isl osl
for cfg in "${PARALLEL_CONFIGS[@]}"; do
read -r tp dp <<< "$cfg"
if ! adaptive_matches_filter "$tp" "$TP_LIST"; then
continue
fi
config_root="${ADAPTIVE_RUN_ROOT}/tp${tp}_dp${dp}"
mkdir -p "${config_root}/raw_outputs" "${config_root}/metrics" \
"${config_root}/logs" "${config_root}/gpu_logs"
server_args="$(engine_build_server_args "$tp" "$dp")"
printf '%s\n' "$server_args" > "${config_root}/server_cmd.txt"
if [[ "$DRY_RUN" == "1" ]]; then
log "DRY_RUN server tp=${tp} dp=${dp}: ${server_args}"
shape_count=0
while IFS=$'\t' read -r mark isl osl; do
[[ "$mark" == "mark" ]] && continue
adaptive_matches_filter "$isl" "$ISL_LIST" || continue
adaptive_matches_filter "$osl" "$OSL_LIST" || continue
if (( GRID_LIMIT > 0 && shape_count >= GRID_LIMIT )); then
break
fi
shape_count=$((shape_count + 1))
if (( resume_mode == 1 )) && adaptive_shape_already_tested "$tp" "$dp" "$mark" "$isl" "$osl" "$resume_shapes_jsonl"; then
log "DRY_RUN RESUME skip already tested shape tp=${tp} dp=${dp} isl=${isl} osl=${osl} mark=${mark}"
continue
fi
log "DRY_RUN tp=${tp} dp=${dp} isl=${isl} osl=${osl} C=[$(adaptive_print_sequence)]"
done < "$shapes_tsv"
continue
fi
ADAPTIVE_ACTIVE_TP="$tp"
ADAPTIVE_ACTIVE_DP="$dp"
if ! engine_start_server "$tp" "$dp"; then
log "ERROR: failed to start ${ENGINE} tp=${tp} dp=${dp}; skipping config"
adaptive_cleanup_active_server
continue
fi
adaptive_warmup "${ADAPTIVE_LOG_DIR}/${ENGINE}_tp${tp}_dp${dp}.warmup.log" || {
log "ERROR: warmup failed for tp=${tp} dp=${dp}; skipping config"
adaptive_cleanup_active_server
continue
}
shape_count=0
while IFS=$'\t' read -r mark isl osl; do
[[ "$mark" == "mark" ]] && continue
adaptive_matches_filter "$isl" "$ISL_LIST" || continue
adaptive_matches_filter "$osl" "$OSL_LIST" || continue
if (( GRID_LIMIT > 0 && shape_count >= GRID_LIMIT )); then
log "GRID_LIMIT=${GRID_LIMIT} reached for tp=${tp} dp=${dp}"
break
fi
shape_count=$((shape_count + 1))
if (( resume_mode == 1 )) && adaptive_shape_already_tested "$tp" "$dp" "$mark" "$isl" "$osl" "$resume_shapes_jsonl"; then
log "RESUME skip already tested shape tp=${tp} dp=${dp} isl=${isl} osl=${osl} mark=${mark}"
continue
fi
log "===== shape START tp=${tp} dp=${dp} isl=${isl} osl=${osl} mark=${mark} ====="
adaptive_run_shape "$tp" "$dp" "$mark" "$isl" "$osl" "$config_root"
if (( ADAPTIVE_FATAL_WORKLOAD == 1 )); then
log "ERROR: invalid synthetic workload; aborting remaining searches"
abort_all=1
break
fi
if (( ADAPTIVE_RESTART_NEEDED == 1 )); then
if ! adaptive_restart_server "$tp" "$dp"; then
log "ERROR: server restart failed; skipping remaining shapes for tp=${tp} dp=${dp}"
break
fi
fi
done < "$shapes_tsv"
adaptive_cleanup_active_server
if (( abort_all == 1 )); then
break
fi
done
"$PYTHON" "$ADAPTIVE_HELPER" summarize \
--points "$ADAPTIVE_POINTS_JSONL" \
--shapes "$ADAPTIVE_SHAPES_JSONL" \
--output-dir "$ADAPTIVE_RUN_ROOT" \
>> "${ADAPTIVE_LOG_DIR}/summarize.log" 2>&1
log "adaptive results saved to ${ADAPTIVE_RUN_ROOT}"
trap - EXIT INT TERM
adaptive_cleanup_active_server
if (( abort_all == 1 )); then
return 2
fi
}