1use anyhow::{bail, Context, Result};
7use serde::{Deserialize, Serialize};
8use std::path::{Path, PathBuf};
9use std::time::Instant;
10
11use crate::core::project::RbtProjectConfig;
12use crate::core::select::SelectMode;
13use crate::engine::TransformationEngine;
14
15pub const SCENARIO_SMOKE_PIPELINE: &str = "smoke_pipeline";
17pub const SCENARIO_VALIDATE_DX: &str = "validate_dx";
18pub const SCENARIO_INCREMENTAL_APPEND: &str = "incremental_append";
19
20#[derive(Debug, Clone, Serialize, Deserialize)]
21pub struct MeasureReport {
22 pub scenario: String,
23 pub project: String,
24 pub package_version: String,
25 pub wall_ms: u128,
26 pub models_executed: usize,
27 pub total_rows: usize,
28 pub bronze_sources: usize,
29 pub peak_rss_kb: Option<u64>,
30 pub notes: Vec<String>,
31 pub ok: bool,
32 pub error: Option<String>,
33}
34
35pub fn read_peak_rss_kb() -> Option<u64> {
37 #[cfg(target_os = "linux")]
38 {
39 let s = std::fs::read_to_string("/proc/self/status").ok()?;
40 for line in s.lines() {
41 if let Some(rest) = line.strip_prefix("VmRSS:") {
42 let kb: u64 = rest.split_whitespace().next()?.parse().ok()?;
43 return Some(kb);
44 }
45 }
46 None
47 }
48 #[cfg(not(target_os = "linux"))]
49 {
50 None
51 }
52}
53
54pub async fn run_measure_scenario(
56 scenario: &str,
57 project_dir: &Path,
58 output_dir: &Path,
59) -> Result<MeasureReport> {
60 match scenario {
61 SCENARIO_SMOKE_PIPELINE | "pipeline" | "smoke" => {
62 measure_pipeline(project_dir, output_dir).await
63 }
64 SCENARIO_VALIDATE_DX | "validate" | "dx" => measure_validate_dx(project_dir).await,
65 SCENARIO_INCREMENTAL_APPEND | "incremental" => {
66 measure_incremental_stub(project_dir, output_dir).await
67 }
68 other => bail!(
69 "E_RBT_MEASURE: unknown scenario '{other}'. \
70 Built-ins: {SCENARIO_SMOKE_PIPELINE}, {SCENARIO_VALIDATE_DX}, {SCENARIO_INCREMENTAL_APPEND}"
71 ),
72 }
73}
74
75async fn measure_pipeline(project_dir: &Path, output_dir: &Path) -> Result<MeasureReport> {
76 let config = RbtProjectConfig::load(project_dir)?;
77 let dag = config.build_dag(project_dir, None)?;
78 let engine = TransformationEngine::new();
79 let rss0 = read_peak_rss_kb();
80 let start = Instant::now();
81 let summary = engine
82 .execute_dag(&dag, project_dir, output_dir)
83 .await
84 .context("E_RBT_MEASURE: pipeline execute failed")?;
85 let wall_ms = start.elapsed().as_millis();
86 let rss1 = read_peak_rss_kb();
87 Ok(MeasureReport {
88 scenario: SCENARIO_SMOKE_PIPELINE.into(),
89 project: config.name,
90 package_version: crate::VERSION.into(),
91 wall_ms,
92 models_executed: summary.models_executed,
93 total_rows: summary.total_rows_produced,
94 bronze_sources: summary.bronze_sources_registered,
95 peak_rss_kb: rss1.or(rss0),
96 notes: vec![
97 "Full DAG materialize (stream default)".into(),
98 format!("select=all models={}", summary.models_executed),
99 ],
100 ok: true,
101 error: None,
102 })
103}
104
105async fn measure_validate_dx(project_dir: &Path) -> Result<MeasureReport> {
106 let config = RbtProjectConfig::load(project_dir)?;
107 let start = Instant::now();
108 let dag = config.build_dag(project_dir, None)?;
109 let _tiers = dag.execution_tiers()?;
110 let report = dag.validate_bronze_sources_with_roots(
111 project_dir,
112 crate::core::frontmatter::BronzeCheckMode::Fail,
113 &config.roots,
114 )?;
115 let wall_ms = start.elapsed().as_millis();
116 let ok = !report.has_errors();
117 Ok(MeasureReport {
118 scenario: SCENARIO_VALIDATE_DX.into(),
119 project: config.name,
120 package_version: crate::VERSION.into(),
121 wall_ms,
122 models_executed: 0,
123 total_rows: 0,
124 bronze_sources: 0,
125 peak_rss_kb: read_peak_rss_kb(),
126 notes: vec![
127 format!("models={}", dag.node_map.len()),
128 format!("bronze_errors={}", report.error_count()),
129 "DX metric: load DAG + bronze check latency".into(),
130 ],
131 ok,
132 error: if ok {
133 None
134 } else {
135 Some(format!("{} bronze errors", report.error_count()))
136 },
137 })
138}
139
140async fn measure_incremental_stub(project_dir: &Path, output_dir: &Path) -> Result<MeasureReport> {
141 let config = RbtProjectConfig::load(project_dir)?;
144 let dag = config.build_dag(project_dir, None)?;
145 let engine = TransformationEngine::new();
146 let start = Instant::now();
147 let s1 = engine
148 .execute_dag(&dag, project_dir, output_dir)
149 .await
150 .context("E_RBT_MEASURE: incremental first pass")?;
151 let s2 = engine
152 .execute_dag(&dag, project_dir, output_dir)
153 .await
154 .context("E_RBT_MEASURE: incremental second pass")?;
155 let wall_ms = start.elapsed().as_millis();
156 Ok(MeasureReport {
157 scenario: SCENARIO_INCREMENTAL_APPEND.into(),
158 project: config.name,
159 package_version: crate::VERSION.into(),
160 wall_ms,
161 models_executed: s1.models_executed + s2.models_executed,
162 total_rows: s1.total_rows_produced + s2.total_rows_produced,
163 bronze_sources: s1.bronze_sources_registered,
164 peak_rss_kb: read_peak_rss_kb(),
165 notes: vec![
166 "Two full DAG runs (baseline for overwrite cost)".into(),
167 "Models with materialization: incremental_append write part files".into(),
168 format!("pass1_rows={} pass2_rows={}", s1.total_rows_produced, s2.total_rows_produced),
169 ],
170 ok: true,
171 error: None,
172 })
173}
174
175pub fn write_measure_report(report: &MeasureReport, path: &Path) -> Result<()> {
177 if let Some(parent) = path.parent() {
178 std::fs::create_dir_all(parent)?;
179 }
180 let body = serde_json::to_vec_pretty(report)?;
181 std::fs::write(path, body)
182 .with_context(|| format!("E_RBT_MEASURE: write report {}", path.display()))?;
183 Ok(())
184}
185
186pub fn default_report_path(project_dir: &Path, scenario: &str) -> PathBuf {
188 project_dir
189 .join(".rbt")
190 .join("measure")
191 .join(format!("{scenario}.json"))
192}
193
194pub fn list_scenarios() -> &'static [&'static str] {
196 &[
197 SCENARIO_SMOKE_PIPELINE,
198 SCENARIO_VALIDATE_DX,
199 SCENARIO_INCREMENTAL_APPEND,
200 ]
201}
202
203#[allow(dead_code)]
205fn _select_mode() -> SelectMode {
206 SelectMode::Execute
207}