diff --git a/apps/mvp-node/tinygrad_worker.py b/apps/mvp-node/tinygrad_worker.py index 7d6330e..921f116 100755 --- a/apps/mvp-node/tinygrad_worker.py +++ b/apps/mvp-node/tinygrad_worker.py @@ -29,6 +29,9 @@ device_objects: dict[int, dict[str, Any]] = {} next_handle = 42 HEADER_LEN = 40 WORKER_GENERATION = 1 +BENCHMARK_SCHEMA = 1 +_benchmark_start = time.monotonic() +_benchmark_seq = 0 class CpuLineSampler: @@ -142,9 +145,37 @@ def env_flag(name: str, default: bool = True) -> bool: return raw.strip().lower() not in {"0", "false", "no", "off"} +def benchmark_stamp() -> dict[str, Any]: + global _benchmark_seq + _benchmark_seq += 1 + return { + "schema": BENCHMARK_SCHEMA, + "component": "tinygrad-worker", + "pid": os.getpid(), + "seq": _benchmark_seq, + "wall_unix_ms": time.time_ns() // 1_000_000, + "mono_ms": int((time.monotonic() - _benchmark_start) * 1000), + } + + +def env_int(name: str) -> int | None: + raw = os.environ.get(name) + if raw is None: + return None + try: + return int(raw) + except ValueError: + return None def control(**event: Any) -> None: + event.setdefault("benchmark", benchmark_stamp()) + if (run_id := env_int("MVP_RUN_ID")) is not None: + event.setdefault("run_id", run_id) + if (node_id := env_int("MVP_LOGICAL_NODE_ID")) is not None: + event.setdefault("node_id", node_id) + if (stage_index := env_int("MVP_STAGE_INDEX")) is not None: + event.setdefault("stage_index", stage_index) print(json.dumps(event, separators=(",", ":")), flush=True) @@ -157,8 +188,6 @@ def fatal(reason: str, **fields: Any) -> None: raise SystemExit(1) -def test_mode() -> bool: - return os.environ.get("MVP_TINYGRAD_TEST_MODE", "").strip().lower() in {"1", "true", "yes", "on"} def configure_tinygrad_cuda_compiler(device: str) -> None: if device.split(":", 1)[0].upper() != "CUDA": @@ -170,31 +199,31 @@ def configure_tinygrad_cuda_compiler(device: str) -> None: os.environ["CUDA_PTX"] = "1" control(type="TinygradCudaCompilerSelected", requested_device=device, compiler="PTX", reason="nvcc_not_found") +def select_tinygrad_device(device: str) -> str: + device_kind = device.split(":", 1)[0].upper() + if device_kind == "CPU" and ":" not in device and shutil.which("clang") is None: + selected = "CPU:X86" + os.environ["DEV"] = selected + control(type="TinygradCpuCompilerSelected", requested_device=device, selected_device=selected, compiler="X86", reason="clang_not_found") + return selected + os.environ["DEV"] = device + configure_tinygrad_cuda_compiler(device) + return device + def initialize(cmd: dict[str, Any]) -> None: global Tensor, dtypes, arena if int(cmd.get("helper_abi_version", 1)) != 1: fatal("UnsupportedHelperAbi", helper_abi_version=cmd.get("helper_abi_version")) - device = str(cmd.get("backend", {}).get("device") or os.environ.get("DEV") or "CUDA") - os.environ["DEV"] = device + requested_device = str(cmd.get("backend", {}).get("device") or os.environ.get("DEV") or "CUDA") + device = select_tinygrad_device(requested_device) arena_fd = os.environ.get("MVP_ARENA_FD") if arena_fd is not None: arena_bytes = int(os.environ.get("MVP_ARENA_BYTES", "0") or "0") if arena_bytes > 0: arena = mmap.mmap(int(arena_fd), arena_bytes) started = time.monotonic() - if test_mode(): - control( - type="WorkerReady", - pid=os.getpid(), - backend={"requested_device": device, "env_DEV": os.environ.get("DEV"), "test_mode": True}, - cuda_probe=[], - test_mode=True, - elapsed_ms=int((time.monotonic() - started) * 1000), - ) - return - configure_tinygrad_cuda_compiler(device) control(type="TinygradImportStarted", requested_device=device, env_DEV=os.environ.get("DEV")) from tinygrad import Tensor as TinyTensor, dtypes as tiny_dtypes @@ -353,31 +382,42 @@ class PipelineStageTinygradModel: vocab_size: int, head_dim: int, rope_theta: float, + rope_dim: int, + v_head_dim: int, max_context: int, qk_norm: int, num_experts: int, num_experts_per_tok: int, + norm_topk_prob: bool, + qkv_bias: bool, + expert_bias: bool, first_stage: bool, final_stage: bool, nn_mod: Any, + config_cls: Any, block_cls: Any, ) -> None: - self.blk = [ - block_cls( - dim, - hidden_dim, - n_heads, - n_kv_heads, - norm_eps, - head_dim, - rope_theta, - max_context, - qk_norm, - num_experts, - num_experts_per_tok, - ) - for _ in range(block_count) - ] + block_config = config_cls( + num_blocks=block_count, + dim=dim, + hidden_dim=hidden_dim, + n_heads=n_heads, + n_kv_heads=n_kv_heads, + norm_eps=norm_eps, + vocab_size=vocab_size, + head_dim=head_dim, + rope_theta=rope_theta, + rope_dim=rope_dim, + v_head_dim=v_head_dim, + max_context=max_context, + qk_norm=qk_norm, + num_experts=num_experts, + num_experts_per_tok=num_experts_per_tok, + norm_topk_prob=norm_topk_prob, + qkv_bias=qkv_bias, + expert_bias=expert_bias, + ) + self.blk = [block_cls(block_config) for _ in range(block_count)] self.max_context = max_context self.hidden_dim = dim self.first_stage = first_stage @@ -389,16 +429,18 @@ class PipelineStageTinygradModel: self.output = nn_mod.Linear(dim, vocab_size, bias=False) def token_hidden(self, tokens_tensor: Any) -> Any: - return self.token_embd(tokens_tensor) + return self.token_embd(tokens_tensor).float() - def forward_hidden(self, hidden: Any, start_pos: int) -> Any: + def forward_hidden(self, hidden: Any, start_pos: Any) -> Any: for block in self.blk: hidden = block(hidden, start_pos) return hidden.contiguous() def next_token(self, hidden: Any) -> Any: - return self.output(self.output_norm(hidden))[:, -1, :].softmax(-1, dtype="float").argmax(-1, keepdim=True) + return self.output(self.output_norm(hidden))[:, -1, :].argmax(-1, keepdim=True) + def __call__(self, tokens_tensor: Any, start_pos: Any) -> Any: + return self.next_token(self.forward_hidden(self.token_hidden(tokens_tensor), start_pos)) def remap_stage_state_dict( state_dict: dict[str, Any], @@ -436,42 +478,66 @@ def load_pipeline_stage_model( ) -> tuple[PipelineStageTinygradModel, dict[str, Any]]: TensorCls = require_tinygrad() from tinygrad import nn - from tinygrad.apps.llm import TransformerBlock + from tinygrad.llm.gguf import gguf_load + from tinygrad.llm.model import TransformerBlock, TransformerConfig - kv, state_dict = nn.state.gguf_load(TensorCls(path).to(None)) + kv, state_dict = gguf_load(path) state_dict = {key: value.cast("float16") if env_flag("HALF", True) else value for key, value in state_dict.items()} + if "output.weight" not in state_dict and "token_embd.weight" in state_dict: + state_dict["output.weight"] = state_dict["token_embd.weight"] arch = kv["general.architecture"] max_context = min(max_context, int(kv[f"{arch}.context_length"])) n_heads = int(kv[f"{arch}.attention.head_count"]) n_kv_heads = int(kv[f"{arch}.attention.head_count_kv"]) - if arch == "llama": - for name in list(state_dict): - if "attn_q.weight" in name: - state_dict[name] = state_dict[name].rearrange("(n h two) d -> (n two h) d", n=n_heads, two=2) - if "attn_k.weight" in name: - state_dict[name] = state_dict[name].rearrange("(n h two) d -> (n two h) d", n=n_kv_heads, two=2) - total_layers = int(kv[f"{arch}.block_count"]) + dim = int(kv[f"{arch}.embedding_length"]) + kv_lora_rank = int(kv.get(f"{arch}.attention.kv_lora_rank", 0)) + head_dim = int(kv.get(f"{arch}.attention.key_length_mla", kv.get(f"{arch}.attention.key_length", dim // n_heads))) + rope_dim = int(kv.get(f"{arch}.rope.dimension_count", head_dim)) + for name in list(state_dict): + if ("attn_q.weight" in name or "attn_q_b.weight" in name) and (arch == "llama" or kv_lora_rank): + weight = state_dict[name].reshape(n_heads, state_dict[name].shape[0] // n_heads, -1) + prefix = head_dim - rope_dim + state_dict[name] = ( + weight[:, :prefix] + .cat(weight[:, prefix:].rearrange("n (h two) d -> n (two h) d", two=2), dim=1) + .reshape(-1, weight.shape[-1]) + ) + elif arch == "llama" and "attn_k.weight" in name: + weight = state_dict[name].reshape(n_kv_heads, state_dict[name].shape[0] // n_kv_heads, -1) + state_dict[name] = weight.rearrange("n (h two) d -> n (two h) d", two=2).reshape(-1, weight.shape[-1]) + elif kv_lora_rank and "attn_kv_a_mqa.weight" in name: + state_dict[name] = state_dict[name][:kv_lora_rank].cat( + state_dict[name][kv_lora_rank:].rearrange("(h two) d -> (two h) d", two=2), + dim=0, + ) + total_layers = int(kv[f"{arch}.block_count"]) - int(kv.get(f"{arch}.nextn_predict_layers", 0)) first_stage = layer_start == 0 final_stage = layer_end_exclusive >= total_layers qk_key = f"blk.{layer_start}.attn_q_norm.weight" qk_norm = int(state_dict[qk_key].shape[0]) if qk_key in state_dict else 0 stage_model = PipelineStageTinygradModel( block_count=layer_end_exclusive - layer_start, - dim=int(kv[f"{arch}.embedding_length"]), - hidden_dim=int(kv.get(f"{arch}.expert_feed_forward_length", kv[f"{arch}.feed_forward_length"])), + dim=dim, + hidden_dim=int(kv.get(f"{arch}.expert_feed_forward_length", kv.get(f"{arch}.feed_forward_length", 0))), n_heads=n_heads, n_kv_heads=n_kv_heads, norm_eps=float(kv[f"{arch}.attention.layer_norm_rms_epsilon"]), vocab_size=len(kv["tokenizer.ggml.tokens"]), - head_dim=int(kv.get(f"{arch}.attention.key_length", int(kv[f"{arch}.embedding_length"]) // n_heads)), + head_dim=head_dim, rope_theta=float(kv[f"{arch}.rope.freq_base"]), + rope_dim=rope_dim, + v_head_dim=int(kv.get(f"{arch}.attention.value_length_mla", kv.get(f"{arch}.attention.value_length", head_dim))), max_context=max_context, qk_norm=qk_norm, num_experts=int(kv.get(f"{arch}.expert_count", 0)), num_experts_per_tok=int(kv.get(f"{arch}.expert_used_count", 0)), + norm_topk_prob=bool(kv.get(f"{arch}.expert_weights_norm", arch in ("qwen3moe", "qwen35moe"))), + qkv_bias="blk.0.attn_q.bias" in state_dict, + expert_bias=f"blk.{int(kv.get(f'{arch}.leading_dense_block_count', 0))}.exp_probs_b.bias" in state_dict, first_stage=first_stage, final_stage=final_stage, nn_mod=nn, + config_cls=TransformerConfig, block_cls=TransformerBlock, ) stage_state = remap_stage_state_dict( @@ -501,24 +567,6 @@ def load_weights(cmd: dict[str, Any]) -> None: layer_start=int(cmd.get("layer_start", 0)), layer_end_exclusive=int(cmd.get("layer_end_exclusive", 0)), ) - if test_mode(): - model = {"test_mode": True} - tokenizer = {"test_mode": True} - loaded.clear() - loaded.update( - model_id=model_id, - path="mvp-tinygrad-test-mode", - layer_start=int(cmd.get("layer_start", 0)), - layer_end_exclusive=int(cmd.get("layer_end_exclusive", 0)), - ) - control( - type="WeightsLoaded", - model_id=model_id, - path=loaded["path"], - test_mode=True, - elapsed_ms=int((time.monotonic() - started) * 1000), - ) - return control(type="GgufResolveStarted", model_id=model_id, source_kind=source_kind(source)) path = fetch_whole(source) model_bytes = path.stat().st_size @@ -527,7 +575,7 @@ def load_weights(cmd: dict[str, Any]) -> None: layer_end_exclusive = int(cmd.get("layer_end_exclusive", 0)) try: control(type="TinygradLlmImportStarted", model_id=model_id) - from tinygrad.apps.llm import SimpleTokenizer + from tinygrad.llm.cli import SimpleTokenizer control(type="TinygradLlmImportReady", model_id=model_id) max_context_raw = os.environ.get("MVP_MAX_CONTEXT", "512") @@ -739,19 +787,6 @@ def infer_prompt(cmd: dict[str, Any]) -> None: prompt_chars=len(prompt), max_tokens=max_tokens, ) - if test_mode(): - text = f"mvp-test response: {prompt}" - control( - type="PromptCompleted", - request_id=request_id, - model_id=loaded.get("model_id"), - prompt_tokens=[], - generated_tokens=list(range(min(max_tokens, 3))), - text=text, - test_mode=True, - elapsed_ms=int((time.monotonic() - started) * 1000), - ) - return model_prompt, prompt_template = model_prompt_text(prompt) control( type="PromptEncodeStarted", @@ -907,7 +942,7 @@ def object_start_pos(sequence: int, token_count: int) -> int: def materialize_object(payload: bytes, sequence: int) -> dict[str, Any]: - if test_mode() or not isinstance(model, PipelineStageTinygradModel): + if not isinstance(model, PipelineStageTinygradModel): return { "kind": "words", "words": payload_words(payload), @@ -1013,7 +1048,7 @@ def execute_step(cmd: dict[str, Any]) -> None: if ring["direction"] != "egress": fatal("WrongRingDirection", ring_id=output_ring_id, direction=ring["direction"]) final_stage = bool(cmd.get("final_stage")) - if test_mode() or not isinstance(model, PipelineStageTinygradModel): + if not isinstance(model, PipelineStageTinygradModel): if final_stage: base = sum(int(word) for word in obj["words"]) + int(role.get("stage_index", 0)) token = 6 if base % 2 else 8 @@ -1067,7 +1102,7 @@ def release_device_object(cmd: dict[str, Any]) -> None: def encode_prompt(cmd: dict[str, Any]) -> None: prompt = str(cmd.get("prompt", "")) - if tokenizer is not None and not test_mode(): + if tokenizer is not None: model_prompt, _ = model_prompt_text(prompt) tokens = [int(token) for token in tokenizer.encode(model_prompt)] else: @@ -1077,7 +1112,7 @@ def encode_prompt(cmd: dict[str, Any]) -> None: def decode_tokens(cmd: dict[str, Any]) -> None: tokens = [int(token) for token in cmd.get("tokens", [])] - if tokenizer is not None and not test_mode(): + if tokenizer is not None: text = strip_chat_stop_markers(tokenizer.decode(tokens)) else: text = "".join(chr(token) if 32 <= token <= 126 else f"" for token in tokens) diff --git a/crates/mvp-system/src/benchmark_observability.rs b/crates/mvp-system/src/benchmark_observability.rs new file mode 100644 index 0000000..867ef8a --- /dev/null +++ b/crates/mvp-system/src/benchmark_observability.rs @@ -0,0 +1,32 @@ +use std::sync::OnceLock; +use std::sync::atomic::{AtomicU64, Ordering}; +use std::time::{Instant, SystemTime, UNIX_EPOCH}; + +use serde_json::{Value, json}; + +pub const BENCHMARK_SCHEMA: u64 = 1; + +static BENCHMARK_START: OnceLock = OnceLock::new(); +static BENCHMARK_SEQ: AtomicU64 = AtomicU64::new(1); + +pub fn unix_ms_now() -> u64 { + let millis = SystemTime::now() + .duration_since(UNIX_EPOCH) + .unwrap_or_default() + .as_millis(); + u64::try_from(millis).unwrap_or(u64::MAX) +} + +pub fn stamp(component: &'static str) -> Value { + let start = BENCHMARK_START.get_or_init(Instant::now); + let mono_ms = u64::try_from(start.elapsed().as_millis()).unwrap_or(u64::MAX); + let seq = BENCHMARK_SEQ.fetch_add(1, Ordering::Relaxed); + json!({ + "schema": BENCHMARK_SCHEMA, + "component": component, + "pid": std::process::id(), + "seq": seq, + "wall_unix_ms": unix_ms_now(), + "mono_ms": mono_ms, + }) +} diff --git a/crates/mvp-system/src/bin/mvp_chat.rs b/crates/mvp-system/src/bin/mvp_chat.rs index 06b6f92..59abb28 100644 --- a/crates/mvp-system/src/bin/mvp_chat.rs +++ b/crates/mvp-system/src/bin/mvp_chat.rs @@ -22,6 +22,7 @@ use signal_hook::consts::signal::{SIGINT, SIGTERM}; #[cfg(target_os = "linux")] use signal_hook::iterator::Signals; +use mvp_system::benchmark_observability; use mvp_system::config as chat_config; use mvp_system::config::ResolvedVastAiConfig; use mvp_system::node_image::{ @@ -73,7 +74,7 @@ where I: IntoIterator, { let config = Config::from_args(args)?; - let mut progress = ChatDatastream::new(1, config.datastream_frame_log.clone())?; + let mut progress = ChatDatastream::new(config.run_id, config.datastream_frame_log.clone())?; progress.emit( CHAT_LIFECYCLE_CHANNEL, "config", @@ -87,29 +88,48 @@ where }), ); confirm_vastai_if_needed(&config)?; - let image_ref = match prepare_runtime(&config) { - Ok(image_ref) => { - progress.emit( - CHAT_RUNTIME_CHANNEL, - "prepare_runtime", - "ready", - json!({"image_ref": image_ref}), - ); - image_ref - } - Err(error) => { - progress.emit( - CHAT_RUNTIME_CHANNEL, - "prepare_runtime", - "failed", - json!({"error": error}), - ); - progress.archive_pending()?; - return Err(error); - } - }; + progress.emit( + CHAT_RUNTIME_CHANNEL, + "prepare_runtime", + "started", + json!({"provider": config.provider.as_str()}), + ); + let image_ref = + match prepare_runtime_with_progress(&config, prepare_node_image, Some(&mut progress)) { + Ok(image_ref) => { + progress.emit( + CHAT_RUNTIME_CHANNEL, + "prepare_runtime", + "ready", + json!({"image_ref": image_ref}), + ); + image_ref + } + Err(error) => { + progress.emit( + CHAT_RUNTIME_CHANNEL, + "prepare_runtime", + "failed", + json!({"error": error}), + ); + progress.archive_pending()?; + return Err(error); + } + }; + progress.emit( + CHAT_COMPONENT_CHANNEL, + "orchestrator_process_spawn", + "started", + json!({"binary": config.orch_bin.to_string_lossy()}), + ); let mut orch = match OrchChild::spawn(&config, &image_ref) { Ok(orch) => { + progress.emit( + CHAT_COMPONENT_CHANNEL, + "orchestrator_process_spawn", + "ready", + json!({"binary": config.orch_bin.to_string_lossy(), "pid": orch.child.id()}), + ); progress.emit( CHAT_COMPONENT_CHANNEL, "orchestrator_process", @@ -119,6 +139,12 @@ where orch } Err(error) => { + progress.emit( + CHAT_COMPONENT_CHANNEL, + "orchestrator_process_spawn", + "failed", + json!({"binary": config.orch_bin.to_string_lossy(), "error": error}), + ); progress.emit( CHAT_COMPONENT_CHANNEL, "orchestrator_process", @@ -129,8 +155,20 @@ where return Err(error); } }; + progress.emit( + CHAT_RUNTIME_CHANNEL, + "prompt_rpc_wait", + "started", + json!({"addr": config.rpc_addr}), + ); let rpc_addr = match orch.wait_ready(config.rpc_addr.clone()) { Ok(addr) => { + progress.emit( + CHAT_RUNTIME_CHANNEL, + "prompt_rpc_wait", + "ready", + json!({"addr": addr}), + ); progress.emit( CHAT_RUNTIME_CHANNEL, "prompt_rpc", @@ -139,7 +177,13 @@ where ); addr } - Err(_) if STOP_REQUESTED.load(Ordering::SeqCst) => { + Err(error) if STOP_REQUESTED.load(Ordering::SeqCst) => { + progress.emit( + CHAT_RUNTIME_CHANNEL, + "prompt_rpc_wait", + "failed", + json!({"addr": config.rpc_addr, "error": error}), + ); progress.emit( CHAT_LIFECYCLE_CHANNEL, "shutdown", @@ -157,6 +201,12 @@ where return Ok(()); } Err(error) => { + progress.emit( + CHAT_RUNTIME_CHANNEL, + "prompt_rpc_wait", + "failed", + json!({"addr": config.rpc_addr, "error": error}), + ); progress.emit( CHAT_RUNTIME_CHANNEL, "prompt_rpc", @@ -201,6 +251,7 @@ struct Config { image_tag: Option, cached_model: Option, datastream_frame_log: Option, + run_id: u64, vastai_yes: bool, vastai: Option, pipeline_stages: u32, @@ -210,6 +261,7 @@ struct Config { struct ChatDatastream { stream: StreamId, + run_id: u64, endpoint: DatastreamEndpoint, producer: DatastreamProducer, channels: BTreeMap, @@ -233,6 +285,7 @@ impl ChatDatastream { let producer = endpoint.producer(); let mut out = Self { stream, + run_id, endpoint, producer, channels: BTreeMap::new(), @@ -272,6 +325,8 @@ impl ChatDatastream { "type": "ChatProgress", "phase": phase, "status": status, + "run_id": self.run_id, + "benchmark": benchmark_observability::stamp("mvp-chat"), "detail": detail, })) .expect("serialize mvp-chat progress event"); @@ -356,6 +411,7 @@ impl ChatFrameArchive { }; let record = json!({ "arrival_seq": self.next_seq, + "arrival_unix_ms": benchmark_observability::unix_ms_now(), "source": source, "stream": stream.to_string(), "channel": channel, @@ -514,6 +570,7 @@ impl Config { image_tag: first_non_empty([toml.image.tag.clone()]), cached_model, datastream_frame_log, + run_id: args.run_id.unwrap_or(1), vastai_yes: args.vastai_yes, pipeline_stages, max_tokens, @@ -534,6 +591,8 @@ impl Config { self.rpc_addr.clone(), "--max-tokens".to_owned(), self.max_tokens.to_string(), + "--run-id".to_owned(), + self.run_id.to_string(), "--pipeline-stages".to_owned(), self.pipeline_stages.to_string(), "--no-dashboard".to_owned(), @@ -617,6 +676,7 @@ struct ParsedArgs { pipeline_stages: Option, dump_logs: bool, dump_log_path: Option, + run_id: Option, skip_rebuild: bool, cached_model: Option, } @@ -658,6 +718,13 @@ impl ParsedArgs { parsed.pipeline_stages = Some(parse_pipeline_stages_value(&mut args, arg.as_str())?) } + "--run-id" => { + let run_id: u64 = parse_next(&mut args, "--run-id")?; + if run_id == 0 { + return Err("--run-id must be greater than 0".to_owned()); + } + parsed.run_id = Some(run_id); + } "--dump-logs" => { parsed.dump_logs = true; } @@ -916,33 +983,193 @@ fn signal_orch_process_group(child: &Child, signal: libc::c_int) -> io::Result<( type PrepareNodeImageFn = fn(NodeImageRequest) -> Result; +#[allow(dead_code)] fn prepare_runtime(config: &Config) -> Result { prepare_runtime_with(config, prepare_node_image) } +#[allow(dead_code)] fn prepare_runtime_with( config: &Config, prepare_node_image_fn: PrepareNodeImageFn, ) -> Result { - ensure_orch_binary(config)?; + prepare_runtime_with_progress(config, prepare_node_image_fn, None) +} + +fn prepare_runtime_with_progress( + config: &Config, + prepare_node_image_fn: PrepareNodeImageFn, + progress: Option<&mut ChatDatastream>, +) -> Result { + let mut progress = progress; + let binary_mode = if config.skip_rebuild { + "existing_artifact" + } else { + "cargo_build" + }; + emit_chat_progress( + &mut progress, + CHAT_RUNTIME_CHANNEL, + "ensure_orch_binary", + "started", + json!({"mode": binary_mode}), + ); + match ensure_orch_binary(config) { + Ok(()) => emit_chat_progress( + &mut progress, + CHAT_RUNTIME_CHANNEL, + "ensure_orch_binary", + "ready", + json!({"mode": binary_mode}), + ), + Err(error) => { + emit_chat_progress( + &mut progress, + CHAT_RUNTIME_CHANNEL, + "ensure_orch_binary", + "failed", + json!({"mode": binary_mode, "error": error.as_str()}), + ); + return Err(error); + } + } + if config.provider == ProviderKind::Process { - ensure_worker_binary(config)?; + emit_chat_progress( + &mut progress, + CHAT_RUNTIME_CHANNEL, + "ensure_worker_binary", + "started", + json!({"mode": binary_mode}), + ); + match ensure_worker_binary(config) { + Ok(()) => emit_chat_progress( + &mut progress, + CHAT_RUNTIME_CHANNEL, + "ensure_worker_binary", + "ready", + json!({"mode": binary_mode}), + ), + Err(error) => { + emit_chat_progress( + &mut progress, + CHAT_RUNTIME_CHANNEL, + "ensure_worker_binary", + "failed", + json!({"mode": binary_mode, "error": error.as_str()}), + ); + return Err(error); + } + } + emit_chat_progress( + &mut progress, + CHAT_RUNTIME_CHANNEL, + "prepare_node_image", + "skipped", + json!({"provider": config.provider.as_str(), "reason": "process_provider"}), + ); return Ok(config.node_image.clone()); } + if config.skip_rebuild { - ensure_worker_binary(config)?; + emit_chat_progress( + &mut progress, + CHAT_RUNTIME_CHANNEL, + "ensure_worker_binary", + "started", + json!({"mode": binary_mode}), + ); + match ensure_worker_binary(config) { + Ok(()) => emit_chat_progress( + &mut progress, + CHAT_RUNTIME_CHANNEL, + "ensure_worker_binary", + "ready", + json!({"mode": binary_mode}), + ), + Err(error) => { + emit_chat_progress( + &mut progress, + CHAT_RUNTIME_CHANNEL, + "ensure_worker_binary", + "failed", + json!({"mode": binary_mode, "error": error.as_str()}), + ); + return Err(error); + } + } + emit_chat_progress( + &mut progress, + CHAT_RUNTIME_CHANNEL, + "prepare_node_image", + "skipped", + json!({"provider": config.provider.as_str(), "reason": "skip_rebuild"}), + ); return Ok(config.node_image.clone()); } - let prepared = prepare_node_image_fn(NodeImageRequest { + + emit_chat_progress( + &mut progress, + CHAT_RUNTIME_CHANNEL, + "prepare_node_image", + "started", + json!({"provider": config.provider.as_str()}), + ); + let node_bin = match node_bin_for_current_profile() { + Ok(path) => path, + Err(error) => { + emit_chat_progress( + &mut progress, + CHAT_RUNTIME_CHANNEL, + "prepare_node_image", + "failed", + json!({"provider": config.provider.as_str(), "error": error.as_str()}), + ); + return Err(error); + } + }; + let provider = match node_image_provider(config.provider) { + Ok(provider) => provider, + Err(error) => { + emit_chat_progress( + &mut progress, + CHAT_RUNTIME_CHANNEL, + "prepare_node_image", + "failed", + json!({"provider": config.provider.as_str(), "error": error.as_str()}), + ); + return Err(error); + } + }; + let prepared = match prepare_node_image_fn(NodeImageRequest { requested_image: config.node_image.clone(), base_image: BASE_NODE_IMAGE.to_owned(), - node_bin: node_bin_for_current_profile()?, - provider: node_image_provider(config.provider)?, + node_bin, + provider, extra_tag: config.image_tag.clone(), push: false, force_refresh: false, enabled: true, - })?; + }) { + Ok(prepared) => prepared, + Err(error) => { + emit_chat_progress( + &mut progress, + CHAT_RUNTIME_CHANNEL, + "prepare_node_image", + "failed", + json!({"provider": config.provider.as_str(), "error": error.as_str()}), + ); + return Err(error); + } + }; + emit_chat_progress( + &mut progress, + CHAT_RUNTIME_CHANNEL, + "prepare_node_image", + "ready", + json!({"provider": config.provider.as_str(), "image_ref": prepared.image_ref}), + ); Ok(prepared.image_ref) } @@ -1254,7 +1481,12 @@ where json!({"request_id": request_id, "text_bytes": text.len()}), ); } - PromptEvent::Done { .. } => { + PromptEvent::Done { + final_text, + tokens_generated, + elapsed_ms, + .. + } => { if response_started { writeln!(output).map_err(|e| format!("write response terminator: {e}"))?; } else { @@ -1266,7 +1498,13 @@ where CHAT_PROMPT_CHANNEL, "request_completed", "ready", - json!({"request_id": request_id, "response_started": response_started}), + json!({ + "request_id": request_id, + "response_started": response_started, + "tokens_generated": tokens_generated, + "elapsed_ms": elapsed_ms, + "final_text_bytes": final_text.len(), + }), ); break; } @@ -1642,6 +1880,7 @@ mod tests { image_tag: None, cached_model: None, datastream_frame_log: None, + run_id: 1, vastai_yes: false, vastai: None, pipeline_stages: 1, @@ -1704,6 +1943,48 @@ mod tests { .collect() } + #[test] + fn benchmark_observability_chat_progress_records_include_run_id_and_stamp() { + let temp = TempDir::new("chat-progress-archive"); + let archive_path = temp.path().join("frames.ndjson"); + let mut progress = ChatDatastream::new(77, Some(archive_path.clone())) + .expect("chat datastream constructs"); + + progress.emit( + CHAT_RUNTIME_CHANNEL, + "unit_phase", + "ready", + serde_json::json!({"ok": true}), + ); + progress.archive_pending().expect("archive pending frames"); + + let archive = fs::read_to_string(&archive_path).expect("read archive"); + let line = archive.lines().next().expect("archive line"); + let outer: serde_json::Value = serde_json::from_str(line).expect("outer archive JSON"); + let inner_text = outer + .get("payload") + .and_then(|payload| payload.get("value")) + .and_then(serde_json::Value::as_str) + .expect("inner event text"); + let inner: serde_json::Value = serde_json::from_str(inner_text).expect("inner event JSON"); + + assert_eq!( + inner.get("type").and_then(serde_json::Value::as_str), + Some("ChatProgress") + ); + assert_eq!( + inner.get("run_id").and_then(serde_json::Value::as_u64), + Some(77) + ); + assert_eq!( + inner + .get("benchmark") + .and_then(|benchmark| benchmark.get("schema")) + .and_then(serde_json::Value::as_u64), + Some(1) + ); + } + #[test] fn parsed_args_accepts_public_flags() { let parsed = ParsedArgs::parse(strings(&[ @@ -1729,6 +2010,32 @@ mod tests { assert!(parsed.skip_rebuild); } + #[test] + fn benchmark_observability_parsed_args_accepts_run_id_and_forwards_to_orchestrator() { + let parsed = ParsedArgs::parse(strings(&["--run-id", "123"])).expect("run id parses"); + assert_eq!(parsed.run_id, Some(123)); + + let temp = TempDir::new("run-id-config"); + with_process_state(&[], Some(temp.path()), || { + let config = + Config::from_args(strings(&["--run-id", "123"])).expect("config resolves run id"); + assert_eq!(config.run_id, 123); + let args = config.orchestrator_cli_args("resolved-image"); + let run_id_arg = args + .windows(2) + .find(|pair| pair[0] == "--run-id") + .map(|pair| pair[1].as_str()); + assert_eq!(run_id_arg, Some("123"), "{args:?}"); + }); + } + + #[test] + fn benchmark_observability_parsed_args_rejects_zero_run_id() { + let error = + ParsedArgs::parse(strings(&["--run-id", "0"])).expect_err("zero run id should fail"); + assert_eq!(error, "--run-id must be greater than 0"); + } + #[test] fn parsed_args_accepts_cached_model_path() { let parsed = ParsedArgs::parse(strings(&["--cached-model=/tmp/model.gguf"])) @@ -1795,6 +2102,7 @@ mod tests { let defaults = Config::from_args(Vec::::new()).expect("defaults resolve"); assert_eq!(defaults.provider, ProviderKind::Process); assert_eq!(defaults.pipeline_stages, 1); + assert_eq!(defaults.run_id, 1); assert!(defaults.datastream_frame_log.is_none()); assert!(defaults.cached_model.is_none()); assert!(defaults.vastai.is_none()); diff --git a/crates/mvp-system/src/bin/orchestrator.rs b/crates/mvp-system/src/bin/orchestrator.rs index 0a636d5..9d61eef 100644 --- a/crates/mvp-system/src/bin/orchestrator.rs +++ b/crates/mvp-system/src/bin/orchestrator.rs @@ -27,6 +27,7 @@ use mvp_system::actors::node_agent::{ }; use mvp_system::actors::orchestrator::{OrchestratorActor, OrchestratorReport}; use mvp_system::actors::register_mvp_actor_codecs; +use mvp_system::benchmark_observability; use mvp_system::config::{DEFAULT_CONFIG_PATH, TomlConfigOverlay}; #[cfg(feature = "dashboard")] use mvp_system::dashboard_view::MvpClusterDashboardView; @@ -415,6 +416,18 @@ fn run() -> Result<(), String> { orchestrator_actor, )?; + orch_datastream.emit_bootstrap( + dashboard.as_ref(), + config.run_id, + config.node_id, + "prompt_rpc", + "started", + json!({ + "bind":config.rpc_bind.to_string(), + "default_max_tokens":config.default_max_tokens, + }), + ); + let rpc_addr = match spawn_prompt_rpc(config.rpc_bind, work_tx, config.default_max_tokens) { Ok(addr) => { orch_datastream.emit_bootstrap( @@ -1700,9 +1713,6 @@ impl Config { if self.provider == ProviderKind::Docker { keys.push("MVP_DOCKER_GPUS"); } - if std::env::var_os("MVP_TINYGRAD_TEST_MODE").is_some() { - keys.push("MVP_TINYGRAD_TEST_MODE"); - } if local_tinygrad_worker_env(self.provider).is_some() { keys.push("MVP_TINYGRAD_WORKER"); } @@ -1800,7 +1810,6 @@ impl Config { if self.provider == ProviderKind::Docker { env.push(("MVP_DOCKER_GPUS".to_owned(), self.docker_gpus.clone())); } - env.extend(optional_env("MVP_TINYGRAD_TEST_MODE")); env.extend(local_tinygrad_worker_env(self.provider)); env.extend(optional_env("MVP_CPU_LINE_PROFILE")); env.extend(optional_env("MVP_CPU_LINE_PROFILE_INTERVAL_MS")); @@ -3132,6 +3141,7 @@ impl FrameArchive { }; let record = json!({ "arrival_seq":self.next_seq, + "arrival_unix_ms":benchmark_observability::unix_ms_now(), "source":source, "stream":stream.to_string(), "channel":channel, @@ -3252,6 +3262,7 @@ impl OrchDatastream { "status":status, "run_id":run_id, "node_id":node_id, + "benchmark":benchmark_observability::stamp("mvp-orchestrator"), "detail":detail, })) .expect("serialize orch bootstrap event"); @@ -3275,6 +3286,7 @@ impl OrchDatastream { "run_id":run_id, "node_id":node_id, "request_id":request_id, + "benchmark":benchmark_observability::stamp("mvp-orchestrator"), "detail":detail, })) .expect("serialize orch prompt event"); @@ -4042,6 +4054,21 @@ impl PipelinePromptRuntime { .unwrap_or(0); let final_text = self.final_text.clone(); let tokens_generated = self.generated_tokens.len() as u32; + orch_datastream.emit_prompt( + dashboard, + run_id, + node_id, + request_id, + "prompt_complete", + "ready", + json!({ + "event":"Done", + "terminal":true, + "tokens_generated":tokens_generated, + "elapsed_ms":elapsed_ms, + "final_text_bytes":final_text.len(), + }), + ); let _ = active.events.send(PromptEvent::Done { request_id, final_text, @@ -4985,7 +5012,6 @@ mod tests { "MVP_PIPELINE_STAGES", "MVP_STAGE_INDEX", "MVP_TOKEN_PROGRESS_EVERY", - "MVP_TINYGRAD_TEST_MODE", "MVP_TINYGRAD_WORKER", "MVP_TOKENIZER_LOCAL_PATH", "MVP_VASTAI_API_KEY", @@ -6070,28 +6096,37 @@ mod tests { .collect::>(); let _ = std::fs::remove_file(&path); + assert_eq!(records.len(), 2); + assert_eq!(records[0]["arrival_seq"], json!(0)); + assert!( + records[0]["arrival_unix_ms"] + .as_u64() + .is_some_and(|value| value > 0) + ); + assert_eq!(records[0]["source"], json!("orchestrator")); + assert_eq!(records[0]["stream"], json!("test-node#42")); + assert_eq!(records[0]["channel"], json!("stdout")); + assert_eq!(records[0]["channel_id"], json!(1)); + assert_eq!(records[0]["position"], json!(7)); assert_eq!( - records, - vec![ - json!({ - "arrival_seq":0, - "source":"orchestrator", - "stream":"test-node#42", - "channel":"stdout", - "channel_id":1, - "position":7, - "payload":{"encoding":"utf8","value":"hello λ"}, - }), - json!({ - "arrival_seq":1, - "source":"orchestrator", - "stream":"test-node#42", - "channel":"stderr", - "channel_id":2, - "position":8, - "payload":{"encoding":"bytes","value":[255,0,65]}, - }), - ] + records[0]["payload"], + json!({"encoding":"utf8","value":"hello λ"}) + ); + + assert_eq!(records[1]["arrival_seq"], json!(1)); + assert!( + records[1]["arrival_unix_ms"] + .as_u64() + .is_some_and(|value| value > 0) + ); + assert_eq!(records[1]["source"], json!("orchestrator")); + assert_eq!(records[1]["stream"], json!("test-node#42")); + assert_eq!(records[1]["channel"], json!("stderr")); + assert_eq!(records[1]["channel_id"], json!(2)); + assert_eq!(records[1]["position"], json!(8)); + assert_eq!( + records[1]["payload"], + json!({"encoding":"bytes","value":[255,0,65]}) ); } diff --git a/crates/mvp-system/src/bin/worker_node.rs b/crates/mvp-system/src/bin/worker_node.rs index aabb288..6e03d5a 100644 --- a/crates/mvp-system/src/bin/worker_node.rs +++ b/crates/mvp-system/src/bin/worker_node.rs @@ -29,6 +29,7 @@ use mvp_system::actors::node_agent::{ }; use mvp_system::actors::register_mvp_actor_codecs; use mvp_system::arena_manager as arena; +use mvp_system::benchmark_observability; use mvp_system::distribution_stack::DistributionRuntimeStack; use mvp_system::driver_pumps as driver_model; use mvp_system::edge_establisher as edge; @@ -72,10 +73,26 @@ fn node_event_payload( "run_id":config.run_id, "node_id":config.logical_node_id, "stage_index":config.stage_index, + "benchmark":benchmark_observability::stamp("mvp-worker-node"), "detail":detail, }) } +fn emit_stdio_datastream_frame(channel: &str, payload: &Value) -> Result<(), String> { + println!( + "{}", + json!({ + "mvp_stdio_event":1, + "kind":"datastream_frame", + "channel":channel, + "payload":payload, + }) + ); + std::io::stdout() + .flush() + .map_err(|e| format!("flush stdio datastream frame: {e}")) +} + fn emit_stdio_node_event( config: &DeploymentConfig, channel: &str, @@ -83,18 +100,8 @@ fn emit_stdio_node_event( status: &str, detail: Value, ) -> Result<(), String> { - println!( - "{}", - json!({ - "mvp_stdio_event":1, - "kind":"datastream_frame", - "channel":channel, - "payload":node_event_payload(config, phase, status, detail), - }) - ); - std::io::stdout() - .flush() - .map_err(|e| format!("flush stdio node event: {e}")) + let payload = node_event_payload(config, phase, status, detail); + emit_stdio_datastream_frame(channel, &payload) } fn emit_node_event( @@ -2956,6 +2963,12 @@ impl DeploymentConfig { .into_owned(), ), }; + let provider = env_string("MVP_NODE_PROVIDER", "process"); + let default_device = if provider == "process" { + "CPU" + } else { + DEFAULT_DEVICE + }; Ok(Self { run_id, logical_node_id, @@ -2966,7 +2979,7 @@ impl DeploymentConfig { debug_join_socket, relay_mode: relay.mode, worker_script: env_string("MVP_TINYGRAD_WORKER", DEFAULT_WORKER_SCRIPT), - device: env_string("DEV", DEFAULT_DEVICE), + device: env_string("DEV", default_device), model_id: env_string("MVP_MODEL_ID", DEFAULT_MODEL_ID), gguf_source: gguf_source_from_env(), tokenizer: tokenizer_from_env(), @@ -2991,6 +3004,9 @@ impl TinygradWorker { let mut child = Command::new("python3") .arg(&config.worker_script) .env("DEV", &config.device) + .env("MVP_RUN_ID", config.run_id.to_string()) + .env("MVP_LOGICAL_NODE_ID", config.logical_node_id.to_string()) + .env("MVP_STAGE_INDEX", config.stage_index.to_string()) .env("MVP_ARENA_FD", arena_fd.to_string()) .env("MVP_ARENA_BYTES", config.arena_bytes.to_string()) .stdin(Stdio::piped()) @@ -3518,6 +3534,8 @@ impl TinygradWorker { json!({"expected_event_type":expected,"channel":channel_name,"line_bytes":n,"worker_event_type":worker_event_type}), ); datastream.submit_text(channel, value.to_string()); + emit_stdio_datastream_frame(channel_name, &value) + .map_err(|e| format!("emit worker stdio datastream frame: {e}"))?; datastream.tick(); pump(); if worker_event_type == "WorkerFatal" { @@ -3651,6 +3669,23 @@ mod tests { DistributionRuntimeStack::new(DistNodeId([1; 32]), DistributedNodeConfig::default()) } + #[test] + fn benchmark_observability_node_event_payload_includes_stamp() { + let config = test_config(None); + let payload = node_event_payload(&config, "phase", "ready", json!({"ok": true})); + + assert_eq!(payload.get("run_id").and_then(Value::as_u64), Some(7)); + assert_eq!(payload.get("node_id").and_then(Value::as_u64), Some(11)); + assert_eq!(payload.get("stage_index").and_then(Value::as_u64), Some(3)); + assert_eq!( + payload + .get("benchmark") + .and_then(|benchmark| benchmark.get("schema")) + .and_then(Value::as_u64), + Some(1) + ); + } + #[test] fn debug_join_client_serializes_endpoint_from_stdin() { let secret = iroh::SecretKey::from_bytes(&[7; 32]); diff --git a/crates/mvp-system/src/lib.rs b/crates/mvp-system/src/lib.rs index 4b0422b..19f7f9c 100644 --- a/crates/mvp-system/src/lib.rs +++ b/crates/mvp-system/src/lib.rs @@ -3,6 +3,7 @@ extern crate self as mvp_system; pub mod actors; pub mod arena_manager; +pub mod benchmark_observability; pub mod bootstrap_datastream; pub mod config; pub mod dashboard_view; diff --git a/crates/mvp-system/tests/python_worker_protocol.rs b/crates/mvp-system/tests/python_worker_protocol.rs index b8a48bd..3e313b3 100755 --- a/crates/mvp-system/tests/python_worker_protocol.rs +++ b/crates/mvp-system/tests/python_worker_protocol.rs @@ -3,11 +3,8 @@ use std::io::{BufRead, BufReader, Write}; use std::process::{Child, ChildStdin, Command, Stdio}; -use mvp_system::arena_manager as arena; use serde_json::{Value, json}; -const HEADER_LEN: usize = 40; - struct WorkerProcess { child: Child, stdin: ChildStdin, @@ -15,22 +12,17 @@ struct WorkerProcess { } impl WorkerProcess { - fn spawn(arena: Option<(&arena::ArenaManager, u64)>) -> Self { + fn spawn() -> Self { let script = worker_script(); - let mut command = Command::new("python3"); - command + let mut child = Command::new("python3") .arg(&script) - .env("MVP_TINYGRAD_TEST_MODE", "1") .env("DEV", "CPU") + .env("MVP_RUN_ID", "9") + .env("MVP_LOGICAL_NODE_ID", "3") + .env("MVP_STAGE_INDEX", "2") .stdin(Stdio::piped()) .stdout(Stdio::piped()) - .stderr(Stdio::piped()); - if let Some((arena, arena_bytes)) = arena { - command - .env("MVP_ARENA_FD", arena.arena_fd().to_string()) - .env("MVP_ARENA_BYTES", arena_bytes.to_string()); - } - let mut child = command + .stderr(Stdio::piped()) .spawn() .unwrap_or_else(|e| panic!("spawn {}: {e}", script.display())); let stdin = child.stdin.take().expect("worker stdin"); @@ -42,9 +34,10 @@ impl WorkerProcess { } } - fn send_expect(&mut self, command: Value, expected_type: &str) -> Value { + fn send_collect_until(&mut self, command: Value, expected_type: &str) -> Vec { writeln!(self.stdin, "{command}").expect("write worker command"); self.stdin.flush().expect("flush worker command"); + let mut events = Vec::new(); loop { let mut line = String::new(); let read = self.stdout.read_line(&mut line).expect("read worker event"); @@ -55,8 +48,10 @@ impl WorkerProcess { actual_type, "WorkerFatal", "worker fatal while waiting for {expected_type}: {event}" ); - if actual_type == expected_type { - return event; + let done = actual_type == expected_type; + events.push(event); + if done { + return events; } } } @@ -70,301 +65,62 @@ impl Drop for WorkerProcess { } #[test] -fn test_mode_tokenizer_commands_round_trip_prompt_bytes_and_visible_tokens() { - let mut worker = WorkerProcess::spawn(None); - worker.send_expect( +fn benchmark_observability_worker_ready_includes_stamps_and_identity() { + if !tinygrad_available() { + eprintln!("skipping Python worker protocol check: tinygrad is unavailable"); + return; + } + + let mut worker = WorkerProcess::spawn(); + let events = worker.send_collect_until( json!({"type":"InitializeWorker","helper_abi_version":1,"backend":{"device":"CPU"}}), "WorkerReady", ); - worker.send_expect( - json!({ - "type":"LoadWeights", - "model_id":"test-model", - "gguf_source":{"LocalPath":"/tmp/not-used-in-test-mode.gguf"}, - "tokenizer":{"EmbeddedGguf":{}}, - "layer_start":0, - "layer_end_exclusive":1 - }), - "WeightsLoaded", - ); - let encoded = worker.send_expect( - json!({"type":"EncodePrompt","request_id":7,"prompt":"Hi!"}), - "PromptEncoded", - ); - assert_eq!(encoded.get("request_id").and_then(Value::as_u64), Some(7)); - assert_eq!( - encoded - .get("tokens") - .and_then(Value::as_array) - .expect("encoded tokens") + assert!( + events .iter() - .map(|value| value.as_u64().expect("token is u64")) - .collect::>(), - vec![72, 105, 33] + .any(|event| event.get("type").and_then(Value::as_str) == Some("TinygradImportStarted")), + "expected tinygrad import milestone in {events:?}" ); - - let decoded = worker.send_expect( - json!({"type":"DecodeTokens","request_id":8,"tokens":[72,105,33,6]}), - "TokensDecoded", + assert!( + events + .iter() + .any(|event| event.get("type").and_then(Value::as_str) == Some("WorkerReady")), + "expected WorkerReady in {events:?}" ); - assert_eq!(decoded.get("request_id").and_then(Value::as_u64), Some(8)); - assert_eq!( - decoded.get("text").and_then(Value::as_str), - Some("Hi!") - ); -} - -#[test] -fn test_mode_worker_executes_single_and_three_stage_mo01_flow_through_real_arena() { - let arena_bytes = 16 * 1024; - let mut arena = arena::ArenaManager::boot(arena::ArenaConfig { - node_id: arena::NodeId(1), - reservation_ceiling: arena_bytes, - base_alignment: 64, - }) - .expect("arena boots"); - let ingress = lease_ring(&mut arena, 1, 1024, 64); - let egress = lease_ring(&mut arena, 2, 1024, 64); - let mut worker = WorkerProcess::spawn(Some((&arena, arena_bytes))); - worker.send_expect( - json!({"type":"InitializeWorker","helper_abi_version":1,"backend":{"device":"CPU"}}), - "WorkerReady", - ); - - let single_token = run_stage( - &mut worker, - &arena, - &ingress, - &egress, - StageFixture { - stage_index: 0, - layer_start: 0, - layer_end_exclusive: 7, - final_stage: true, - }, - 900, - &[2, 3], - ); - assert_eq!( - single_token, - vec![6], - "N=1 stage should emit a token record" - ); - - let stage0 = run_stage( - &mut worker, - &arena, - &ingress, - &egress, - StageFixture { - stage_index: 0, - layer_start: 0, - layer_end_exclusive: 3, - final_stage: false, - }, - 901, - &[2, 3], - ); - assert_eq!(stage0, vec![8], "stage 0 activation fixture word"); - - let stage1 = run_stage( - &mut worker, - &arena, - &ingress, - &egress, - StageFixture { - stage_index: 1, - layer_start: 3, - layer_end_exclusive: 5, - final_stage: false, - }, - 902, - &stage0, - ); - assert_eq!(stage1, vec![17], "stage 1 activation fixture word"); - - let stage2 = run_stage( - &mut worker, - &arena, - &ingress, - &egress, - StageFixture { - stage_index: 2, - layer_start: 5, - layer_end_exclusive: 7, - final_stage: true, - }, - 903, - &stage1, - ); - assert_eq!(stage2, vec![6], "final stage should emit a token record"); -} - -#[derive(Clone, Copy)] -struct StageFixture { - stage_index: u32, - layer_start: u32, - layer_end_exclusive: u32, - final_stage: bool, -} - -fn run_stage( - worker: &mut WorkerProcess, - arena: &arena::ArenaManager, - ingress: &arena::RingLease, - egress: &arena::RingLease, - stage: StageFixture, - output_object_id: u64, - input_words: &[u32], -) -> Vec { - worker.send_expect( - json!({ - "type":"ConfigureRole", - "role_id":1, - "config":{ - "run_id":1, - "stage_index":stage.stage_index, - "layer_start":stage.layer_start, - "layer_end_exclusive":stage.layer_end_exclusive - } - }), - "RoleConfigured", - ); - worker.send_expect( - json!({ - "type":"LoadWeights", - "model_id":"test-model", - "gguf_source":{"LocalPath":"/tmp/not-used-in-test-mode.gguf"}, - "tokenizer":{"EmbeddedGguf":{}}, - "layer_start":stage.layer_start, - "layer_end_exclusive":stage.layer_end_exclusive - }), - "WeightsLoaded", - ); - install_ring(worker, ingress, 1, 10, "ingress", 4096, 4); - install_ring(worker, egress, 2, 11, "egress", 4096, 4); - - let input_payload = words_payload(input_words); - arena - .write_arena( - ingress.layout.data_offset, - &object_record(100 + u64::from(stage.stage_index), 0, 0, &input_payload), - ) - .expect("write ingress record"); - let loaded = worker.send_expect(json!({"type":"RingReadable","ring_id":1}), "ObjectLoaded"); - let handle_id = loaded - .get("handle_id") - .and_then(Value::as_u64) - .expect("handle id"); - worker.send_expect( - json!({ - "type":"ExecuteStep", - "role_id":1, - "step_id":u64::from(stage.stage_index) + 1, - "input_handle_id":handle_id, - "input_object_id":100 + u64::from(stage.stage_index), - "input_sequence":0, - "output_ring_id":2, - "output_object_id":output_object_id, - "output_sequence":0, - "final_stage":stage.final_stage - }), - "StepExecuted", - ); - let output = arena - .read_arena(egress.layout.data_offset, HEADER_LEN + 4) - .expect("read egress record"); - decode_record_words(&output, output_object_id) -} - -fn install_ring( - worker: &mut WorkerProcess, - lease: &arena::RingLease, - ring_id: u64, - edge_id: u64, - direction: &str, - max_extent: u64, - alignment: u32, -) { - worker.send_expect( - json!({ - "type":"InstallRing", - "ring_id":ring_id, - "edge_id":edge_id, - "port":direction, - "direction":direction, - "layout":{ - "data_offset":lease.layout.data_offset, - "data_bytes":lease.layout.data_bytes - }, - "object_spec":{ - "max_extent":max_extent, - "alignment":alignment - } - }), - "RingInstalled", - ); -} - -fn lease_ring( - arena: &mut arena::ArenaManager, - request_id: u64, - data_bytes: u64, - alignment: u64, -) -> arena::RingLease { - let events = arena.request(arena::ArenaRequest::LeaseRing(arena::LeaseRing { - request_id: arena::LeaseRequestId(request_id), - ring_spec: arena::RingSpec { - header_bytes: 64, - data_bytes, - alignment, - }, - })); - match events.into_iter().next().expect("arena event") { - arena::ArenaEvent::RingLeased { lease } => lease, - event => panic!("expected ring lease, got {event:?}"), + for event in &events { + assert_eq!( + event + .get("benchmark") + .and_then(|benchmark| benchmark.get("schema")) + .and_then(Value::as_u64), + Some(1), + "{event}" + ); + assert_eq!( + event.get("run_id").and_then(Value::as_u64), + Some(9), + "{event}" + ); + assert_eq!( + event.get("node_id").and_then(Value::as_u64), + Some(3), + "{event}" + ); + assert_eq!( + event.get("stage_index").and_then(Value::as_u64), + Some(2), + "{event}" + ); } } -fn object_record(object_id: u64, sequence: u64, flags: u32, payload: &[u8]) -> Vec { - let mut bytes = vec![0_u8; HEADER_LEN]; - bytes[0..4].copy_from_slice(b"MO01"); - bytes[4..6].copy_from_slice(&1_u16.to_le_bytes()); - bytes[6..8].copy_from_slice(&(HEADER_LEN as u16).to_le_bytes()); - bytes[8..16].copy_from_slice(&object_id.to_le_bytes()); - bytes[16..24].copy_from_slice(&sequence.to_le_bytes()); - bytes[24..32].copy_from_slice(&(payload.len() as u64).to_le_bytes()); - bytes[32..36].copy_from_slice(&flags.to_le_bytes()); - bytes.extend_from_slice(payload); - bytes -} - -fn words_payload(words: &[u32]) -> Vec { - words - .iter() - .flat_map(|word| word.to_le_bytes()) - .collect::>() -} - -fn decode_record_words(record: &[u8], expected_object_id: u64) -> Vec { - assert!(record.len() >= HEADER_LEN); - assert_eq!(&record[0..4], b"MO01"); - assert_eq!(u16::from_le_bytes(record[4..6].try_into().unwrap()), 1); - assert_eq!( - u16::from_le_bytes(record[6..8].try_into().unwrap()), - HEADER_LEN as u16 - ); - assert_eq!( - u64::from_le_bytes(record[8..16].try_into().unwrap()), - expected_object_id - ); - let extent = u64::from_le_bytes(record[24..32].try_into().unwrap()) as usize; - assert_eq!(extent % 4, 0); - record[HEADER_LEN..HEADER_LEN + extent] - .chunks_exact(4) - .map(|chunk| u32::from_le_bytes(chunk.try_into().unwrap())) - .collect() +fn tinygrad_available() -> bool { + Command::new("python3") + .args(["-c", "import tinygrad"]) + .status() + .is_ok_and(|status| status.success()) } fn worker_script() -> std::path::PathBuf { diff --git a/xtask/src/main.rs b/xtask/src/main.rs index 5311bcc..dc5469c 100644 --- a/xtask/src/main.rs +++ b/xtask/src/main.rs @@ -1,14 +1,17 @@ +use std::collections::BTreeMap; use std::fs; use std::io::{Read, Write}; use std::path::{Path, PathBuf}; use std::process::{Child, Command, ExitCode, ExitStatus, Stdio}; +use std::sync::LazyLock; +use std::sync::atomic::{AtomicU64, Ordering}; use std::thread; use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH}; #[cfg(target_os = "linux")] use std::os::unix::process::CommandExt; -use serde_json::Value; +use serde_json::{Value, json}; struct TestStep { label: &'static str, @@ -27,6 +30,7 @@ struct MvpChatCheckPaths { struct MvpChatCheckOutput { status: ExitStatus, + child_elapsed_ms: u64, stdout: String, stderr: String, timed_out: bool, @@ -142,9 +146,44 @@ fn run_mvp_chat(args: Vec) -> ExitCode { } else { args }; - command.args(forwarded); + let dump_log_path = explicit_dump_log_path_from_mvp_chat_args(&forwarded); + let run_id = run_id_from_mvp_chat_args(&forwarded); + let benchmark_target = dump_log_path.as_deref().zip(run_id); + if let Some((path, run_id)) = benchmark_target { + let event = xtask_mvp_chat_benchmark_event( + run_id, + "started", + json!({ + "program": "cargo", + "args": ["run", "--package", "mvp-system", "--bin", "mvp-chat", "--"], + }), + ); + if let Err(error) = + append_synthetic_benchmark_frame(path, "xtask-mvp-chat", "mvp.xtask.benchmark", event) + { + eprintln!("Failed to write mvp-chat benchmark frame: {error}"); + return ExitCode::from(1); + } + } + command.args(&forwarded); - match command.status() { + let status_result = command.status(); + if let Some((path, run_id)) = benchmark_target { + let (status, detail) = match &status_result { + Ok(status) if status.success() => ("ready", json!({"exit_status": status.to_string()})), + Ok(status) => ("failed", json!({"exit_status": status.to_string()})), + Err(error) => ("failed", json!({"error": error.to_string()})), + }; + let event = xtask_mvp_chat_benchmark_event(run_id, status, detail); + if let Err(error) = + append_synthetic_benchmark_frame(path, "xtask-mvp-chat", "mvp.xtask.benchmark", event) + { + eprintln!("Failed to write mvp-chat benchmark frame: {error}"); + return ExitCode::from(1); + } + } + + match status_result { Ok(status) if status.success() => ExitCode::SUCCESS, Ok(status) => ExitCode::from( status @@ -159,6 +198,124 @@ fn run_mvp_chat(args: Vec) -> ExitCode { } } +fn explicit_dump_log_path_from_mvp_chat_args(args: &[String]) -> Option { + let args = strip_leading_double_dash(args); + let mut index = 0; + while index < args.len() { + let arg = args[index].as_str(); + if let Some(path) = arg.strip_prefix("--dump-logs=") { + if !path.is_empty() { + return Some(PathBuf::from(path)); + } + } else if arg == "--dump-logs" { + if let Some(path) = args.get(index + 1) + && !path.starts_with("--") + { + return Some(PathBuf::from(path)); + } + } + index += 1; + } + None +} + +fn run_id_from_mvp_chat_args(args: &[String]) -> Option { + let args = strip_leading_double_dash(args); + let mut index = 0; + while index < args.len() { + if args[index] == "--run-id" { + return args + .get(index + 1) + .and_then(|value| value.parse::().ok()) + .filter(|value| *value != 0); + } + index += 1; + } + None +} + +fn strip_leading_double_dash(args: &[String]) -> &[String] { + if args.first().is_some_and(|arg| arg == "--") { + &args[1..] + } else { + args + } +} + +fn append_synthetic_benchmark_frame( + path: &Path, + source: &str, + channel: &str, + event: Value, +) -> Result<(), String> { + if let Some(parent) = path.parent() + && !parent.as_os_str().is_empty() + { + fs::create_dir_all(parent).map_err(|e| { + format!( + "create synthetic benchmark frame dir {}: {e}", + parent.display() + ) + })?; + } + let run_id = event + .get("run_id") + .and_then(Value::as_u64) + .ok_or_else(|| "synthetic benchmark event missing run_id".to_owned())?; + let inner = serde_json::to_string(&event) + .map_err(|e| format!("serialize synthetic benchmark event: {e}"))?; + let record = json!({ + "arrival_seq": 0, + "arrival_unix_ms": unix_ms_now(), + "source": source, + "stream": format!("xtask#{run_id}"), + "channel": channel, + "channel_id": 0, + "position": 0, + "payload": {"encoding": "utf8", "value": inner}, + }); + let mut file = fs::OpenOptions::new() + .create(true) + .append(true) + .open(path) + .map_err(|e| format!("open synthetic benchmark frame {}: {e}", path.display()))?; + let mut line = serde_json::to_vec(&record) + .map_err(|e| format!("serialize synthetic benchmark frame: {e}"))?; + line.push(b'\n'); + file.write_all(&line) + .map_err(|e| format!("write synthetic benchmark frame {}: {e}", path.display()))?; + file.flush() + .map_err(|e| format!("flush synthetic benchmark frame {}: {e}", path.display())) +} + +fn xtask_mvp_chat_benchmark_event(run_id: u64, status: &str, detail: Value) -> Value { + json!({ + "type": "XtaskBenchmark", + "phase": "cargo_run_mvp_chat", + "status": status, + "run_id": run_id, + "detail": detail, + "benchmark": xtask_benchmark_stamp(), + }) +} + +const XTASK_BENCHMARK_SCHEMA: u64 = 1; +static XTASK_BENCHMARK_START: LazyLock = LazyLock::new(Instant::now); +static XTASK_BENCHMARK_SEQ: AtomicU64 = AtomicU64::new(1); + +fn xtask_benchmark_stamp() -> Value { + let mono_ms = u64::try_from(XTASK_BENCHMARK_START.elapsed().as_millis()).unwrap_or(u64::MAX); + let seq = XTASK_BENCHMARK_SEQ.fetch_add(1, Ordering::Relaxed); + json!({ + "schema": XTASK_BENCHMARK_SCHEMA, + "component": "xtask", + "pid": std::process::id(), + "seq": seq, + "wall_unix_ms": unix_ms_now(), + "mono_ms": mono_ms, + }) +} + fn workspace_root() -> PathBuf { PathBuf::from(env!("CARGO_MANIFEST_DIR")) .parent() @@ -173,16 +330,29 @@ fn unique_temp_dir(prefix: &str) -> PathBuf { .duration_since(UNIX_EPOCH) .unwrap_or_default() .as_nanos(); - let root = - std::env::temp_dir().join(format!("{prefix}-{pid}-{timestamp_nanos}-{attempt}")); + let root = std::env::temp_dir().join(format!("{prefix}-{pid}-{timestamp_nanos}-{attempt}")); match fs::create_dir(&root) { Ok(()) => return root, Err(error) if error.kind() == std::io::ErrorKind::AlreadyExists => continue, - Err(error) => panic!("mvp-chat-check: create temp dir {}: {error}", root.display()), + Err(error) => panic!( + "mvp-chat-check: create temp dir {}: {error}", + root.display() + ), } } panic!("mvp-chat-check: could not allocate unique temp dir for prefix {prefix}"); } +fn unix_ms_now() -> u64 { + let millis = SystemTime::now() + .duration_since(UNIX_EPOCH) + .unwrap_or_default() + .as_millis(); + u64::try_from(millis).unwrap_or(u64::MAX) +} + +fn mvp_chat_check_run_id() -> u64 { + unix_ms_now().max(1) +} fn write_mvp_chat_check_paths(root: &Path) -> Result { if !root.is_dir() { @@ -218,8 +388,9 @@ fn run_mvp_chat_check() -> ExitCode { return ExitCode::from(1); } }; + let run_id = mvp_chat_check_run_id(); - let output = match run_mvp_chat_check_process(&workspace, &paths) { + let output = match run_mvp_chat_check_process(&workspace, &paths, run_id) { Ok(output) => output, Err(error) => return fail_mvp_chat_check(&error, &paths, "", "", None), }; @@ -265,14 +436,32 @@ fn run_mvp_chat_check() -> ExitCode { ); } }; - if let Err(error) = assert_dump_log_facts(&paths.dump_log) { - return fail_mvp_chat_check( - &error, - &paths, - &output.stdout, - &output.stderr, - Some(&output.status), - ); + let events = match assert_dump_log_facts(&paths.dump_log) { + Ok(events) => events, + Err(error) => { + return fail_mvp_chat_check( + &error, + &paths, + &output.stdout, + &output.stderr, + Some(&output.status), + ); + } + }; + let report = match build_benchmark_report(&events, output.child_elapsed_ms, run_id) { + Ok(report) => report, + Err(error) => { + return fail_mvp_chat_check( + &error, + &paths, + &output.stdout, + &output.stderr, + Some(&output.status), + ); + } + }; + for line in &report.lines { + println!("{line}"); } if let Err(error) = fs::remove_dir_all(&paths.root) { @@ -293,11 +482,13 @@ fn run_mvp_chat_check() -> ExitCode { fn run_mvp_chat_check_process( workspace: &Path, paths: &MvpChatCheckPaths, + run_id: u64, ) -> Result { let mut command = Command::new(cargo_bin()); command .current_dir(workspace) - .args(["mvp-chat", "--", "--cached-model"]) + .args(["mvp-chat", "--", "--cached-model", "--run-id"]) + .arg(run_id.to_string()) .arg(format!("--dump-logs={}", paths.dump_log.display())) .stdin(Stdio::piped()) .stdout(Stdio::piped()) @@ -315,6 +506,7 @@ fn run_mvp_chat_check_process( }); } + let child_started = Instant::now(); let mut child = command .spawn() .map_err(|e| format!("mvp-chat-check: spawn cargo mvp-chat: {e}"))?; @@ -345,10 +537,12 @@ fn run_mvp_chat_check_process( } else { wait_mvp_chat_check_child(&mut child)? }; + let child_elapsed_ms = duration_ms_u64(child_started.elapsed()); let stdout = join_reader(stdout_reader, "stdout")?; let stderr = join_reader(stderr_reader, "stderr")?; Ok(MvpChatCheckOutput { + child_elapsed_ms, status, stdout, stderr, @@ -383,9 +577,7 @@ fn terminate_mvp_chat_child(child: &mut Child) -> Result { Ok(Some(status)) => return Ok(status), Ok(None) => thread::sleep(Duration::from_millis(MVP_CHAT_CHECK_POLL_MS)), Err(error) => { - return Err(format!( - "mvp-chat-check: poll child after SIGTERM: {error}" - )); + return Err(format!("mvp-chat-check: poll child after SIGTERM: {error}")); } } } @@ -519,29 +711,299 @@ fn find_stdout_marker( stdout[start..] .find(marker) .map(|offset| start + offset) - .ok_or_else(|| { - format!("mvp-chat-check: missing {label} marker for prompt cycle {cycle}") - }) + .ok_or_else(|| format!("mvp-chat-check: missing {label} marker for prompt cycle {cycle}")) +} + +#[derive(Clone)] +struct DumpLogEvent { + source: String, + channel: String, + arrival_unix_ms: Option, + event: Value, +} + +#[derive(Clone)] +struct BenchmarkPoint { + component: Option, + wall_unix_ms: Option, + mono_ms: Option, + arrival_unix_ms: Option, +} + +#[derive(Debug)] +struct BenchmarkReport { + lines: Vec, } #[derive(Default)] -struct DumpLogFacts { - chat_config_ready: bool, - prepare_runtime_ready: bool, - prompt_rpc_ready: bool, - orch_iroh_driver_ready: bool, - node_iroh_driver_ready: bool, - node_worker_initialize_ready: bool, - orch_weights_loaded_ready: bool, - response_text_1: bool, - request_completed_1: bool, - response_text_2: bool, - request_completed_2: bool, - shutdown_requested: bool, - orchestrator_stopped: bool, +struct BenchmarkFacts { + spans: BTreeMap<(String, String, String, String), BenchmarkPoint>, + prompts: BTreeMap, } -fn assert_dump_log_facts(path: &Path) -> Result<(), String> { +#[derive(Default)] +struct PromptBenchmarkFacts { + request_id: u64, + chat_submitted: Option, + chat_completed: Option, + worker_started: Option, + worker_completed: Option, + encode_started: Option, + encode_ready: Option, + decode_started: Option, + first_token_ready: Option, + decode_ready: Option, + text_decode_started: Option, + text_decode_ready: Option, + tokens_generated: Option, +} + +impl PromptBenchmarkFacts { + fn new(request_id: u64) -> Self { + Self { + request_id, + ..Self::default() + } + } +} + +impl BenchmarkFacts { + fn from_events(events: &[DumpLogEvent], run_id: u64) -> Self { + let mut facts = Self::default(); + for record in events { + if !event_matches_run_id(&record.event, run_id) { + continue; + } + let point = BenchmarkPoint::from_record(record); + let event_type = record.event.get("type").and_then(Value::as_str); + let phase = record.event.get("phase").and_then(Value::as_str); + let status = record.event.get("status").and_then(Value::as_str); + if let (Some(event_type), Some(phase), Some(status)) = (event_type, phase, status) { + facts + .spans + .entry(( + record.channel.clone(), + event_type.to_owned(), + phase.to_owned(), + status.to_owned(), + )) + .or_insert_with(|| point.clone()); + } + + match (record.channel.as_str(), event_type, phase, status) { + ( + "mvp.chat.prompt", + Some("ChatProgress"), + Some("prompt_submitted"), + Some("ready"), + ) => { + if let Some(request_id) = dump_log_request_id(&record.event) { + facts.prompt_mut(request_id).chat_submitted = Some(point); + } + } + ( + "mvp.chat.prompt", + Some("ChatProgress"), + Some("request_completed"), + Some("ready"), + ) => { + if let Some(request_id) = dump_log_request_id(&record.event) { + let prompt = facts.prompt_mut(request_id); + prompt.chat_completed = Some(point); + if prompt.tokens_generated.is_none() { + prompt.tokens_generated = record + .event + .get("detail") + .and_then(|detail| detail.get("tokens_generated")) + .and_then(Value::as_u64); + } + } + } + ("mvp.worker.prompt", Some(worker_type), _, _) => { + let Some(request_id) = benchmark_request_id(&record.event) else { + continue; + }; + let prompt = facts.prompt_mut(request_id); + match worker_type { + "PromptStarted" => prompt.worker_started = Some(point), + "PromptCompleted" => { + prompt.worker_completed = Some(point); + if let Some(tokens) = record + .event + .get("generated_tokens") + .and_then(Value::as_array) + { + prompt.tokens_generated = Some(tokens.len() as u64); + } + } + "PromptEncodeStarted" => prompt.encode_started = Some(point), + "PromptEncodeReady" => prompt.encode_ready = Some(point), + "DecodeStarted" => prompt.decode_started = Some(point), + "FirstTokenReady" => prompt.first_token_ready = Some(point), + "DecodeReady" => prompt.decode_ready = Some(point), + "TextDecodeStarted" => prompt.text_decode_started = Some(point), + "TextDecodeReady" => prompt.text_decode_ready = Some(point), + _ => {} + } + } + ( + "mvp.orch.prompt", + Some("OrchPromptEvent"), + Some("pipeline_tokenizer_encode"), + Some("started"), + ) => { + if let Some(request_id) = benchmark_request_id(&record.event) { + let prompt = facts.prompt_mut(request_id); + if prompt.encode_started.is_none() { + prompt.encode_started = Some(point); + } + } + } + ( + "mvp.orch.prompt", + Some("OrchPromptEvent"), + Some("pipeline_tokenizer_encode"), + Some("ready"), + ) => { + if let Some(request_id) = benchmark_request_id(&record.event) { + facts.prompt_mut(request_id).encode_ready = Some(point); + } + } + ( + "mvp.orch.prompt", + Some("OrchPromptEvent"), + Some("pipeline_token_in"), + Some("started"), + ) => { + if let Some(request_id) = benchmark_request_id(&record.event) { + let prompt = facts.prompt_mut(request_id); + if prompt.decode_started.is_none() { + prompt.decode_started = Some(point); + } + } + } + ( + "mvp.orch.prompt", + Some("OrchPromptEvent"), + Some("pipeline_token_out"), + Some("observed"), + ) => { + if let Some(request_id) = benchmark_request_id(&record.event) { + let prompt = facts.prompt_mut(request_id); + if prompt.first_token_ready.is_none() { + prompt.first_token_ready = Some(point.clone()); + } + prompt.decode_ready = Some(point); + prompt.tokens_generated = + Some(prompt.tokens_generated.unwrap_or(0).saturating_add(1)); + } + } + ( + "mvp.orch.prompt", + Some("OrchPromptEvent"), + Some("pipeline_tokenizer_decode"), + Some("started"), + ) => { + if let Some(request_id) = benchmark_request_id(&record.event) { + let prompt = facts.prompt_mut(request_id); + if prompt.text_decode_started.is_none() { + prompt.text_decode_started = Some(point); + } + } + } + ( + "mvp.orch.prompt", + Some("OrchPromptEvent"), + Some("pipeline_tokenizer_decode"), + Some("ready"), + ) => { + if let Some(request_id) = benchmark_request_id(&record.event) { + facts.prompt_mut(request_id).text_decode_ready = Some(point); + } + } + ( + "mvp.orch.prompt", + Some("OrchPromptEvent"), + Some("prompt_complete"), + Some("ready"), + ) => { + if let Some(request_id) = benchmark_request_id(&record.event) { + facts.prompt_mut(request_id).worker_completed = Some(point); + } + } + _ => {} + } + } + facts + } + + fn prompt_mut(&mut self, request_id: u64) -> &mut PromptBenchmarkFacts { + self.prompts + .entry(request_id) + .or_insert_with(|| PromptBenchmarkFacts::new(request_id)) + } + + fn span_point( + &self, + channel: &str, + event_type: &str, + phase: &str, + status: &str, + ) -> Option<&BenchmarkPoint> { + self.spans.get(&( + channel.to_owned(), + event_type.to_owned(), + phase.to_owned(), + status.to_owned(), + )) + } + + fn require_span( + &self, + channel: &str, + event_type: &str, + phase: &str, + status: &str, + ) -> Result<(), String> { + if self + .span_point(channel, event_type, phase, status) + .is_some() + { + Ok(()) + } else { + Err(missing_benchmark_event(format!( + "{channel}/{event_type}/{phase}/{status}" + ))) + } + } +} + +impl BenchmarkPoint { + fn from_record(record: &DumpLogEvent) -> Self { + let benchmark = record.event.get("benchmark"); + Self { + component: benchmark + .and_then(|value| value.get("component")) + .and_then(Value::as_str) + .map(str::to_owned), + wall_unix_ms: benchmark + .and_then(|value| value.get("wall_unix_ms")) + .and_then(Value::as_u64), + mono_ms: benchmark + .and_then(|value| value.get("mono_ms")) + .and_then(Value::as_u64), + arrival_unix_ms: record.arrival_unix_ms, + } + } +} + +#[derive(Clone, Copy)] +struct DurationRender { + value_ms: Option, + clock_skew: bool, +} + +fn parse_dump_log_events(path: &Path) -> Result, String> { let content = fs::read_to_string(path) .map_err(|e| format!("mvp-chat-check: read dump log {}: {e}", path.display()))?; if content.lines().next().is_none() { @@ -551,7 +1013,7 @@ fn assert_dump_log_facts(path: &Path) -> Result<(), String> { )); } - let mut facts = DumpLogFacts::default(); + let mut events = Vec::new(); for (line_index, line) in content.lines().enumerate() { if line.trim().is_empty() { continue; @@ -571,47 +1033,424 @@ fn assert_dump_log_facts(path: &Path) -> Result<(), String> { "mvp-chat-check: dump log line {} missing channel", line_index + 1 ) - })?; + })? + .to_owned(); let payload = outer.get("payload").ok_or_else(|| { format!( "mvp-chat-check: dump log line {} missing payload", line_index + 1 ) })?; - if payload.get("encoding").and_then(Value::as_str) != Some("utf8") { + let Some(inner_text) = dump_log_inner_payload_text(payload, line_index)? else { continue; - } - let inner_text = payload.get("value").and_then(Value::as_str).ok_or_else(|| { - format!( - "mvp-chat-check: dump log line {} missing utf8 payload value", - line_index + 1 - ) - })?; + }; let event: Value = serde_json::from_str(inner_text).map_err(|e| { format!( "mvp-chat-check: parse inner event on dump log line {}: {e}", line_index + 1 ) })?; - record_dump_log_event(channel, &event, &mut facts)?; + events.push(DumpLogEvent { + source: outer + .get("source") + .and_then(Value::as_str) + .unwrap_or_default() + .to_owned(), + channel, + arrival_unix_ms: outer.get("arrival_unix_ms").and_then(Value::as_u64), + event, + }); + } + Ok(events) +} + +fn dump_log_inner_payload_text<'a>( + payload: &'a Value, + line_index: usize, +) -> Result, String> { + if let Some(text) = payload.as_str() { + return Ok(Some(text)); + } + if payload.get("encoding").and_then(Value::as_str) != Some("utf8") { + return Ok(None); + } + payload + .get("value") + .and_then(Value::as_str) + .map(Some) + .ok_or_else(|| { + format!( + "mvp-chat-check: dump log line {} missing utf8 payload value", + line_index + 1 + ) + }) +} + +fn assert_dump_log_facts(path: &Path) -> Result, String> { + let events = parse_dump_log_events(path)?; + let mut facts = DumpLogFacts::default(); + for record in &events { + let _source = record.source.as_str(); + record_dump_log_event(&record.channel, &record.event, &mut facts)?; } require_dump_log_fact(facts.chat_config_ready, "config ready")?; require_dump_log_fact(facts.prepare_runtime_ready, "prepare_runtime ready")?; require_dump_log_fact(facts.prompt_rpc_ready, "prompt_rpc ready")?; - require_dump_log_fact(facts.orch_iroh_driver_ready, "OrchBootstrap iroh_driver ready")?; + require_dump_log_fact( + facts.orch_iroh_driver_ready, + "OrchBootstrap iroh_driver ready", + )?; require_dump_log_fact(facts.node_iroh_driver_ready, "NodeEvent iroh_driver ready")?; require_dump_log_fact( facts.node_worker_initialize_ready, "NodeEvent worker_initialize ready", )?; - require_dump_log_fact(facts.orch_weights_loaded_ready, "OrchBootstrap weights_loaded ready")?; + require_dump_log_fact( + facts.orch_weights_loaded_ready, + "OrchBootstrap weights_loaded ready", + )?; require_dump_log_fact(facts.response_text_1, "response_text request_id=1")?; require_dump_log_fact(facts.request_completed_1, "request_completed request_id=1")?; require_dump_log_fact(facts.response_text_2, "response_text request_id=2")?; require_dump_log_fact(facts.request_completed_2, "request_completed request_id=2")?; require_dump_log_fact(facts.shutdown_requested, "shutdown requested")?; - require_dump_log_fact(facts.orchestrator_stopped, "orchestrator_process stopped") + require_dump_log_fact(facts.orchestrator_stopped, "orchestrator_process stopped")?; + Ok(events) +} + +fn build_benchmark_report( + events: &[DumpLogEvent], + child_elapsed_ms: u64, + run_id: u64, +) -> Result { + let facts = BenchmarkFacts::from_events(events, run_id); + facts.require_span( + "mvp.xtask.benchmark", + "XtaskBenchmark", + "cargo_run_mvp_chat", + "started", + )?; + facts.require_span( + "mvp.xtask.benchmark", + "XtaskBenchmark", + "cargo_run_mvp_chat", + "ready", + )?; + facts.require_span( + "mvp.chat.runtime", + "ChatProgress", + "prepare_runtime", + "started", + )?; + facts.require_span( + "mvp.chat.runtime", + "ChatProgress", + "prepare_runtime", + "ready", + )?; + facts.require_span( + "mvp.chat.runtime", + "ChatProgress", + "ensure_orch_binary", + "ready", + )?; + facts.require_span( + "mvp.chat.runtime", + "ChatProgress", + "ensure_worker_binary", + "ready", + )?; + facts.require_span( + "mvp.orch.bootstrap", + "OrchBootstrap", + "weights_loaded", + "ready", + )?; + facts.require_span("mvp.chat.runtime", "ChatProgress", "prompt_rpc", "ready")?; + + for request_id in 1..=2 { + let prompt = facts.prompts.get(&request_id).ok_or_else(|| { + missing_benchmark_event(format!( + "mvp.chat.prompt/ChatProgress/request_completed/ready request_id={request_id}" + )) + })?; + require_prompt_point( + prompt.chat_completed.is_some(), + "mvp.chat.prompt/ChatProgress/request_completed/ready", + request_id, + )?; + require_prompt_point( + prompt.encode_ready.is_some(), + "mvp.worker.prompt/PromptEncodeReady or mvp.orch.prompt/pipeline_tokenizer_encode/ready", + request_id, + )?; + require_prompt_point( + prompt.decode_started.is_some(), + "mvp.worker.prompt/DecodeStarted or mvp.orch.prompt/pipeline_token_in/started", + request_id, + )?; + require_prompt_point( + prompt.first_token_ready.is_some(), + "mvp.worker.prompt/FirstTokenReady or mvp.orch.prompt/pipeline_token_out/observed", + request_id, + )?; + require_prompt_point( + prompt.decode_ready.is_some(), + "mvp.worker.prompt/DecodeReady or mvp.orch.prompt/pipeline_token_out/observed", + request_id, + )?; + require_prompt_point( + prompt.text_decode_ready.is_some(), + "mvp.worker.prompt/TextDecodeReady or mvp.orch.prompt/pipeline_tokenizer_decode/ready", + request_id, + )?; + } + + let cargo_run = duration_between( + facts.span_point( + "mvp.xtask.benchmark", + "XtaskBenchmark", + "cargo_run_mvp_chat", + "started", + ), + facts.span_point( + "mvp.xtask.benchmark", + "XtaskBenchmark", + "cargo_run_mvp_chat", + "ready", + ), + ); + let prepare_runtime = duration_between( + facts.span_point( + "mvp.chat.runtime", + "ChatProgress", + "prepare_runtime", + "started", + ), + facts.span_point( + "mvp.chat.runtime", + "ChatProgress", + "prepare_runtime", + "ready", + ), + ); + let standup_start = facts.span_point( + "mvp.chat.runtime", + "ChatProgress", + "prepare_runtime", + "ready", + ); + let weights_loaded = facts.span_point( + "mvp.orch.bootstrap", + "OrchBootstrap", + "weights_loaded", + "ready", + ); + let prompt_rpc = facts.span_point("mvp.chat.runtime", "ChatProgress", "prompt_rpc", "ready"); + let standup_to_weights = duration_between(standup_start, weights_loaded); + let standup_to_prompt_rpc = duration_between(standup_start, prompt_rpc); + + let mut lines = vec![ + format!("mvp-chat-check: benchmark: run_id={run_id}"), + format!("mvp-chat-check: benchmark total_child_ms={child_elapsed_ms}"), + benchmark_span_line("cargo_run_mvp_chat_ms", cargo_run), + benchmark_span_line("prepare_runtime_ms", prepare_runtime), + benchmark_span_line("standup_to_weights_loaded_ms", standup_to_weights), + benchmark_span_line("standup_to_prompt_rpc_ms", standup_to_prompt_rpc), + ]; + for request_id in 1..=2 { + let prompt = facts + .prompts + .get(&request_id) + .expect("prompt facts were required above"); + lines.push(prompt_benchmark_line(prompt)); + } + Ok(BenchmarkReport { lines }) +} + +fn benchmark_span_line(name: &str, duration: DurationRender) -> String { + let mut line = format!( + "mvp-chat-check: benchmark {name}={}", + render_duration_value(duration) + ); + if duration.clock_skew { + line.push_str(" clock_skew=true"); + } + line +} + +fn prompt_benchmark_line(prompt: &PromptBenchmarkFacts) -> String { + let roundtrip = duration_between( + prompt.chat_submitted.as_ref(), + prompt.chat_completed.as_ref(), + ); + let worker_start = prompt + .worker_started + .as_ref() + .or(prompt.encode_started.as_ref()); + let worker_end = prompt + .worker_completed + .as_ref() + .or(prompt.chat_completed.as_ref()); + let worker_total = duration_between(worker_start, worker_end); + let encode = duration_between(prompt.encode_started.as_ref(), prompt.encode_ready.as_ref()); + let first_token = duration_between( + prompt.decode_started.as_ref(), + prompt.first_token_ready.as_ref(), + ); + let decode = duration_between(prompt.decode_started.as_ref(), prompt.decode_ready.as_ref()); + let text_decode = duration_between( + prompt.text_decode_started.as_ref(), + prompt.text_decode_ready.as_ref(), + ); + let mut fields = vec![ + format!("roundtrip_ms={}", render_duration_value(roundtrip)), + format!("worker_total_ms={}", render_duration_value(worker_total)), + format!("encode_ms={}", render_duration_value(encode)), + format!("first_token_ms={}", render_duration_value(first_token)), + format!("decode_ms={}", render_duration_value(decode)), + format!("text_decode_ms={}", render_duration_value(text_decode)), + format!( + "tokens_generated={}", + prompt + .tokens_generated + .map(|value| value.to_string()) + .unwrap_or_else(|| "unavailable".to_owned()) + ), + format!( + "tokens_per_sec={}", + tokens_per_sec(prompt.tokens_generated, decode.value_ms) + ), + ]; + if [ + roundtrip, + worker_total, + encode, + first_token, + decode, + text_decode, + ] + .iter() + .any(|duration| duration.clock_skew) + { + fields.push("clock_skew=true".to_owned()); + } + format!( + "mvp-chat-check: benchmark prompt {} {}", + prompt.request_id, + fields.join(" ") + ) +} + +fn tokens_per_sec(tokens_generated: Option, decode_ms: Option) -> String { + match (tokens_generated, decode_ms) { + (Some(tokens), Some(ms)) if tokens > 0 && ms > 0 => { + format!("{:.2}", tokens as f64 / (ms as f64 / 1000.0)) + } + _ => "unavailable".to_owned(), + } +} + +fn render_duration_value(duration: DurationRender) -> String { + duration + .value_ms + .map(|value| value.to_string()) + .unwrap_or_else(|| "unavailable".to_owned()) +} + +fn duration_between( + start: Option<&BenchmarkPoint>, + end: Option<&BenchmarkPoint>, +) -> DurationRender { + let (Some(start), Some(end)) = (start, end) else { + return DurationRender { + value_ms: None, + clock_skew: false, + }; + }; + if start.component.is_some() + && start.component == end.component + && let (Some(start_ms), Some(end_ms)) = (start.mono_ms, end.mono_ms) + { + return DurationRender { + value_ms: end_ms.checked_sub(start_ms), + clock_skew: end_ms < start_ms, + }; + } + if let (Some(start_ms), Some(end_ms)) = (start.wall_unix_ms, end.wall_unix_ms) { + if let Some(value_ms) = end_ms.checked_sub(start_ms) { + return DurationRender { + value_ms: Some(value_ms), + clock_skew: false, + }; + } + return DurationRender { + value_ms: match (start.arrival_unix_ms, end.arrival_unix_ms) { + (Some(start_arrival), Some(end_arrival)) => end_arrival.checked_sub(start_arrival), + _ => None, + }, + clock_skew: true, + }; + } + DurationRender { + value_ms: match (start.arrival_unix_ms, end.arrival_unix_ms) { + (Some(start_arrival), Some(end_arrival)) => end_arrival.checked_sub(start_arrival), + _ => None, + }, + clock_skew: false, + } +} + +fn duration_ms_u64(duration: Duration) -> u64 { + u64::try_from(duration.as_millis()).unwrap_or(u64::MAX) +} + +fn event_matches_run_id(event: &Value, run_id: u64) -> bool { + event + .get("run_id") + .and_then(Value::as_u64) + .is_none_or(|event_run_id| event_run_id == run_id) +} + +fn benchmark_request_id(event: &Value) -> Option { + event.get("request_id").and_then(Value::as_u64).or_else(|| { + event + .get("detail") + .and_then(|detail| detail.get("request_id")) + .and_then(Value::as_u64) + }) +} + +fn require_prompt_point(found: bool, event: &str, request_id: u64) -> Result<(), String> { + if found { + Ok(()) + } else { + Err(missing_benchmark_event(format!( + "{event} request_id={request_id}" + ))) + } +} + +fn missing_benchmark_event(event: String) -> String { + format!("mvp-chat-check: missing benchmark event {event}") +} + +#[derive(Default)] +struct DumpLogFacts { + chat_config_ready: bool, + prepare_runtime_ready: bool, + prompt_rpc_ready: bool, + orch_iroh_driver_ready: bool, + node_iroh_driver_ready: bool, + node_worker_initialize_ready: bool, + orch_weights_loaded_ready: bool, + response_text_1: bool, + request_completed_1: bool, + response_text_2: bool, + request_completed_2: bool, + shutdown_requested: bool, + orchestrator_stopped: bool, } fn record_dump_log_event( @@ -698,6 +1537,554 @@ fn require_dump_log_fact(found: bool, fact: &str) -> Result<(), String> { } } +#[cfg(test)] +mod tests { + use super::*; + + static NEXT_TEST_FILE: AtomicU64 = AtomicU64::new(1); + + fn strings(values: &[&str]) -> Vec { + values.iter().map(|value| (*value).to_owned()).collect() + } + + fn temp_path(label: &str) -> PathBuf { + let id = NEXT_TEST_FILE.fetch_add(1, Ordering::Relaxed); + std::env::temp_dir().join(format!( + "xtask-benchmark-observability-{label}-{}-{id}.ndjson", + std::process::id() + )) + } + + fn benchmark_observability_archive_record( + channel: &str, + payload: Value, + arrival_unix_ms: u64, + ) -> String { + json!({ + "arrival_seq": 0, + "arrival_unix_ms": arrival_unix_ms, + "source": "unit-test", + "stream": "unit#9", + "channel": channel, + "channel_id": 0, + "position": 0, + "payload": payload, + }) + .to_string() + } + + fn write_dump_log(label: &str, lines: Vec) -> PathBuf { + let path = temp_path(label); + fs::write(&path, format!("{}\n", lines.join("\n"))).expect("write synthetic dump log"); + path + } + + fn parse_synthetic_events( + label: &str, + events: Vec<(&'static str, Value)>, + ) -> Vec { + let lines = events + .into_iter() + .enumerate() + .map(|(index, (channel, event))| { + benchmark_observability_archive_record( + channel, + json!({"encoding": "utf8", "value": event.to_string()}), + 10_000 + index as u64, + ) + }) + .collect::>(); + let path = write_dump_log(label, lines); + let parsed = parse_dump_log_events(&path).expect("parse synthetic dump log"); + let _ = fs::remove_file(path); + parsed + } + + fn benchmark(component: &str, wall_unix_ms: u64, mono_ms: u64) -> Value { + json!({ + "schema": 1, + "component": component, + "pid": 1, + "seq": wall_unix_ms, + "wall_unix_ms": wall_unix_ms, + "mono_ms": mono_ms, + }) + } + + fn stamped(mut event: Value, component: &str, wall_unix_ms: u64, mono_ms: u64) -> Value { + let object = event.as_object_mut().expect("event object"); + object.insert( + "benchmark".to_owned(), + benchmark(component, wall_unix_ms, mono_ms), + ); + event + } + + fn chat_span(phase: &str, status: &str, wall_unix_ms: u64, mono_ms: u64) -> Value { + stamped( + json!({ + "type": "ChatProgress", + "phase": phase, + "status": status, + "run_id": 9, + "detail": {}, + }), + "mvp-chat", + wall_unix_ms, + mono_ms, + ) + } + + fn prompt_chat_span( + phase: &str, + status: &str, + request_id: u64, + wall_unix_ms: u64, + mono_ms: u64, + ) -> Value { + stamped( + json!({ + "type": "ChatProgress", + "phase": phase, + "status": status, + "run_id": 9, + "detail": { + "request_id": request_id, + "tokens_generated": 3, + "elapsed_ms": 50, + "final_text_bytes": 5, + "response_started": true, + }, + }), + "mvp-chat", + wall_unix_ms, + mono_ms, + ) + } + + fn worker_prompt_event( + event_type: &str, + request_id: u64, + wall_unix_ms: u64, + mono_ms: u64, + ) -> Value { + let mut event = stamped( + json!({ + "type": event_type, + "run_id": 9, + "node_id": 3, + "stage_index": 2, + "request_id": request_id, + }), + "tinygrad-worker", + wall_unix_ms, + mono_ms, + ); + if event_type == "PromptCompleted" { + event + .as_object_mut() + .expect("event object") + .insert("generated_tokens".to_owned(), json!([1, 2, 3])); + } + event + } + + fn pipeline_prompt_event( + phase: &str, + status: &str, + request_id: u64, + wall_unix_ms: u64, + mono_ms: u64, + detail: Value, + ) -> Value { + stamped( + json!({ + "type": "OrchPromptEvent", + "phase": phase, + "status": status, + "run_id": 9, + "node_id": 1, + "request_id": request_id, + "detail": detail, + }), + "mvp-orchestrator", + wall_unix_ms, + mono_ms, + ) + } + + fn benchmark_report_events(include_first_token: bool) -> Vec<(&'static str, Value)> { + let mut events = benchmark_report_base_events(); + for request_id in 1..=2 { + let base = 1_200 + request_id * 100; + events.push(( + "mvp.chat.prompt", + prompt_chat_span("prompt_submitted", "ready", request_id, base, base - 1_000), + )); + events.push(( + "mvp.worker.prompt", + worker_prompt_event("PromptStarted", request_id, base + 5, request_id * 100), + )); + events.push(( + "mvp.worker.prompt", + worker_prompt_event( + "PromptEncodeStarted", + request_id, + base + 10, + request_id * 100 + 10, + ), + )); + events.push(( + "mvp.worker.prompt", + worker_prompt_event( + "PromptEncodeReady", + request_id, + base + 15, + request_id * 100 + 15, + ), + )); + events.push(( + "mvp.worker.prompt", + worker_prompt_event( + "DecodeStarted", + request_id, + base + 20, + request_id * 100 + 20, + ), + )); + if include_first_token { + events.push(( + "mvp.worker.prompt", + worker_prompt_event( + "FirstTokenReady", + request_id, + base + 25, + request_id * 100 + 25, + ), + )); + } + events.push(( + "mvp.worker.prompt", + worker_prompt_event("DecodeReady", request_id, base + 40, request_id * 100 + 40), + )); + events.push(( + "mvp.worker.prompt", + worker_prompt_event( + "TextDecodeStarted", + request_id, + base + 41, + request_id * 100 + 41, + ), + )); + events.push(( + "mvp.worker.prompt", + worker_prompt_event( + "TextDecodeReady", + request_id, + base + 45, + request_id * 100 + 45, + ), + )); + events.push(( + "mvp.worker.prompt", + worker_prompt_event( + "PromptCompleted", + request_id, + base + 50, + request_id * 100 + 50, + ), + )); + events.push(( + "mvp.chat.prompt", + prompt_chat_span( + "request_completed", + "ready", + request_id, + base + 60, + base - 940, + ), + )); + } + events + } + + fn pipeline_benchmark_report_events(include_first_token: bool) -> Vec<(&'static str, Value)> { + let mut events = benchmark_report_base_events(); + for request_id in 1..=2 { + let base = 1_200 + request_id * 100; + events.push(( + "mvp.chat.prompt", + prompt_chat_span("prompt_submitted", "ready", request_id, base, base - 1_000), + )); + events.push(( + "mvp.orch.prompt", + pipeline_prompt_event( + "pipeline_tokenizer_encode", + "started", + request_id, + base + 5, + base + 5, + json!({"prompt_bytes":4}), + ), + )); + events.push(( + "mvp.orch.prompt", + pipeline_prompt_event( + "pipeline_tokenizer_encode", + "ready", + request_id, + base + 10, + base + 10, + json!({"tokens":4}), + ), + )); + events.push(( + "mvp.orch.prompt", + pipeline_prompt_event( + "pipeline_token_in", + "started", + request_id, + base + 12, + base + 12, + json!({"tokens":4}), + ), + )); + if include_first_token { + events.push(( + "mvp.orch.prompt", + pipeline_prompt_event( + "pipeline_token_out", + "observed", + request_id, + base + 20, + base + 20, + json!({"sequence":0,"token_id":7}), + ), + )); + } + events.push(( + "mvp.orch.prompt", + pipeline_prompt_event( + "pipeline_tokenizer_decode", + "started", + request_id, + base + 21, + base + 21, + json!({"token_id":7}), + ), + )); + events.push(( + "mvp.orch.prompt", + pipeline_prompt_event( + "pipeline_tokenizer_decode", + "ready", + request_id, + base + 25, + base + 25, + json!({"text_bytes":1}), + ), + )); + events.push(( + "mvp.chat.prompt", + prompt_chat_span( + "request_completed", + "ready", + request_id, + base + 60, + base - 940, + ), + )); + } + events + } + + fn benchmark_report_base_events() -> Vec<(&'static str, Value)> { + vec![ + ( + "mvp.xtask.benchmark", + stamped( + json!({ + "type": "XtaskBenchmark", + "phase": "cargo_run_mvp_chat", + "status": "started", + "run_id": 9, + "detail": {}, + }), + "xtask", + 1_000, + 0, + ), + ), + ( + "mvp.xtask.benchmark", + stamped( + json!({ + "type": "XtaskBenchmark", + "phase": "cargo_run_mvp_chat", + "status": "ready", + "run_id": 9, + "detail": {}, + }), + "xtask", + 1_025, + 25, + ), + ), + ( + "mvp.chat.runtime", + chat_span("prepare_runtime", "started", 1_030, 30), + ), + ( + "mvp.chat.runtime", + chat_span("ensure_orch_binary", "ready", 1_035, 35), + ), + ( + "mvp.chat.runtime", + chat_span("ensure_worker_binary", "ready", 1_036, 36), + ), + ( + "mvp.chat.runtime", + chat_span("prepare_runtime", "ready", 1_040, 40), + ), + ( + "mvp.orch.bootstrap", + stamped( + json!({ + "type": "OrchBootstrap", + "phase": "weights_loaded", + "status": "ready", + "run_id": 9, + "node_id": 1, + "detail": {}, + }), + "mvp-orchestrator", + 1_100, + 100, + ), + ), + ( + "mvp.chat.runtime", + chat_span("prompt_rpc", "ready", 1_120, 120), + ), + ] + } + + #[test] + fn benchmark_observability_dump_log_path_parser_finds_equals_and_separate_forms() { + assert_eq!( + explicit_dump_log_path_from_mvp_chat_args(&strings(&["--dump-logs=/tmp/a.ndjson"])), + Some(PathBuf::from("/tmp/a.ndjson")) + ); + assert_eq!( + explicit_dump_log_path_from_mvp_chat_args(&strings(&[ + "--", + "--run-id", + "7", + "--dump-logs", + "-logs.ndjson", + ])), + Some(PathBuf::from("-logs.ndjson")) + ); + assert_eq!( + explicit_dump_log_path_from_mvp_chat_args(&strings(&["--dump-logs"])), + None + ); + assert_eq!( + explicit_dump_log_path_from_mvp_chat_args(&strings(&["--dump-logs", "--run-id"])), + None + ); + assert_eq!( + run_id_from_mvp_chat_args(&strings(&["--", "--run-id", "42"])), + Some(42) + ); + assert_eq!( + run_id_from_mvp_chat_args(&strings(&["--run-id", "0"])), + None + ); + } + + #[test] + fn benchmark_observability_parse_dump_log_events_accepts_utf8_wrapper_and_plain_string_payload() + { + let wrapped = json!({"type": "Wrapped", "run_id": 9}); + let plain = json!({"type": "Plain", "run_id": 9}); + let path = write_dump_log( + "payload-shapes", + vec![ + benchmark_observability_archive_record( + "wrapped", + json!({"encoding": "utf8", "value": wrapped.to_string()}), + 111, + ), + benchmark_observability_archive_record("plain", json!(plain.to_string()), 112), + benchmark_observability_archive_record( + "bytes", + json!({"encoding": "bytes", "value": [0, 1]}), + 113, + ), + ], + ); + + let events = parse_dump_log_events(&path).expect("parse dump log events"); + let _ = fs::remove_file(path); + + assert_eq!(events.len(), 2); + assert_eq!(events[0].channel, "wrapped"); + assert_eq!(events[0].arrival_unix_ms, Some(111)); + assert_eq!( + events[0].event.get("type").and_then(Value::as_str), + Some("Wrapped") + ); + assert_eq!(events[1].channel, "plain"); + assert_eq!( + events[1].event.get("type").and_then(Value::as_str), + Some("Plain") + ); + } + + #[test] + fn benchmark_observability_report_requires_granular_decode_events() { + let events = parse_synthetic_events("missing-first-token", benchmark_report_events(false)); + + let error = + build_benchmark_report(&events, 80, 9).expect_err("missing first token should fail"); + + assert!( + error.starts_with("mvp-chat-check: missing benchmark event "), + "{error}" + ); + assert!( + error.contains("FirstTokenReady") && error.contains("request_id=1"), + "{error}" + ); + } + + #[test] + fn benchmark_observability_report_uses_pipeline_prompt_events() { + let events = + parse_synthetic_events("pipeline-report", pipeline_benchmark_report_events(true)); + + let report = build_benchmark_report(&events, 80, 9).expect("pipeline report builds"); + + assert!( + report + .lines + .iter() + .any(|line| line == "mvp-chat-check: benchmark: run_id=9") + ); + assert!(report.lines.iter().any(|line| { + line.starts_with("mvp-chat-check: benchmark prompt 1 ") + && line.contains("first_token_ms=8") + && line.contains("decode_ms=8") + })); + assert!(report.lines.iter().any(|line| { + line.starts_with("mvp-chat-check: benchmark prompt 2 ") + && line.contains("first_token_ms=8") + && line.contains("decode_ms=8") + })); + } +} + fn main() -> ExitCode { let mut args = std::env::args().skip(1); match args.next().as_deref() {