Skip to main content

rbt/core/
project.rs

1use anyhow::{Context, Result};
2use serde::{Deserialize, Serialize};
3use std::collections::HashMap;
4use std::fs;
5use std::path::{Path, PathBuf};
6use walkdir::WalkDir;
7
8use super::dag::{Materialization, ModelDag, ModelLayer, OutputFormat};
9use super::paths::{resolve_configured_path, resolve_project_path};
10
11/// Default MemTable row cutoff when `ref_strategy: memtable` and max rows omitted.
12pub const DEFAULT_MEMTABLE_MAX_ROWS: usize = 50_000;
13
14/// How completed models are exposed to downstream `{{ ref() }}` in the same run.
15///
16/// Default is lake-as-truth **parquet / file re-read** (no long-lived MemTable).
17/// Opt into MemTable via `materialize.ref_strategy: memtable` in `rbt_project.yml`.
18#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
19#[serde(rename_all = "snake_case")]
20pub enum RefStrategy {
21    /// Always re-read the written lake file for `ref()` (default).
22    #[default]
23    #[serde(alias = "parquet_reread", alias = "lake", alias = "file")]
24    Parquet,
25    /// Keep an in-memory `MemTable` when `row_count < memtable_max_rows`; else re-read file.
26    #[serde(alias = "mem_table", alias = "memory", alias = "arc")]
27    Memtable,
28}
29
30/// Chosen backend after applying strategy + row cutoff.
31#[derive(Debug, Clone, Copy, PartialEq, Eq)]
32pub enum RefBackend {
33    /// DataFusion `MemTable` holding Arrow batches (Arc-shared).
34    MemTable,
35    /// Re-read model output from the lake path (`register_parquet` / json / csv).
36    LakeFile,
37}
38
39/// How model SQL results are written to the lake.
40///
41/// Default is **`stream`**: `execute_stream` → batch write → drop batch (no full
42/// `Vec<RecordBatch>` retention). Use **`collect`** only for debugging / emergency.
43#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
44#[serde(rename_all = "snake_case")]
45pub enum MaterializeMode {
46    /// Pull DataFusion stream batch-by-batch; atomic publish; bounded peak RAM.
47    #[default]
48    #[serde(alias = "streaming")]
49    Stream,
50    /// `DataFrame::collect` then write (legacy; holds full result in RAM).
51    #[serde(alias = "batch", alias = "legacy")]
52    Collect,
53}
54
55/// Default Parquet max row-group row count for streaming writers.
56pub const DEFAULT_MAX_ROW_GROUP_ROWS: usize = 1_000_000;
57/// Default Parquet in-progress size threshold before `flush()` (128 MiB).
58pub const DEFAULT_MAX_ROW_GROUP_BYTES: usize = 128 * 1024 * 1024;
59
60/// Optional materialization / `ref()` registration policy (`materialize:` in yml).
61///
62/// All fields are optional; omitting the whole block keeps lake-as-truth Parquet re-read
63/// and **stream** write mode.
64#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
65pub struct MaterializeConfig {
66    /// `stream` (default) | `collect`
67    #[serde(default)]
68    pub mode: MaterializeMode,
69    /// `parquet` (default) | `memtable`
70    #[serde(default)]
71    pub ref_strategy: RefStrategy,
72    /// Used only when `ref_strategy: memtable`. Defaults to [`DEFAULT_MEMTABLE_MAX_ROWS`].
73    #[serde(default = "default_memtable_max_rows")]
74    pub memtable_max_rows: usize,
75    /// Parquet `WriterProperties` max rows per row group (stream + collect writers).
76    #[serde(default = "default_max_row_group_rows")]
77    pub max_row_group_rows: usize,
78    /// Soft flush threshold for Parquet `in_progress_size` (bytes).
79    #[serde(default = "default_max_row_group_bytes")]
80    pub max_row_group_bytes: usize,
81}
82
83fn default_memtable_max_rows() -> usize {
84    DEFAULT_MEMTABLE_MAX_ROWS
85}
86
87fn default_max_row_group_rows() -> usize {
88    DEFAULT_MAX_ROW_GROUP_ROWS
89}
90
91fn default_max_row_group_bytes() -> usize {
92    DEFAULT_MAX_ROW_GROUP_BYTES
93}
94
95impl Default for MaterializeConfig {
96    fn default() -> Self {
97        Self {
98            mode: MaterializeMode::Stream,
99            ref_strategy: RefStrategy::Parquet,
100            memtable_max_rows: DEFAULT_MEMTABLE_MAX_ROWS,
101            max_row_group_rows: DEFAULT_MAX_ROW_GROUP_ROWS,
102            max_row_group_bytes: DEFAULT_MAX_ROW_GROUP_BYTES,
103        }
104    }
105}
106
107/// Default max size for a single opaque protobuf bronze file (1 GiB).
108pub const DEFAULT_PROTOBUF_MAX_PAYLOAD_BYTES: u64 = 1024 * 1024 * 1024;
109
110/// Optional scan / bronze ingest limits (`scan:` in yml).
111#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
112pub struct ScanConfig {
113    /// Max bytes for one `source_format: protobuf` file. Default: 1 GiB.
114    ///
115    /// Override to raise/lower the safety cap for opaque `payload` columns.
116    #[serde(default = "default_protobuf_max_payload_bytes")]
117    pub protobuf_max_payload_bytes: u64,
118    /// Spill Arrow IPC bronze (hive / multi-file) to a single Parquet cache, then
119    /// register via DataFusion listing — avoids holding every IPC partition in a MemTable.
120    ///
121    /// Default **true**. Set `false` only for tiny trees or debugging MemTable path.
122    #[serde(default = "default_spill_arrow_ipc")]
123    pub spill_arrow_ipc: bool,
124    /// Directory (project-relative or absolute / `$root`) for bronze spill files.
125    /// Default: `.rbt/bronze_spill`.
126    #[serde(default = "default_spill_dir")]
127    pub spill_dir: String,
128}
129
130fn default_protobuf_max_payload_bytes() -> u64 {
131    DEFAULT_PROTOBUF_MAX_PAYLOAD_BYTES
132}
133
134fn default_spill_arrow_ipc() -> bool {
135    true
136}
137
138fn default_spill_dir() -> String {
139    ".rbt/bronze_spill".into()
140}
141
142impl Default for ScanConfig {
143    fn default() -> Self {
144        Self {
145            protobuf_max_payload_bytes: DEFAULT_PROTOBUF_MAX_PAYLOAD_BYTES,
146            spill_arrow_ipc: true,
147            spill_dir: default_spill_dir(),
148        }
149    }
150}
151
152impl MaterializeConfig {
153    /// Resolve write mode: env `RBT_MATERIALIZE_MODE=stream|collect` overrides yml.
154    pub fn effective_mode(&self) -> MaterializeMode {
155        match std::env::var("RBT_MATERIALIZE_MODE")
156            .or_else(|_| std::env::var("RBT_STREAM_MATERIALIZE"))
157            .ok()
158            .as_deref()
159        {
160            // RBT_STREAM_MATERIALIZE=1 / true → stream; 0 / false → collect
161            Some("1") | Some("true") | Some("TRUE") | Some("yes") | Some("on") => {
162                MaterializeMode::Stream
163            }
164            Some("0") | Some("false") | Some("FALSE") | Some("no") | Some("off") => {
165                MaterializeMode::Collect
166            }
167            Some(s) if s.eq_ignore_ascii_case("stream") || s.eq_ignore_ascii_case("streaming") => {
168                MaterializeMode::Stream
169            }
170            Some(s) if s.eq_ignore_ascii_case("collect") || s.eq_ignore_ascii_case("batch") => {
171                MaterializeMode::Collect
172            }
173            _ => self.mode,
174        }
175    }
176
177    /// Decide MemTable vs lake file for a model that produced `row_count` rows.
178    pub fn choose_ref_backend(&self, row_count: usize) -> RefBackend {
179        match self.ref_strategy {
180            RefStrategy::Parquet => RefBackend::LakeFile,
181            RefStrategy::Memtable if row_count < self.memtable_max_rows => RefBackend::MemTable,
182            RefStrategy::Memtable => RefBackend::LakeFile,
183        }
184    }
185}
186
187/// Layer-specific target storage & path configuration.
188#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
189pub struct LayerConfig {
190    pub path: PathBuf,
191    pub target_path: PathBuf,
192    pub default_format: Option<String>,
193}
194
195/// Project-wide `rbt_project.yml` configuration schema.
196#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
197pub struct RbtProjectConfig {
198    pub name: String,
199    pub version: String,
200    pub models_dir: PathBuf,
201    pub target_path: PathBuf,
202    #[serde(default)]
203    pub layers: HashMap<String, LayerConfig>,
204    /// Optional; defaults to lake-as-truth Parquet re-read for `ref()`.
205    #[serde(default)]
206    pub materialize: MaterializeConfig,
207    /// Optional bronze scan limits (e.g. protobuf payload cap).
208    #[serde(default)]
209    pub scan: ScanConfig,
210    /// Named absolute (or relative) roots for multi-root lakes.
211    ///
212    /// Referenced in paths as `$name` or `${name}` (e.g. `$nonprod_lake/lz/runs`).
213    #[serde(default)]
214    pub roots: HashMap<String, String>,
215}
216
217impl Default for RbtProjectConfig {
218    fn default() -> Self {
219        let mut layers = HashMap::new();
220        layers.insert(
221            "staging".to_string(),
222            LayerConfig {
223                path: PathBuf::from("models/staging"),
224                target_path: PathBuf::from("lake/silver"),
225                default_format: Some("parquet".to_string()),
226            },
227        );
228        layers.insert(
229            "transforms".to_string(),
230            LayerConfig {
231                path: PathBuf::from("models/transforms"),
232                target_path: PathBuf::from("lake/gold"),
233                default_format: Some("parquet".to_string()),
234            },
235        );
236        layers.insert(
237            "marts".to_string(),
238            LayerConfig {
239                path: PathBuf::from("models/marts"),
240                target_path: PathBuf::from("lake/gold"),
241                default_format: Some("parquet_and_iceberg".to_string()),
242            },
243        );
244
245        Self {
246            name: "rbt_project".to_string(),
247            version: "1.0.0".to_string(),
248            models_dir: PathBuf::from("models"),
249            target_path: PathBuf::from("lake/gold"),
250            layers,
251            materialize: MaterializeConfig::default(),
252            scan: ScanConfig::default(),
253            roots: HashMap::new(),
254        }
255    }
256}
257
258impl RbtProjectConfig {
259    /// Loads `rbt_project.yml` from project directory or returns default configuration.
260    pub fn load(project_dir: &Path) -> Result<Self> {
261        let project_file = project_dir.join("rbt_project.yml");
262        if project_file.exists() {
263            let content = fs::read_to_string(&project_file).with_context(|| {
264                format!(
265                    "E_RBT_PROJECT_LOAD: cannot read project file {}",
266                    project_file.display()
267                )
268            })?;
269            let mut config: RbtProjectConfig = serde_yaml::from_str(&content).with_context(|| {
270                format!(
271                    "E_RBT_PROJECT_LOAD: failed to parse {}. \
272                     Check required keys (name, version, models_dir, target_path) and \
273                     optional materialize:/scan:/roots:/layers blocks.",
274                    project_file.display()
275                )
276            })?;
277
278            let defaults = Self::default();
279            for (key, val) in defaults.layers {
280                config.layers.entry(key).or_insert(val);
281            }
282            Ok(config)
283        } else {
284            Ok(Self::default())
285        }
286    }
287
288    /// Resolve a configured path (absolute, relative, or `$root/...`) against the project.
289    pub fn resolve_path(&self, project_dir: &Path, configured: &str) -> Result<PathBuf> {
290        resolve_project_path(project_dir, configured, &self.roots)
291    }
292
293    /// Layer output directory (file parent for flat parquet, or table root parent).
294    pub fn resolve_layer_target_dir(
295        &self,
296        project_dir: &Path,
297        layer: ModelLayer,
298    ) -> Result<PathBuf> {
299        let layer_key = match layer {
300            ModelLayer::Staging => "staging",
301            ModelLayer::Transform => "transforms",
302            ModelLayer::Mart => "marts",
303        };
304        if let Some(layer_cfg) = self.layers.get(layer_key) {
305            resolve_configured_path(project_dir, &layer_cfg.target_path, &self.roots)
306        } else {
307            resolve_configured_path(project_dir, &self.target_path, &self.roots)
308        }
309    }
310
311    /// Resolves destination output file path for a model based on its layer configuration.
312    ///
313    /// Supports absolute `target_path` and `$root` templates — never nests an absolute
314    /// lake path under `project_dir`.
315    pub fn resolve_model_target_path(
316        &self,
317        project_dir: &Path,
318        model_name: &str,
319        layer: ModelLayer,
320        ext: &str,
321    ) -> Result<PathBuf> {
322        let dir = self
323            .resolve_layer_target_dir(project_dir, layer)
324            .with_context(|| {
325                format!(
326                    "E_RBT_MODEL_TARGET: cannot resolve output directory for model '{model_name}' \
327                 (layer={layer:?}). Check `layers.*.target_path`, top-level `target_path`, and \
328                 `roots:` in rbt_project.yml."
329                )
330            })?;
331        Ok(dir.join(format!("{model_name}.{ext}")))
332    }
333
334    /// Directory target for Iceberg-style table roots (no file extension).
335    pub fn resolve_model_target_dir(
336        &self,
337        project_dir: &Path,
338        model_name: &str,
339        layer: ModelLayer,
340    ) -> Result<PathBuf> {
341        let dir = self
342            .resolve_layer_target_dir(project_dir, layer)
343            .with_context(|| {
344                format!(
345                    "E_RBT_MODEL_TARGET: cannot resolve table directory for model '{model_name}' \
346                 (layer={layer:?}). Check layer target_path and roots:."
347                )
348            })?;
349        Ok(dir.join(model_name))
350    }
351
352    /// Discovers all `.sql` models under `models/` directory, resolves layer target paths, and constructs `ModelDag`.
353    pub fn build_dag(
354        &self,
355        project_dir: &Path,
356        cli_format_override: Option<OutputFormat>,
357    ) -> Result<ModelDag> {
358        let models_dir = project_dir.join(&self.models_dir);
359        let mut dag = ModelDag::new();
360
361        if !models_dir.exists() {
362            let default_fmt = cli_format_override.unwrap_or(OutputFormat::Parquet);
363            dag.add_model_with_format(
364                "stg_users",
365                "SELECT 1 AS id, 'Alice' AS name, 'admin' AS role",
366                Materialization::Table,
367                default_fmt,
368                None,
369                "",
370            )?;
371            dag.build_graph()?;
372            return Ok(dag);
373        }
374
375        let mut model_count = 0;
376        for entry in WalkDir::new(&models_dir).into_iter().filter_map(|e| e.ok()) {
377            let path = entry.path();
378            if path.is_file() && path.extension().is_some_and(|ext| ext == "sql") {
379                let stem = path.file_stem().and_then(|s| s.to_str()).with_context(|| {
380                    format!(
381                        "E_RBT_MODEL_NAME: invalid file stem for model path {}",
382                        path.display()
383                    )
384                })?;
385                let raw_sql = fs::read_to_string(path).with_context(|| {
386                    format!(
387                        "E_RBT_MODEL_IO: failed reading model SQL {}",
388                        path.display()
389                    )
390                })?;
391
392                let layer = ModelLayer::from_name(stem);
393                let format = cli_format_override.clone().unwrap_or_else(|| {
394                    let layer_key = match layer {
395                        ModelLayer::Staging => "staging",
396                        ModelLayer::Transform => "transforms",
397                        ModelLayer::Mart => "marts",
398                    };
399                    if let Some(l_cfg) = self.layers.get(layer_key) {
400                        match l_cfg.default_format.as_deref() {
401                            Some("parquet") => OutputFormat::Parquet,
402                            Some("jsonl") => OutputFormat::Jsonl,
403                            Some("csv") => OutputFormat::Csv,
404                            Some("iceberg") => OutputFormat::Iceberg,
405                            Some("parquet_and_iceberg") => OutputFormat::ParquetAndIceberg,
406                            _ => OutputFormat::Parquet,
407                        }
408                    } else {
409                        OutputFormat::Parquet
410                    }
411                });
412
413                let target_file_path = match format {
414                    OutputFormat::Iceberg => self
415                        .resolve_model_target_dir(project_dir, stem, layer)
416                        .with_context(|| {
417                            format!(
418                                "E_RBT_MODEL_TARGET: model '{stem}' (Iceberg) — \
419                                 failed resolving layer target. \
420                                 layer={layer:?}; project={}",
421                                project_dir.display()
422                            )
423                        })?,
424                    OutputFormat::Parquet
425                    | OutputFormat::ParquetAndIceberg
426                    | OutputFormat::ZeroCopyClone => self
427                        .resolve_model_target_path(project_dir, stem, layer, "parquet")
428                        .with_context(|| {
429                            format!(
430                                "E_RBT_MODEL_TARGET: model '{stem}' (parquet) — \
431                                 failed resolving layer target. \
432                                 layer={layer:?}; project={}",
433                                project_dir.display()
434                            )
435                        })?,
436                    OutputFormat::Jsonl => self
437                        .resolve_model_target_path(project_dir, stem, layer, "jsonl")
438                        .with_context(|| {
439                            format!(
440                                "E_RBT_MODEL_TARGET: model '{stem}' (jsonl) — \
441                                 failed resolving layer target"
442                            )
443                        })?,
444                    OutputFormat::Csv => self
445                        .resolve_model_target_path(project_dir, stem, layer, "csv")
446                        .with_context(|| {
447                            format!(
448                                "E_RBT_MODEL_TARGET: model '{stem}' (csv) — \
449                                 failed resolving layer target"
450                            )
451                        })?,
452                };
453
454                dag.add_model_with_format(
455                    stem,
456                    &raw_sql,
457                    Materialization::Table,
458                    format,
459                    Some(target_file_path.to_string_lossy().to_string()),
460                    "",
461                )?;
462                model_count += 1;
463            }
464        }
465
466        if model_count == 0 {
467            let default_fmt = cli_format_override.unwrap_or(OutputFormat::Parquet);
468            dag.add_model_with_format(
469                "stg_users",
470                "SELECT 1 AS id, 'Alice' AS name, 'admin' AS role",
471                Materialization::Table,
472                default_fmt,
473                None,
474                "",
475            )?;
476        }
477
478        dag.build_graph()?;
479        Ok(dag)
480    }
481}
482
483#[cfg(test)]
484mod tests {
485    use super::*;
486
487    #[test]
488    fn test_layer_target_path_resolution() -> Result<()> {
489        let config = RbtProjectConfig::default();
490        let project_dir = Path::new("/tmp/test_project");
491
492        let stg_path = config.resolve_model_target_path(
493            project_dir,
494            "stg_trades",
495            ModelLayer::Staging,
496            "parquet",
497        )?;
498        assert_eq!(stg_path, project_dir.join("lake/silver/stg_trades.parquet"));
499
500        let tf_path = config.resolve_model_target_path(
501            project_dir,
502            "tf_1m_bars",
503            ModelLayer::Transform,
504            "parquet",
505        )?;
506        assert_eq!(tf_path, project_dir.join("lake/gold/tf_1m_bars.parquet"));
507
508        let mart_path = config.resolve_model_target_path(
509            project_dir,
510            "fact_1d_bars",
511            ModelLayer::Mart,
512            "parquet",
513        )?;
514        assert_eq!(
515            mart_path,
516            project_dir.join("lake/gold/fact_1d_bars.parquet")
517        );
518
519        Ok(())
520    }
521
522    #[test]
523    fn materialize_defaults_to_parquet_reread() {
524        let cfg = MaterializeConfig::default();
525        assert_eq!(cfg.ref_strategy, RefStrategy::Parquet);
526        assert_eq!(cfg.mode, MaterializeMode::Stream);
527        assert_eq!(cfg.memtable_max_rows, DEFAULT_MEMTABLE_MAX_ROWS);
528        assert_eq!(cfg.max_row_group_rows, DEFAULT_MAX_ROW_GROUP_ROWS);
529        assert_eq!(cfg.choose_ref_backend(0), RefBackend::LakeFile);
530        assert_eq!(cfg.choose_ref_backend(1_000_000), RefBackend::LakeFile);
531    }
532
533    #[test]
534    fn materialize_mode_from_yaml() -> Result<()> {
535        let yml = r#"
536name: t
537version: "1"
538models_dir: models
539target_path: lake/gold
540materialize:
541  mode: collect
542  max_row_group_rows: 1000
543"#;
544        let cfg: RbtProjectConfig = serde_yaml::from_str(yml)?;
545        assert_eq!(cfg.materialize.mode, MaterializeMode::Collect);
546        assert_eq!(cfg.materialize.max_row_group_rows, 1000);
547        Ok(())
548    }
549
550    #[test]
551    fn materialize_memtable_respects_cutoff() {
552        let cfg = MaterializeConfig {
553            ref_strategy: RefStrategy::Memtable,
554            memtable_max_rows: 50_000,
555            ..Default::default()
556        };
557        assert_eq!(cfg.choose_ref_backend(49_999), RefBackend::MemTable);
558        assert_eq!(cfg.choose_ref_backend(50_000), RefBackend::LakeFile);
559        assert_eq!(cfg.choose_ref_backend(50_001), RefBackend::LakeFile);
560    }
561
562    #[test]
563    fn parse_materialize_block_from_yaml() -> Result<()> {
564        let yml = r#"
565name: t
566version: "1"
567models_dir: models
568target_path: lake/gold
569materialize:
570  ref_strategy: memtable
571  memtable_max_rows: 10000
572"#;
573        let cfg: RbtProjectConfig = serde_yaml::from_str(yml)?;
574        assert_eq!(cfg.materialize.ref_strategy, RefStrategy::Memtable);
575        assert_eq!(cfg.materialize.memtable_max_rows, 10_000);
576        assert_eq!(
577            cfg.materialize.choose_ref_backend(9_999),
578            RefBackend::MemTable
579        );
580        assert_eq!(
581            cfg.materialize.choose_ref_backend(10_000),
582            RefBackend::LakeFile
583        );
584        Ok(())
585    }
586
587    #[test]
588    fn parse_project_without_materialize_uses_defaults() -> Result<()> {
589        let yml = r#"
590name: t
591version: "1"
592models_dir: models
593target_path: lake/gold
594"#;
595        let cfg: RbtProjectConfig = serde_yaml::from_str(yml)?;
596        assert_eq!(cfg.materialize, MaterializeConfig::default());
597        Ok(())
598    }
599
600    #[test]
601    fn parse_memtable_without_max_rows_defaults_cutoff() -> Result<()> {
602        let yml = r#"
603name: t
604version: "1"
605models_dir: models
606target_path: lake/gold
607materialize:
608  ref_strategy: memtable
609"#;
610        let cfg: RbtProjectConfig = serde_yaml::from_str(yml)?;
611        assert_eq!(cfg.materialize.ref_strategy, RefStrategy::Memtable);
612        assert_eq!(cfg.materialize.memtable_max_rows, DEFAULT_MEMTABLE_MAX_ROWS);
613        Ok(())
614    }
615
616    #[test]
617    fn absolute_layer_target_not_nested_under_project() -> Result<()> {
618        let yml = r#"
619name: multi_root_demo
620version: "1"
621models_dir: models
622target_path: /mnt/datalake/acme/nonprod/lake_us/lake/gold
623layers:
624  staging:
625    path: models/staging
626    target_path: /mnt/datalake/acme/nonprod/lake_us/lake/silver/stage
627    default_format: parquet
628"#;
629        let cfg: RbtProjectConfig = serde_yaml::from_str(yml)?;
630        let project = Path::new("/home/dev/rbt_projects/demo");
631        let stg =
632            cfg.resolve_model_target_path(project, "stg_events", ModelLayer::Staging, "parquet")?;
633        assert_eq!(
634            stg,
635            PathBuf::from(
636                "/mnt/datalake/acme/nonprod/lake_us/lake/silver/stage/stg_events.parquet"
637            )
638        );
639        assert!(!stg.starts_with(project));
640        Ok(())
641    }
642
643    #[test]
644    fn multi_root_template_in_layer_target() -> Result<()> {
645        let yml = r#"
646name: multi_root_demo
647version: "1"
648models_dir: models
649target_path: $nonprod_lake/gold
650roots:
651  nonprod_lake: /mnt/datalake/acme/nonprod/lake_us/lake
652layers:
653  staging:
654    path: models/staging
655    target_path: $nonprod_lake/silver/stage
656    default_format: parquet
657"#;
658        let cfg: RbtProjectConfig = serde_yaml::from_str(yml)?;
659        let project = Path::new("/home/dev/proj");
660        let dir = cfg.resolve_layer_target_dir(project, ModelLayer::Staging)?;
661        assert_eq!(
662            dir,
663            PathBuf::from("/mnt/datalake/acme/nonprod/lake_us/lake/silver/stage")
664        );
665        Ok(())
666    }
667
668    #[test]
669    fn bad_root_in_layer_target_is_error() {
670        let yml = r#"
671name: t
672version: "1"
673models_dir: models
674target_path: lake/gold
675layers:
676  staging:
677    path: models/staging
678    target_path: $missing_root/silver
679    default_format: parquet
680"#;
681        let cfg: RbtProjectConfig = serde_yaml::from_str(yml).unwrap();
682        let err = cfg
683            .resolve_layer_target_dir(Path::new("/proj"), ModelLayer::Staging)
684            .unwrap_err()
685            .to_string();
686        assert!(err.contains("E_RBT_ROOT_UNKNOWN") || err.contains("E_RBT_LAYER_PATH"));
687    }
688
689    #[test]
690    fn scan_config_defaults_protobuf_cap() {
691        let cfg = ScanConfig::default();
692        assert_eq!(
693            cfg.protobuf_max_payload_bytes,
694            DEFAULT_PROTOBUF_MAX_PAYLOAD_BYTES
695        );
696        assert_eq!(DEFAULT_PROTOBUF_MAX_PAYLOAD_BYTES, 1024 * 1024 * 1024);
697    }
698
699    #[test]
700    fn scan_config_override_from_yml() -> Result<()> {
701        let yml = r#"
702name: t
703version: "1"
704models_dir: models
705target_path: lake/gold
706scan:
707  protobuf_max_payload_bytes: 4096
708"#;
709        let cfg: RbtProjectConfig = serde_yaml::from_str(yml)?;
710        assert_eq!(cfg.scan.protobuf_max_payload_bytes, 4096);
711        // omit scan: → default 1 GiB
712        let yml2 = r#"
713name: t
714version: "1"
715models_dir: models
716target_path: lake/gold
717"#;
718        let cfg2: RbtProjectConfig = serde_yaml::from_str(yml2)?;
719        assert_eq!(
720            cfg2.scan.protobuf_max_payload_bytes,
721            DEFAULT_PROTOBUF_MAX_PAYLOAD_BYTES
722        );
723        Ok(())
724    }
725
726    /// Workspace examples stay loadable / 0.3.7-shaped (roots + defaults).
727    #[test]
728    fn load_workspace_example_projects() -> Result<()> {
729        let manifest = PathBuf::from(env!("CARGO_MANIFEST_DIR"));
730        let repo = manifest.join("../..");
731        for (rel, name, expect_root) in [
732            ("examples/smoke_fixture", "smoke_fixture", "lake"),
733            ("examples/full_e2e_rbt_example", "market_bars", "lake"),
734        ] {
735            let dir = repo.join(rel);
736            if !dir.join("rbt_project.yml").is_file() {
737                // crates.io source package may omit large e2e bronze; skip if missing
738                continue;
739            }
740            let cfg = RbtProjectConfig::load(&dir)?;
741            assert_eq!(cfg.name, name, "example {rel}");
742            assert_eq!(
743                cfg.roots.get("lake").map(String::as_str),
744                Some(expect_root),
745                "example {rel} should declare roots.lake"
746            );
747            assert_eq!(cfg.materialize.ref_strategy, RefStrategy::Parquet);
748            assert_eq!(
749                cfg.scan.protobuf_max_payload_bytes,
750                DEFAULT_PROTOBUF_MAX_PAYLOAD_BYTES
751            );
752            let silver = cfg.resolve_layer_target_dir(&dir, ModelLayer::Staging)?;
753            assert!(
754                silver.ends_with("lake/silver") || silver.ends_with("lake\\silver"),
755                "staging target for {rel}: {}",
756                silver.display()
757            );
758            // DAG builds for smoke always; e2e only if models present
759            if dir.join("models").is_dir() {
760                let dag = cfg.build_dag(&dir, None)?;
761                assert!(
762                    dag.graph.node_count() >= 3,
763                    "example {rel} expected ≥3 models, got {}",
764                    dag.graph.node_count()
765                );
766            }
767        }
768        Ok(())
769    }
770}