1use anyhow::{bail, Context, Result};
13use serde::{Deserialize, Serialize};
14use std::fs;
15use std::io::Write;
16use std::path::{Path, PathBuf};
17use std::time::Instant;
18
19use crate::core::frontmatter::BronzeCheckMode;
20use crate::core::project::{MaterializeConfig, MaterializeMode, RbtProjectConfig};
21use crate::core::run_scope::RunScope;
22use crate::engine::TransformationEngine;
23
24pub const SCENARIO_SMOKE_PIPELINE: &str = "smoke_pipeline";
26pub const SCENARIO_VALIDATE_DX: &str = "validate_dx";
27pub const SCENARIO_INCREMENTAL_APPEND: &str = "incremental_append";
28pub const SCENARIO_STREAM_VS_COLLECT: &str = "stream_vs_collect";
29pub const SCENARIO_WHALE_SYNTHETIC: &str = "whale_synthetic";
30pub const SCENARIO_COMPLEX_BRONZE: &str = "complex_bronze";
31
32pub const DEFAULT_WHALE_ROWS: usize = 100_000;
34pub const DEFAULT_WHALE_PARTS: usize = 20;
36
37#[derive(Debug, Clone, Serialize, Deserialize)]
38pub struct MeasureReport {
39 pub scenario: String,
40 pub project: String,
41 pub package_version: String,
42 pub wall_ms: u128,
43 pub models_executed: usize,
44 pub total_rows: usize,
45 pub bronze_sources: usize,
46 pub peak_rss_kb: Option<u64>,
47 pub notes: Vec<String>,
48 pub ok: bool,
49 pub error: Option<String>,
50 #[serde(default, skip_serializing_if = "Option::is_none")]
52 pub mode_compare: Option<ModeCompare>,
53 #[serde(default, skip_serializing_if = "Option::is_none")]
55 pub synthetic_rows: Option<usize>,
56 #[serde(default, skip_serializing_if = "Option::is_none")]
57 pub synthetic_parts: Option<usize>,
58}
59
60#[derive(Debug, Clone, Serialize, Deserialize)]
62pub struct ModeCompare {
63 pub stream_wall_ms: u128,
64 pub collect_wall_ms: u128,
65 pub stream_rss_kb: Option<u64>,
66 pub collect_rss_kb: Option<u64>,
67 pub rows: usize,
68 pub models: usize,
69}
70
71pub fn read_peak_rss_kb() -> Option<u64> {
73 #[cfg(target_os = "linux")]
74 {
75 let s = std::fs::read_to_string("/proc/self/status").ok()?;
76 for line in s.lines() {
77 if let Some(rest) = line.strip_prefix("VmRSS:") {
78 let kb: u64 = rest.split_whitespace().next()?.parse().ok()?;
79 return Some(kb);
80 }
81 }
82 None
83 }
84 #[cfg(not(target_os = "linux"))]
85 {
86 None
87 }
88}
89
90fn env_usize(name: &str, default: usize) -> usize {
91 std::env::var(name)
92 .ok()
93 .and_then(|s| s.parse().ok())
94 .filter(|&n| n > 0)
95 .unwrap_or(default)
96}
97
98pub async fn run_measure_scenario(
100 scenario: &str,
101 project_dir: &Path,
102 output_dir: &Path,
103) -> Result<MeasureReport> {
104 match scenario {
105 SCENARIO_SMOKE_PIPELINE | "pipeline" | "smoke" => {
106 measure_pipeline(project_dir, output_dir).await
107 }
108 SCENARIO_VALIDATE_DX | "validate" | "dx" => measure_validate_dx(project_dir).await,
109 SCENARIO_INCREMENTAL_APPEND | "incremental" => {
110 measure_incremental_stub(project_dir, output_dir).await
111 }
112 SCENARIO_STREAM_VS_COLLECT | "stream_collect" | "stream-vs-collect" => {
113 measure_stream_vs_collect(project_dir, output_dir).await
114 }
115 SCENARIO_WHALE_SYNTHETIC | "whale" | "synthetic" => {
116 measure_whale_synthetic(output_dir).await
117 }
118 SCENARIO_COMPLEX_BRONZE | "complex" | "multi_artifact" => {
119 measure_complex_bronze(project_dir, output_dir).await
120 }
121 other => bail!(
122 "E_RBT_MEASURE: unknown scenario '{other}'. Built-ins: {}",
123 list_scenarios().join(", ")
124 ),
125 }
126}
127
128async fn measure_pipeline(project_dir: &Path, output_dir: &Path) -> Result<MeasureReport> {
129 let config = RbtProjectConfig::load(project_dir)?;
130 let dag = config.build_dag(project_dir, None)?;
131 let engine = TransformationEngine::new();
132 let rss0 = read_peak_rss_kb();
133 let start = Instant::now();
134 let summary = engine
135 .execute_dag(&dag, project_dir, output_dir)
136 .await
137 .context("E_RBT_MEASURE: pipeline execute failed")?;
138 let wall_ms = start.elapsed().as_millis();
139 let rss1 = read_peak_rss_kb();
140 Ok(MeasureReport {
141 scenario: SCENARIO_SMOKE_PIPELINE.into(),
142 project: config.name,
143 package_version: crate::VERSION.into(),
144 wall_ms,
145 models_executed: summary.models_executed,
146 total_rows: summary.total_rows_produced,
147 bronze_sources: summary.bronze_sources_registered,
148 peak_rss_kb: rss1.or(rss0),
149 notes: vec![
150 "Full DAG materialize (stream default)".into(),
151 format!("select=all models={}", summary.models_executed),
152 ],
153 ok: true,
154 error: None,
155 mode_compare: None,
156 synthetic_rows: None,
157 synthetic_parts: None,
158 })
159}
160
161async fn measure_validate_dx(project_dir: &Path) -> Result<MeasureReport> {
162 let config = RbtProjectConfig::load(project_dir)?;
163 let start = Instant::now();
164 let dag = config.build_dag(project_dir, None)?;
165 let _tiers = dag.execution_tiers()?;
166 let report = dag.validate_bronze_sources_with_roots(
167 project_dir,
168 BronzeCheckMode::Fail,
169 &config.roots,
170 )?;
171 let wall_ms = start.elapsed().as_millis();
172 let ok = !report.has_errors();
173 Ok(MeasureReport {
174 scenario: SCENARIO_VALIDATE_DX.into(),
175 project: config.name,
176 package_version: crate::VERSION.into(),
177 wall_ms,
178 models_executed: 0,
179 total_rows: 0,
180 bronze_sources: 0,
181 peak_rss_kb: read_peak_rss_kb(),
182 notes: vec![
183 format!("models={}", dag.node_map.len()),
184 format!("bronze_errors={}", report.error_count()),
185 "DX metric: load DAG + bronze check latency".into(),
186 ],
187 ok,
188 error: if ok {
189 None
190 } else {
191 Some(format!("{} bronze errors", report.error_count()))
192 },
193 mode_compare: None,
194 synthetic_rows: None,
195 synthetic_parts: None,
196 })
197}
198
199async fn measure_incremental_stub(project_dir: &Path, output_dir: &Path) -> Result<MeasureReport> {
200 let config = RbtProjectConfig::load(project_dir)?;
201 let dag = config.build_dag(project_dir, None)?;
202 let engine = TransformationEngine::new();
203 let start = Instant::now();
204 let s1 = engine
205 .execute_dag(&dag, project_dir, output_dir)
206 .await
207 .context("E_RBT_MEASURE: incremental first pass")?;
208 let s2 = engine
209 .execute_dag(&dag, project_dir, output_dir)
210 .await
211 .context("E_RBT_MEASURE: incremental second pass")?;
212 let wall_ms = start.elapsed().as_millis();
213 Ok(MeasureReport {
214 scenario: SCENARIO_INCREMENTAL_APPEND.into(),
215 project: config.name,
216 package_version: crate::VERSION.into(),
217 wall_ms,
218 models_executed: s1.models_executed + s2.models_executed,
219 total_rows: s1.total_rows_produced + s2.total_rows_produced,
220 bronze_sources: s1.bronze_sources_registered,
221 peak_rss_kb: read_peak_rss_kb(),
222 notes: vec![
223 "Two full DAG runs (baseline for overwrite cost)".into(),
224 "Models with materialization: incremental_append write part files".into(),
225 format!(
226 "pass1_rows={} pass2_rows={}",
227 s1.total_rows_produced, s2.total_rows_produced
228 ),
229 ],
230 ok: true,
231 error: None,
232 mode_compare: None,
233 synthetic_rows: None,
234 synthetic_parts: None,
235 })
236}
237
238async fn measure_stream_vs_collect(
240 project_dir: &Path,
241 output_dir: &Path,
242) -> Result<MeasureReport> {
243 let config = RbtProjectConfig::load(project_dir)?;
244 let dag = config.build_dag(project_dir, None)?;
245
246 let stream_out = output_dir.join("measure_stream");
247 let collect_out = output_dir.join("measure_collect");
248 let _ = fs::remove_dir_all(&stream_out);
249 let _ = fs::remove_dir_all(&collect_out);
250
251 let mut stream_cfg = config.clone();
252 stream_cfg.materialize.mode = MaterializeMode::Stream;
253 let engine_s = TransformationEngine::new();
255 let rss_before_s = read_peak_rss_kb();
256 let t0 = Instant::now();
257 let sum_s = engine_s
258 .execute_dag_with_config(&dag, project_dir, &stream_out, &stream_cfg)
259 .await
260 .context("E_RBT_MEASURE: stream pass failed")?;
261 let stream_wall_ms = t0.elapsed().as_millis();
262 let stream_rss = read_peak_rss_kb().or(rss_before_s);
263
264 let mut collect_cfg = config.clone();
265 collect_cfg.materialize.mode = MaterializeMode::Collect;
266 let engine_c = TransformationEngine::new();
267 let rss_before_c = read_peak_rss_kb();
268 let t1 = Instant::now();
269 let sum_c = engine_c
270 .execute_dag_with_config(&dag, project_dir, &collect_out, &collect_cfg)
271 .await
272 .context("E_RBT_MEASURE: collect pass failed")?;
273 let collect_wall_ms = t1.elapsed().as_millis();
274 let collect_rss = read_peak_rss_kb().or(rss_before_c);
275
276 let rows = sum_s.total_rows_produced.max(sum_c.total_rows_produced);
277 let compare = ModeCompare {
278 stream_wall_ms,
279 collect_wall_ms,
280 stream_rss_kb: stream_rss,
281 collect_rss_kb: collect_rss,
282 rows,
283 models: sum_s.models_executed,
284 };
285
286 let mut notes = vec![
287 "Same DAG: materialize.mode=stream then collect (fresh SessionContext each)".into(),
288 format!(
289 "stream_wall_ms={} collect_wall_ms={} ratio_collect_over_stream={:.3}",
290 stream_wall_ms,
291 collect_wall_ms,
292 if stream_wall_ms == 0 {
293 0.0
294 } else {
295 collect_wall_ms as f64 / stream_wall_ms as f64
296 }
297 ),
298 ];
299 if let (Some(sr), Some(cr)) = (stream_rss, collect_rss) {
300 notes.push(format!(
301 "stream_rss_kb={sr} collect_rss_kb={cr} delta_kb={}",
302 cr as i64 - sr as i64
303 ));
304 notes.push(
305 "RSS is process VmRSS after each pass (not allocator peak); treat as directional".into(),
306 );
307 } else {
308 notes.push("RSS unavailable on this platform (Linux VmRSS only)".into());
309 }
310
311 Ok(MeasureReport {
312 scenario: SCENARIO_STREAM_VS_COLLECT.into(),
313 project: config.name,
314 package_version: crate::VERSION.into(),
315 wall_ms: stream_wall_ms + collect_wall_ms,
316 models_executed: sum_s.models_executed + sum_c.models_executed,
317 total_rows: sum_s.total_rows_produced + sum_c.total_rows_produced,
318 bronze_sources: sum_s.bronze_sources_registered,
319 peak_rss_kb: collect_rss.or(stream_rss),
320 notes,
321 ok: sum_s.total_rows_produced == sum_c.total_rows_produced,
322 error: if sum_s.total_rows_produced != sum_c.total_rows_produced {
323 Some(format!(
324 "row mismatch stream={} collect={}",
325 sum_s.total_rows_produced, sum_c.total_rows_produced
326 ))
327 } else {
328 None
329 },
330 mode_compare: Some(compare),
331 synthetic_rows: None,
332 synthetic_parts: None,
333 })
334}
335
336async fn measure_whale_synthetic(output_dir: &Path) -> Result<MeasureReport> {
338 let rows = env_usize("RBT_MEASURE_ROWS", DEFAULT_WHALE_ROWS);
339 let parts = env_usize("RBT_MEASURE_PARTS", DEFAULT_WHALE_PARTS).max(1);
340 let work = output_dir.join("whale_synth_project");
341 let _ = fs::remove_dir_all(&work);
342 write_whale_project(&work, rows, parts)?;
343
344 let config = RbtProjectConfig::load(&work)?;
345 let dag = config.build_dag(&work, None)?;
346 let lake_out = work.join("lake").join("silver");
347
348 let mut stream_cfg = config.clone();
350 stream_cfg.materialize = MaterializeConfig {
351 mode: MaterializeMode::Stream,
352 ..config.materialize.clone()
353 };
354 let engine_s = TransformationEngine::new();
355 let rss0 = read_peak_rss_kb();
356 let t0 = Instant::now();
357 let sum_s = engine_s
358 .execute_dag_with_config(&dag, &work, &lake_out.join("stream_run"), &stream_cfg)
359 .await
360 .context("E_RBT_MEASURE: whale stream pass")?;
361 let stream_wall_ms = t0.elapsed().as_millis();
362 let stream_rss = read_peak_rss_kb().or(rss0);
363
364 let mut collect_cfg = config.clone();
366 collect_cfg.materialize = MaterializeConfig {
367 mode: MaterializeMode::Collect,
368 ..config.materialize.clone()
369 };
370 let engine_c = TransformationEngine::new();
371 let rss1 = read_peak_rss_kb();
372 let t1 = Instant::now();
373 let sum_c = engine_c
374 .execute_dag_with_config(&dag, &work, &lake_out.join("collect_run"), &collect_cfg)
375 .await
376 .context("E_RBT_MEASURE: whale collect pass")?;
377 let collect_wall_ms = t1.elapsed().as_millis();
378 let collect_rss = read_peak_rss_kb().or(rss1);
379
380 let compare = ModeCompare {
381 stream_wall_ms,
382 collect_wall_ms,
383 stream_rss_kb: stream_rss,
384 collect_rss_kb: collect_rss,
385 rows: sum_s.total_rows_produced,
386 models: sum_s.models_executed,
387 };
388
389 let ok = sum_s.total_rows_produced == rows && sum_c.total_rows_produced == rows;
390 let mut notes = vec![
391 format!("Synthetic JSONL bronze: {rows} rows across {parts} part files"),
392 "Env: RBT_MEASURE_ROWS, RBT_MEASURE_PARTS".into(),
393 format!(
394 "stream_wall_ms={stream_wall_ms} collect_wall_ms={collect_wall_ms} stream_rss_kb={stream_rss:?} collect_rss_kb={collect_rss:?}"
395 ),
396 "Whale-ish default is 100k rows / 20 parts — raise RBT_MEASURE_ROWS for larger packs".into(),
397 ];
398 if !ok {
399 notes.push(format!(
400 "expected_rows={rows} stream_rows={} collect_rows={}",
401 sum_s.total_rows_produced, sum_c.total_rows_produced
402 ));
403 }
404
405 Ok(MeasureReport {
406 scenario: SCENARIO_WHALE_SYNTHETIC.into(),
407 project: config.name,
408 package_version: crate::VERSION.into(),
409 wall_ms: stream_wall_ms + collect_wall_ms,
410 models_executed: sum_s.models_executed + sum_c.models_executed,
411 total_rows: sum_s.total_rows_produced + sum_c.total_rows_produced,
412 bronze_sources: sum_s.bronze_sources_registered,
413 peak_rss_kb: collect_rss.or(stream_rss),
414 notes,
415 ok,
416 error: if ok {
417 None
418 } else {
419 Some("row count did not match synthetic generator".into())
420 },
421 mode_compare: Some(compare),
422 synthetic_rows: Some(rows),
423 synthetic_parts: Some(parts),
424 })
425}
426
427fn write_whale_project(root: &Path, total_rows: usize, parts: usize) -> Result<()> {
428 let bronze = root.join("lake/bronze/events");
429 fs::create_dir_all(&bronze)?;
430 fs::create_dir_all(root.join("models/staging"))?;
431 fs::create_dir_all(root.join("models/transforms"))?;
432 fs::create_dir_all(root.join("models/marts"))?;
433
434 let per = total_rows.div_ceil(parts);
435 let mut written = 0usize;
436 for p in 0..parts {
437 if written >= total_rows {
438 break;
439 }
440 let n = (total_rows - written).min(per);
441 let path = bronze.join(format!("part-{p:04}.jsonl"));
442 let mut f = fs::File::create(&path)?;
443 for i in 0..n {
444 let id = written + i;
445 writeln!(
446 f,
447 "{{\"id\":{id},\"domain\":\"d{}\",\"amount\":{}}}",
448 id % 97,
449 (id % 1000) as f64 * 0.01
450 )?;
451 }
452 written += n;
453 }
454
455 fs::write(
456 root.join("rbt_project.yml"),
457 r#"name: whale_synthetic
458version: "0.1.0"
459contract_version: "measure-whale-v1"
460models_dir: models
461target_path: lake/gold
462roots:
463 lake: lake
464layers:
465 staging:
466 path: models/staging
467 target_path: lake/silver
468 default_format: parquet
469 transforms:
470 path: models/transforms
471 target_path: lake/silver
472 default_format: parquet
473 marts:
474 path: models/marts
475 target_path: lake/gold
476 default_format: parquet
477"#,
478 )?;
479
480 fs::write(
482 root.join("models/staging/stg_events.sql"),
483 r#"---
484description: Full-refresh silver mirror of multi-part bronze events (whale measure).
485source_format: jsonl
486scan_path: $lake/bronze/events
487path_glob: "*.jsonl"
488stage_mode: full_refresh
489columns:
490 id: { dtype: int64 }
491 domain: { dtype: utf8 }
492 amount: { dtype: float64 }
493tests:
494 not_null: [id]
495---
496SELECT id, domain, amount FROM {{ source('bronze', 'events') }}
497"#,
498 )?;
499
500 Ok(())
501}
502
503async fn measure_complex_bronze(project_dir: &Path, output_dir: &Path) -> Result<MeasureReport> {
505 let project = if project_dir.join("models/staging/stg_plan.sql").exists() {
507 project_dir.to_path_buf()
508 } else {
509 let candidate = project_dir.join("examples/complex_bronze_landing");
510 if candidate.join("rbt_project.yml").exists() {
511 candidate
512 } else {
513 let alt = PathBuf::from("examples/complex_bronze_landing");
515 if alt.join("rbt_project.yml").exists() {
516 alt
517 } else {
518 bail!(
519 "E_RBT_MEASURE: complex_bronze needs examples/complex_bronze_landing \
520 (got project_dir={})",
521 project_dir.display()
522 );
523 }
524 }
525 };
526
527 let config = RbtProjectConfig::load(&project)?;
528 let dag = config.build_dag(&project, None)?;
529
530 let mut domain = "semicon-ai-research".to_string();
532 let mut report_date = "2026-08-01".to_string();
533 let mut run_id = "run20260802T022258Z".to_string();
534 let pointer = project.join("lake/lz/LATEST_RUN.json");
535 if pointer.is_file() {
536 if let Ok(v) = serde_json::from_str::<serde_json::Value>(&std::fs::read_to_string(&pointer)?)
537 {
538 if let Some(s) = v.get("domain").and_then(|x| x.as_str()) {
539 domain = s.to_string();
540 }
541 if let Some(s) = v.get("report_date").and_then(|x| x.as_str()) {
542 report_date = s.to_string();
543 }
544 if let Some(s) = v.get("run_id").and_then(|x| x.as_str()) {
545 run_id = s.to_string();
546 }
547 }
548 }
549
550 let mut scope = RunScope::new()
551 .with_var("domain", domain.clone())
552 .with_var("report_date", report_date.clone())
553 .with_var("run_id", run_id.clone());
554 scope.write_receipt = true;
555 scope.skip_if_fingerprint_match = false;
556
557 let engine = TransformationEngine::new();
558 let rss0 = read_peak_rss_kb();
559 let start = Instant::now();
560 let summary = engine
561 .execute_dag_with_scope(&dag, &project, output_dir, &config, &scope)
562 .await
563 .context("E_RBT_MEASURE: complex_bronze execute failed")?;
564 let wall_ms = start.elapsed().as_millis();
565
566 let ok = summary.total_rows_produced >= 10 && summary.models_executed >= 5;
568 Ok(MeasureReport {
569 scenario: SCENARIO_COMPLEX_BRONZE.into(),
570 project: config.name,
571 package_version: crate::VERSION.into(),
572 wall_ms,
573 models_executed: summary.models_executed,
574 total_rows: summary.total_rows_produced,
575 bronze_sources: summary.bronze_sources_registered,
576 peak_rss_kb: read_peak_rss_kb().or(rss0),
577 notes: vec![
578 "Research papers mini-lake: PubMed/Crossref bronze → silver stg → gold tf/marts".into(),
579 format!("fingerprint={:?}", summary.bronze_fingerprint),
580 format!("receipt={:?}", summary.receipt_path),
581 format!("Run vars: domain={domain} report_date={report_date} run_id={run_id}"),
582 ],
583 ok,
584 error: if ok {
585 None
586 } else {
587 Some("expected multi-model research lake materialize (works+dims+fact)".into())
588 },
589 mode_compare: None,
590 synthetic_rows: None,
591 synthetic_parts: None,
592 })
593}
594
595pub fn write_measure_report(report: &MeasureReport, path: &Path) -> Result<()> {
597 if let Some(parent) = path.parent() {
598 std::fs::create_dir_all(parent)?;
599 }
600 let body = serde_json::to_vec_pretty(report)?;
601 std::fs::write(path, body)
602 .with_context(|| format!("E_RBT_MEASURE: write report {}", path.display()))?;
603 Ok(())
604}
605
606pub fn default_report_path(project_dir: &Path, scenario: &str) -> PathBuf {
608 project_dir
609 .join(".rbt")
610 .join("measure")
611 .join(format!("{scenario}.json"))
612}
613
614pub fn list_scenarios() -> Vec<&'static str> {
616 vec![
617 SCENARIO_SMOKE_PIPELINE,
618 SCENARIO_VALIDATE_DX,
619 SCENARIO_INCREMENTAL_APPEND,
620 SCENARIO_STREAM_VS_COLLECT,
621 SCENARIO_WHALE_SYNTHETIC,
622 SCENARIO_COMPLEX_BRONZE,
623 ]
624}
625
626#[cfg(test)]
627mod tests {
628 use super::*;
629
630 #[tokio::test]
631 async fn whale_synthetic_small_ok() {
632 std::env::set_var("RBT_MEASURE_ROWS", "500");
633 std::env::set_var("RBT_MEASURE_PARTS", "4");
634 let dir = tempfile::tempdir().unwrap();
635 let report = measure_whale_synthetic(dir.path()).await.unwrap();
636 assert!(report.ok, "{:?}", report.error);
637 assert_eq!(report.synthetic_rows, Some(500));
638 assert!(report.mode_compare.is_some());
639 assert_eq!(report.mode_compare.as_ref().unwrap().rows, 500);
640 std::env::remove_var("RBT_MEASURE_ROWS");
641 std::env::remove_var("RBT_MEASURE_PARTS");
642 }
643
644 #[tokio::test]
645 async fn stream_vs_collect_on_smoke() {
646 let root = PathBuf::from(env!("CARGO_MANIFEST_DIR")).join("../../examples/smoke_fixture");
647 if !root.exists() {
648 return;
649 }
650 let out = tempfile::tempdir().unwrap();
651 let report = measure_stream_vs_collect(&root, out.path()).await.unwrap();
652 assert!(report.ok, "{:?}", report.error);
653 assert!(report.mode_compare.is_some());
654 }
655}