Skip to main content

rbt/measure/
mod.rs

1//! Measure packs — thesis proof harness (honest numbers only).
2//!
3//! Scenarios produce a machine-readable [`MeasureReport`] (wall time, rows, optional RSS).
4//! Public “beats Spark” claims require checked-in packs + reports; this module is the pack runner.
5//!
6//! ## P5c scenarios
7//!
8//! * [`SCENARIO_STREAM_VS_COLLECT`] — same DAG under stream vs collect; compare wall + RSS
9//! * [`SCENARIO_WHALE_SYNTHETIC`] — synthetic multi-file bronze (row count via env) + stream materialize
10//! * [`SCENARIO_COMPLEX_BRONZE`] — multi-artifact outer-join example with run scope
11
12use 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
24/// Built-in scenario names.
25pub 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
32/// Default synthetic row count for whale scenario (override with `RBT_MEASURE_ROWS`).
33pub const DEFAULT_WHALE_ROWS: usize = 100_000;
34/// Default number of bronze part files for whale scenario (`RBT_MEASURE_PARTS`).
35pub 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    /// Present for stream-vs-collect comparisons (P5c).
51    #[serde(default, skip_serializing_if = "Option::is_none")]
52    pub mode_compare: Option<ModeCompare>,
53    /// Synthetic generator settings when applicable.
54    #[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/// Side-by-side stream vs collect timings / RSS (Linux VmRSS when available).
61#[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
71/// Linux-only VmRSS from /proc/self/status (no new deps).
72pub 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
98/// Run a named measure scenario against a project.
99pub 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
238/// Run the project DAG twice: stream (default) then collect; report both.
239async 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    // Fresh engine per mode so DF caches do not muddy RSS.
254    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
336/// Synthetic multi-part bronze + single staging→silver model at configurable scale.
337async 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    // Stream pass (product default)
349    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    // Collect pass (memory pressure comparison)
365    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    // Single model keeps total_rows == synthetic row count (clean measure identity).
481    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
503/// Multi-artifact complex bronze example (when project_dir points at it).
504async fn measure_complex_bronze(project_dir: &Path, output_dir: &Path) -> Result<MeasureReport> {
505    // Prefer explicit path; if caller pointed at workspace root, use example.
506    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            // Try CWD-relative from crate workspace
514            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    // Prefer lake/lz/LATEST_RUN.json from fetch_bronze.py; fallback vars for empty CI.
531    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    // Research mini-lake: expect multi-model star with works + dims + fact.
567    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
595/// Write report JSON next to project or to given path.
596pub 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
606/// Default report path under project.
607pub 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
614/// List built-in scenarios for CLI help.
615pub 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}