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