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