diff --git a/.research/experiments/qwen30-fixed-pd-ttft-admission-diagnosis-20260719.md b/.research/experiments/qwen30-fixed-pd-ttft-admission-diagnosis-20260719.md new file mode 100644 index 0000000..c3ac4f1 --- /dev/null +++ b/.research/experiments/qwen30-fixed-pd-ttft-admission-diagnosis-20260719.md @@ -0,0 +1,45 @@ +# 实验:Qwen30 Fixed-PD TTFT admission diagnosis + +> **状态:** 已批准,运行中 +> +> 用户要求分析 Frontier 在 Qwen30 Fixed-PD 高压 case 上 56--58% TTFT +> selection regret 的根因,并给出简洁结论。 + +## Claim 与决策 + +- **Parent claim:** Frontier 在 capacity knee 附近的配置排序是否会因 state transition error 失效。 +- **目的:** 区分 TP2/TP4 conditional prefill-time 错误、mixed-step composition 错误与 admission queue feedback。 +- **Competing hypotheses:** H1:Frontier 把 TP4 prefill execution 相对 TP2 算慢;H2:decode service time 的绝对误差使 `arrival_rate × residence_time` 越过 MNS cap,首次调度等待被阈值放大;H3:即使固定 admission state,mixed prefill/decode composition 仍反序。 +- **事前预测:** H1 下去掉 queue 后 TP2 仍有更低 prefill time;H2 下去掉 queue 后 TP4 恢复更快,且只有 simulator 的 required concurrency 超过 MNS;H3 下 state-matched stage contrast 仍支持 TP2。 +- **判定规则:** 只有 state ledger/scorer 等价、queue counterfactual 和真机 Running/Waiting 同时支持时才归因 H2;否则保留 H1/H3 并补最小 telemetry。 + +## Setup + +- **自变量:** config 为 Frontier winner `TP2/MNS64`、real mean winner `TP4/MNS32` 与 real p90 winner `TP4/MNS64`。 +- **控制变量:** Qwen3-30B-A3B BF16、community vLLM 0.20、H20、Fixed-PD 4096→256、1.125 req/s/GPU、MBT8192、piecewise graph、原 measured profiles/collectives 和原 257-request traces全部冻结。 +- **Workload:** uniform open-loop arrival;global rate 随 TP 为 2.25/4.5 req/s;prefix cache off;每个真机 cell 三次 fresh-server。 +- **Baselines:** 完整 12-cell frozen real/sim surface;真机三轮 pooled metrics。 +- **Metrics:** TTFT=`first scheduling delay + prefill execution`;request execution/residence time;`arrival_rate × service_time` 相对 MNS;Running/Waiting;stage ledger composition。 + +## 预期产物与 review + +- **预期数据:** 三个 scorer-equivalent Frontier state replays;service/admission decomposition;H1--H3 verdict。 +- **Figure prototype:** `../../runs/frontier-fidelity-envelope-v1/qwen30-fixed-pd-ttft-admission-mock.png`;左图区分 TTFT execution 与 queue,右图显示 required concurrency 是否跨越 MNS。 +- **人工 review:** 已批准(用户要求直接分析清楚该 case)。 +- **Review 意见:** 先复用 existing artifacts 和 CPU replay;只有现有真机 periodic queue proxy 不足时才增加 GPU telemetry。 + +## 复现信息 + +- **Code:** AITuner `feature/sim`;Frontier frozen campaign commit 见原 provenance。 +- **Environment:** dash0;Frontier replay CPU-only,GPU visibility disabled。 +- **输入:** `/home/admin/cpfs/wjh/aituner/qwen30-fixed-pressure-surface-20260719-r1`。 +- **产物路径:** `/home/admin/cpfs/wjh/aituner/qwen30-fixed-pd-ttft-diagnosis-20260719-r1`。 +- **已知 deviation:** 原真机日志只有 10 秒 periodic Running/Waiting,没有 per-iteration ledger;它可验证 steady queue 是否积压,但不用于细粒度 stage timing。 + +## 结果 + +- **观察事实:** 待完成。 +- **异常:** 待完成。 +- **含义:** 待完成。 +- **Claim update:** unchanged +- **下一步:** 复跑三项 Frontier state ledger,并核对 request metrics byte-identical。 diff --git a/runs/frontier-fidelity-envelope-v1/plot_qwen30_fixed_pd_ttft_admission_mock.py b/runs/frontier-fidelity-envelope-v1/plot_qwen30_fixed_pd_ttft_admission_mock.py new file mode 100644 index 0000000..67d76ea --- /dev/null +++ b/runs/frontier-fidelity-envelope-v1/plot_qwen30_fixed_pd_ttft_admission_mock.py @@ -0,0 +1,56 @@ +#!/usr/bin/env python3 +"""Render the schematic for Qwen30 Fixed-PD TTFT admission diagnosis.""" + +from pathlib import Path + +import matplotlib.pyplot as plt +import numpy as np + + +OUTPUT = Path(__file__).with_name("qwen30-fixed-pd-ttft-admission-mock.png") + + +def main() -> None: + figure, axes = plt.subplots(1, 2, figsize=(10.8, 4.2)) + + labels = ["TP2/MNS64", "TP4/MNS64"] + x = np.arange(len(labels)) + axes[0].bar(x, [190, 135], label="Prefill execution (mock)", color="#4c78a8") + axes[0].bar( + x, + [30, 6000], + bottom=[190, 135], + label="Admission wait (mock)", + color="#f58518", + ) + axes[0].set_xticks(x, labels) + axes[0].set_yscale("log") + axes[0].set_ylabel("TTFT decomposition (ms, log scale)") + axes[0].set_title("(a) Execution error or queue amplification?") + axes[0].legend(frameon=False) + + width = 0.34 + axes[1].bar(x - width / 2, [8, 15], width, label="Real required slots (mock)") + axes[1].bar(x + width / 2, [58, 80], width, label="Frontier required slots (mock)") + axes[1].axhline(64, color="#d62728", linestyle="--", label="MNS=64") + axes[1].set_xticks(x, labels) + axes[1].set_ylabel("arrival rate × residence time") + axes[1].set_title("(b) Does modeled state cross admission cap?") + axes[1].legend(frameon=False) + axes[1].text( + 0.02, + 0.97, + "SCHEMATIC / MOCK DATA", + transform=axes[1].transAxes, + va="top", + fontsize=9, + color="#8c2d04", + ) + + figure.suptitle("Qwen30 Fixed-PD TTFT diagnosis", fontsize=13) + figure.tight_layout() + figure.savefig(OUTPUT, dpi=180, bbox_inches="tight") + + +if __name__ == "__main__": + main() diff --git a/runs/frontier-fidelity-envelope-v1/qwen30-fixed-pd-ttft-admission-mock.png b/runs/frontier-fidelity-envelope-v1/qwen30-fixed-pd-ttft-admission-mock.png new file mode 100644 index 0000000..9f39ca1 Binary files /dev/null and b/runs/frontier-fidelity-envelope-v1/qwen30-fixed-pd-ttft-admission-mock.png differ diff --git a/runs/frontier-fidelity-envelope-v1/run_qwen30_fixed_pd_state_replay.py b/runs/frontier-fidelity-envelope-v1/run_qwen30_fixed_pd_state_replay.py new file mode 100644 index 0000000..14bb639 --- /dev/null +++ b/runs/frontier-fidelity-envelope-v1/run_qwen30_fixed_pd_state_replay.py @@ -0,0 +1,229 @@ +#!/usr/bin/env python3 +"""Replay the decisive Qwen30 Fixed-PD cells with Frontier state outputs.""" + +from __future__ import annotations + +import argparse +import importlib.util +import json +import subprocess +import sys +import time +from pathlib import Path +from typing import Any + + +HERE = Path(__file__).resolve().parent +CONFIGS = { + "tp2_mns64": "tp2", + "tp4_mns32": "tp4", + "tp4_mns64": "tp4", +} + + +def load_module(name: str, path: Path): + spec = importlib.util.spec_from_file_location(name, path) + if spec is None or spec.loader is None: + raise ImportError(path) + module = importlib.util.module_from_spec(spec) + sys.modules[spec.name] = module + spec.loader.exec_module(module) + return module + + +BASE = load_module( + "qwen30_fixed_pd_state_base", HERE / "run_qwen235_fixed_pd_state_replay.py" +) +Q30 = load_module( + "qwen30_fixed_pd_surface", HERE / "run_frontier_qwen30_exact_trace_surface.py" +) + + +def parse_args() -> argparse.Namespace: + parser = argparse.ArgumentParser() + parser.add_argument("--base-sim-root", type=Path, required=True) + parser.add_argument("--output-root", type=Path, required=True) + parser.add_argument("--frontier-source", type=Path, required=True) + parser.add_argument("--python-deps", type=Path, required=True) + parser.add_argument("--config", action="append", choices=CONFIGS) + parser.add_argument("--timeout-seconds", type=float, default=1800) + parser.add_argument("--resume", action="store_true") + parser.add_argument("--op-trace", action="store_true") + return parser.parse_args() + + +def run_cell( + *, + config: str, + base_sim_root: Path, + output_root: Path, + frontier_source: Path, + python_deps: Path, + timeout_seconds: float, + resume: bool, + op_trace: bool, +) -> dict[str, Any]: + base_run = base_sim_root / "runs" / config / CONFIGS[config] + base_command_path = base_run / "command.json" + base_result_path = base_run / "result.json" + if not base_command_path.is_file() or not base_result_path.is_file(): + raise FileNotFoundError(f"frozen inputs missing under {base_run}") + base_command = json.loads(base_command_path.read_text()) + base_result = json.loads(base_result_path.read_text()) + if base_result.get("status") != "completed": + raise ValueError(f"frozen result is not completed: {base_result_path}") + + run_root = output_root / config + result_path = run_root / "result.json" + if resume and result_path.is_file(): + previous = json.loads(result_path.read_text()) + if ( + previous.get("status") == "PASS" + and previous.get("inputs", {}).get("base_command_sha256") + == BASE.sha256_file(base_command_path) + and previous.get("inputs", {}).get("base_result_sha256") + == BASE.sha256_file(base_result_path) + and bool(previous.get("op_trace_enabled")) == op_trace + ): + return previous + if run_root.exists() and any(run_root.iterdir()): + raise FileExistsError(f"refusing non-empty output: {run_root}") + run_root.mkdir(parents=True, exist_ok=True) + + command = BASE.transform_command( + base_command, + metrics_root=run_root / "frontier_metrics", + run_id=f"qwen30_fixed_pd_state_{config}", + op_trace=op_trace, + ) + BASE.atomic_json(run_root / "command.json", command) + inputs = { + "base_run": str(base_run.resolve()), + "base_command_sha256": BASE.sha256_file(base_command_path), + "base_result_sha256": BASE.sha256_file(base_result_path), + "base_request_metrics_sha256": base_result["request_metrics_sha256"], + } + manifest = { + "schema": "qwen30-fixed-pd-state-replay-v1", + "config": config, + "inputs": inputs, + "command_sha256": BASE.sha256_file(run_root / "command.json"), + "frontier_git_head": subprocess.check_output( + ["git", "-C", str(frontier_source), "rev-parse", "HEAD"], text=True + ).strip(), + "op_trace_enabled": op_trace, + "controlled_changes": [ + "metrics output directory", + "metrics run id", + "full Frontier stage/batch ledger enabled", + "individual batch metrics enabled", + ] + + (["op-level tracing enabled"] if op_trace else []), + } + BASE.atomic_json(run_root / "run_manifest.json", manifest) + + started = time.monotonic() + with (run_root / "stdout.log").open("w") as stdout, ( + run_root / "stderr.log" + ).open("w") as stderr: + try: + completed = subprocess.run( + command, + cwd=frontier_source, + env=BASE.environment(frontier_source, python_deps), + stdout=stdout, + stderr=stderr, + timeout=timeout_seconds, + check=False, + ) + returncode = int(completed.returncode) + except subprocess.TimeoutExpired: + returncode = 124 + elapsed_seconds = time.monotonic() - started + if returncode != 0: + failure = { + "status": "STOP", + "config": config, + "returncode": returncode, + "elapsed_seconds": elapsed_seconds, + } + BASE.atomic_json(run_root / "failure.json", failure) + raise RuntimeError(f"state replay failed: {failure}") + + fallback_evidence = Q30.collective_fallback_evidence(run_root) + if fallback_evidence: + raise RuntimeError(f"collective-profile fallback detected: {fallback_evidence}") + paths = BASE.find_state_metrics(run_root) + request_metrics_sha256 = BASE.sha256_file(paths["requests"]) + scorer_equivalent = request_metrics_sha256 == base_result["request_metrics_sha256"] + if not scorer_equivalent: + raise RuntimeError( + f"observation changed scorer input for {config}: " + f"{request_metrics_sha256} != {base_result['request_metrics_sha256']}" + ) + state = BASE.summarize_frontier( + system_metrics_path=paths["system"], + request_metrics_path=paths["requests"], + batch_metrics_path=paths["batches"], + ledger_path=paths["ledger"], + ) + BASE.atomic_json(run_root / "common-state.json", state) + if op_trace: + matches = sorted(run_root.glob("frontier_metrics/**/op_traces.jsonl")) + if len(matches) != 1: + raise ValueError(f"expected one op trace, found {len(matches)}") + paths["op_trace"] = matches[0] + result = { + "schema": "qwen30-fixed-pd-state-replay-result-v1", + "status": "PASS", + "config": config, + "elapsed_seconds": elapsed_seconds, + "returncode": returncode, + "inputs": inputs, + "scorer_equivalence": { + "request_metrics_byte_identical": scorer_equivalent, + "request_metrics_sha256": request_metrics_sha256, + "base_metrics": base_result["metrics"], + }, + "state_artifacts": { + name: {"path": str(path.resolve()), "sha256": BASE.sha256_file(path)} + for name, path in paths.items() + }, + "common_state": state, + "collective_fallback_evidence": fallback_evidence, + "op_trace_enabled": op_trace, + } + BASE.atomic_json(result_path, result) + return result + + +def main() -> None: + args = parse_args() + for name in ("base_sim_root", "output_root", "frontier_source", "python_deps"): + setattr(args, name, getattr(args, name).resolve()) + results = [] + for config in tuple(args.config or CONFIGS): + result = run_cell( + config=config, + base_sim_root=args.base_sim_root, + output_root=args.output_root, + frontier_source=args.frontier_source, + python_deps=args.python_deps, + timeout_seconds=args.timeout_seconds, + resume=args.resume, + op_trace=args.op_trace, + ) + results.append(result) + print(json.dumps({"config": config, "status": result["status"]}), flush=True) + BASE.atomic_json( + args.output_root / "state_replay.json", + { + "schema": "qwen30-fixed-pd-state-replay-aggregate-v1", + "status": "PASS", + "results": results, + }, + ) + + +if __name__ == "__main__": + main()