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
11pub const DEFAULT_MEMTABLE_MAX_ROWS: usize = 50_000;
13
14#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
19#[serde(rename_all = "snake_case")]
20pub enum RefStrategy {
21 #[default]
23 #[serde(alias = "parquet_reread", alias = "lake", alias = "file")]
24 Parquet,
25 #[serde(alias = "mem_table", alias = "memory", alias = "arc")]
27 Memtable,
28}
29
30#[derive(Debug, Clone, Copy, PartialEq, Eq)]
32pub enum RefBackend {
33 MemTable,
35 LakeFile,
37}
38
39#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
44#[serde(rename_all = "snake_case")]
45pub enum MaterializeMode {
46 #[default]
48 #[serde(alias = "streaming")]
49 Stream,
50 #[serde(alias = "batch", alias = "legacy")]
52 Collect,
53}
54
55pub const DEFAULT_MAX_ROW_GROUP_ROWS: usize = 1_000_000;
57pub const DEFAULT_MAX_ROW_GROUP_BYTES: usize = 128 * 1024 * 1024;
59
60#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
62#[serde(rename_all = "snake_case")]
63pub enum IcebergWriteMode {
64 #[default]
66 #[serde(alias = "catalog_commit", alias = "sor")]
67 Catalog,
68 #[serde(alias = "fs", alias = "layout")]
70 Filesystem,
71}
72
73#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
75pub struct IcebergConfig {
76 #[serde(default)]
78 pub mode: IcebergWriteMode,
79 #[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#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
102pub struct MaterializeConfig {
103 #[serde(default)]
105 pub mode: MaterializeMode,
106 #[serde(default)]
108 pub ref_strategy: RefStrategy,
109 #[serde(default = "default_memtable_max_rows")]
111 pub memtable_max_rows: usize,
112 #[serde(default = "default_max_row_group_rows")]
114 pub max_row_group_rows: usize,
115 #[serde(default = "default_max_row_group_bytes")]
117 pub max_row_group_bytes: usize,
118 #[serde(default)]
120 pub iceberg: IcebergConfig,
121 #[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
153pub const DEFAULT_PROTOBUF_MAX_PAYLOAD_BYTES: u64 = 1024 * 1024 * 1024;
155
156#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
158pub struct ScanConfig {
159 #[serde(default = "default_protobuf_max_payload_bytes")]
163 pub protobuf_max_payload_bytes: u64,
164 #[serde(default = "default_spill_arrow_ipc")]
169 pub spill_arrow_ipc: bool,
170 #[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 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 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 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#[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#[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 #[serde(default)]
249 pub layers: HashMap<String, LayerConfig>,
250 #[serde(default)]
252 pub materialize: MaterializeConfig,
253 #[serde(default)]
255 pub scan: ScanConfig,
256 #[serde(default)]
260 pub roots: HashMap<String, String>,
261}
262
263impl Default for RbtProjectConfig {
264 fn default() -> Self {
265 let mut layers = HashMap::new();
266 layers.insert(
267 "staging".to_string(),
268 LayerConfig {
269 path: PathBuf::from("models/staging"),
270 target_path: PathBuf::from("lake/silver"),
271 default_format: Some("parquet".to_string()),
272 },
273 );
274 layers.insert(
275 "transforms".to_string(),
276 LayerConfig {
277 path: PathBuf::from("models/transforms"),
278 target_path: PathBuf::from("lake/gold"),
279 default_format: Some("parquet".to_string()),
280 },
281 );
282 layers.insert(
283 "marts".to_string(),
284 LayerConfig {
285 path: PathBuf::from("models/marts"),
286 target_path: PathBuf::from("lake/gold"),
287 default_format: Some("parquet_and_iceberg".to_string()),
288 },
289 );
290
291 Self {
292 name: "rbt_project".to_string(),
293 version: "1.0.0".to_string(),
294 models_dir: PathBuf::from("models"),
295 target_path: PathBuf::from("lake/gold"),
296 layers,
297 materialize: MaterializeConfig::default(),
298 scan: ScanConfig::default(),
299 roots: HashMap::new(),
300 }
301 }
302}
303
304impl RbtProjectConfig {
305 pub fn load(project_dir: &Path) -> Result<Self> {
307 let project_file = project_dir.join("rbt_project.yml");
308 if project_file.exists() {
309 let content = fs::read_to_string(&project_file).with_context(|| {
310 format!(
311 "E_RBT_PROJECT_LOAD: cannot read project file {}",
312 project_file.display()
313 )
314 })?;
315 let mut config: RbtProjectConfig = serde_yaml::from_str(&content).with_context(|| {
316 format!(
317 "E_RBT_PROJECT_LOAD: failed to parse {}. \
318 Check required keys (name, version, models_dir, target_path) and \
319 optional materialize:/scan:/roots:/layers blocks.",
320 project_file.display()
321 )
322 })?;
323
324 let defaults = Self::default();
325 for (key, val) in defaults.layers {
326 config.layers.entry(key).or_insert(val);
327 }
328 Ok(config)
329 } else {
330 Ok(Self::default())
331 }
332 }
333
334 pub fn resolve_path(&self, project_dir: &Path, configured: &str) -> Result<PathBuf> {
336 resolve_project_path(project_dir, configured, &self.roots)
337 }
338
339 pub fn resolve_layer_target_dir(
341 &self,
342 project_dir: &Path,
343 layer: ModelLayer,
344 ) -> Result<PathBuf> {
345 let layer_key = match layer {
346 ModelLayer::Staging => "staging",
347 ModelLayer::Transform => "transforms",
348 ModelLayer::Mart => "marts",
349 };
350 if let Some(layer_cfg) = self.layers.get(layer_key) {
351 resolve_configured_path(project_dir, &layer_cfg.target_path, &self.roots)
352 } else {
353 resolve_configured_path(project_dir, &self.target_path, &self.roots)
354 }
355 }
356
357 pub fn resolve_model_target_path(
362 &self,
363 project_dir: &Path,
364 model_name: &str,
365 layer: ModelLayer,
366 ext: &str,
367 ) -> Result<PathBuf> {
368 let dir = self
369 .resolve_layer_target_dir(project_dir, layer)
370 .with_context(|| {
371 format!(
372 "E_RBT_MODEL_TARGET: cannot resolve output directory for model '{model_name}' \
373 (layer={layer:?}). Check `layers.*.target_path`, top-level `target_path`, and \
374 `roots:` in rbt_project.yml."
375 )
376 })?;
377 Ok(dir.join(format!("{model_name}.{ext}")))
378 }
379
380 pub fn resolve_model_target_dir(
382 &self,
383 project_dir: &Path,
384 model_name: &str,
385 layer: ModelLayer,
386 ) -> Result<PathBuf> {
387 let dir = self
388 .resolve_layer_target_dir(project_dir, layer)
389 .with_context(|| {
390 format!(
391 "E_RBT_MODEL_TARGET: cannot resolve table directory for model '{model_name}' \
392 (layer={layer:?}). Check layer target_path and roots:."
393 )
394 })?;
395 Ok(dir.join(model_name))
396 }
397
398 pub fn build_dag(
400 &self,
401 project_dir: &Path,
402 cli_format_override: Option<OutputFormat>,
403 ) -> Result<ModelDag> {
404 let models_dir = project_dir.join(&self.models_dir);
405 let mut dag = ModelDag::new();
406
407 if !models_dir.exists() {
408 let default_fmt = cli_format_override.unwrap_or(OutputFormat::Parquet);
409 dag.add_model_with_format(
410 "stg_users",
411 "SELECT 1 AS id, 'Alice' AS name, 'admin' AS role",
412 Materialization::Table,
413 default_fmt,
414 None,
415 "",
416 )?;
417 dag.build_graph()?;
418 return Ok(dag);
419 }
420
421 let mut model_count = 0;
422 for entry in WalkDir::new(&models_dir).into_iter().filter_map(|e| e.ok()) {
423 let path = entry.path();
424 if path.is_file() && path.extension().is_some_and(|ext| ext == "sql") {
425 let stem = path.file_stem().and_then(|s| s.to_str()).with_context(|| {
426 format!(
427 "E_RBT_MODEL_NAME: invalid file stem for model path {}",
428 path.display()
429 )
430 })?;
431 let raw_sql = fs::read_to_string(path).with_context(|| {
432 format!(
433 "E_RBT_MODEL_IO: failed reading model SQL {}",
434 path.display()
435 )
436 })?;
437
438 let layer = ModelLayer::from_name(stem);
439 let format = cli_format_override.clone().unwrap_or_else(|| {
440 let layer_key = match layer {
441 ModelLayer::Staging => "staging",
442 ModelLayer::Transform => "transforms",
443 ModelLayer::Mart => "marts",
444 };
445 if let Some(l_cfg) = self.layers.get(layer_key) {
446 match l_cfg.default_format.as_deref() {
447 Some("parquet") => OutputFormat::Parquet,
448 Some("jsonl") => OutputFormat::Jsonl,
449 Some("csv") => OutputFormat::Csv,
450 Some("iceberg") => OutputFormat::Iceberg,
451 Some("parquet_and_iceberg") => OutputFormat::ParquetAndIceberg,
452 _ => OutputFormat::Parquet,
453 }
454 } else {
455 OutputFormat::Parquet
456 }
457 });
458
459 let target_file_path = match format {
460 OutputFormat::Iceberg => self
461 .resolve_model_target_dir(project_dir, stem, layer)
462 .with_context(|| {
463 format!(
464 "E_RBT_MODEL_TARGET: model '{stem}' (Iceberg) — \
465 failed resolving layer target. \
466 layer={layer:?}; project={}",
467 project_dir.display()
468 )
469 })?,
470 OutputFormat::Parquet
471 | OutputFormat::ParquetAndIceberg
472 | OutputFormat::ZeroCopyClone => self
473 .resolve_model_target_path(project_dir, stem, layer, "parquet")
474 .with_context(|| {
475 format!(
476 "E_RBT_MODEL_TARGET: model '{stem}' (parquet) — \
477 failed resolving layer target. \
478 layer={layer:?}; project={}",
479 project_dir.display()
480 )
481 })?,
482 OutputFormat::Jsonl => self
483 .resolve_model_target_path(project_dir, stem, layer, "jsonl")
484 .with_context(|| {
485 format!(
486 "E_RBT_MODEL_TARGET: model '{stem}' (jsonl) — \
487 failed resolving layer target"
488 )
489 })?,
490 OutputFormat::Csv => self
491 .resolve_model_target_path(project_dir, stem, layer, "csv")
492 .with_context(|| {
493 format!(
494 "E_RBT_MODEL_TARGET: model '{stem}' (csv) — \
495 failed resolving layer target"
496 )
497 })?,
498 };
499
500 dag.add_model_with_format(
501 stem,
502 &raw_sql,
503 Materialization::Table,
504 format,
505 Some(target_file_path.to_string_lossy().to_string()),
506 "",
507 )?;
508 model_count += 1;
509 }
510 }
511
512 if model_count == 0 {
513 let default_fmt = cli_format_override.unwrap_or(OutputFormat::Parquet);
514 dag.add_model_with_format(
515 "stg_users",
516 "SELECT 1 AS id, 'Alice' AS name, 'admin' AS role",
517 Materialization::Table,
518 default_fmt,
519 None,
520 "",
521 )?;
522 }
523
524 dag.build_graph()?;
525 Ok(dag)
526 }
527}
528
529#[cfg(test)]
530mod tests {
531 use super::*;
532
533 #[test]
534 fn test_layer_target_path_resolution() -> Result<()> {
535 let config = RbtProjectConfig::default();
536 let project_dir = Path::new("/tmp/test_project");
537
538 let stg_path = config.resolve_model_target_path(
539 project_dir,
540 "stg_trades",
541 ModelLayer::Staging,
542 "parquet",
543 )?;
544 assert_eq!(stg_path, project_dir.join("lake/silver/stg_trades.parquet"));
545
546 let tf_path = config.resolve_model_target_path(
547 project_dir,
548 "tf_1m_bars",
549 ModelLayer::Transform,
550 "parquet",
551 )?;
552 assert_eq!(tf_path, project_dir.join("lake/gold/tf_1m_bars.parquet"));
553
554 let mart_path = config.resolve_model_target_path(
555 project_dir,
556 "fact_1d_bars",
557 ModelLayer::Mart,
558 "parquet",
559 )?;
560 assert_eq!(
561 mart_path,
562 project_dir.join("lake/gold/fact_1d_bars.parquet")
563 );
564
565 Ok(())
566 }
567
568 #[test]
569 fn materialize_defaults_to_parquet_reread() {
570 let cfg = MaterializeConfig::default();
571 assert_eq!(cfg.ref_strategy, RefStrategy::Parquet);
572 assert_eq!(cfg.mode, MaterializeMode::Stream);
573 assert_eq!(cfg.memtable_max_rows, DEFAULT_MEMTABLE_MAX_ROWS);
574 assert_eq!(cfg.max_row_group_rows, DEFAULT_MAX_ROW_GROUP_ROWS);
575 assert_eq!(cfg.choose_ref_backend(0), RefBackend::LakeFile);
576 assert_eq!(cfg.choose_ref_backend(1_000_000), RefBackend::LakeFile);
577 }
578
579 #[test]
580 fn materialize_mode_from_yaml() -> Result<()> {
581 let yml = r#"
582name: t
583version: "1"
584models_dir: models
585target_path: lake/gold
586materialize:
587 mode: collect
588 max_row_group_rows: 1000
589"#;
590 let cfg: RbtProjectConfig = serde_yaml::from_str(yml)?;
591 assert_eq!(cfg.materialize.mode, MaterializeMode::Collect);
592 assert_eq!(cfg.materialize.max_row_group_rows, 1000);
593 Ok(())
594 }
595
596 #[test]
597 fn materialize_memtable_respects_cutoff() {
598 let cfg = MaterializeConfig {
599 ref_strategy: RefStrategy::Memtable,
600 memtable_max_rows: 50_000,
601 ..Default::default()
602 };
603 assert_eq!(cfg.choose_ref_backend(49_999), RefBackend::MemTable);
604 assert_eq!(cfg.choose_ref_backend(50_000), RefBackend::LakeFile);
605 assert_eq!(cfg.choose_ref_backend(50_001), RefBackend::LakeFile);
606 }
607
608 #[test]
609 fn parse_materialize_block_from_yaml() -> Result<()> {
610 let yml = r#"
611name: t
612version: "1"
613models_dir: models
614target_path: lake/gold
615materialize:
616 ref_strategy: memtable
617 memtable_max_rows: 10000
618"#;
619 let cfg: RbtProjectConfig = serde_yaml::from_str(yml)?;
620 assert_eq!(cfg.materialize.ref_strategy, RefStrategy::Memtable);
621 assert_eq!(cfg.materialize.memtable_max_rows, 10_000);
622 assert_eq!(
623 cfg.materialize.choose_ref_backend(9_999),
624 RefBackend::MemTable
625 );
626 assert_eq!(
627 cfg.materialize.choose_ref_backend(10_000),
628 RefBackend::LakeFile
629 );
630 Ok(())
631 }
632
633 #[test]
634 fn parse_project_without_materialize_uses_defaults() -> Result<()> {
635 let yml = r#"
636name: t
637version: "1"
638models_dir: models
639target_path: lake/gold
640"#;
641 let cfg: RbtProjectConfig = serde_yaml::from_str(yml)?;
642 assert_eq!(cfg.materialize, MaterializeConfig::default());
643 Ok(())
644 }
645
646 #[test]
647 fn parse_memtable_without_max_rows_defaults_cutoff() -> Result<()> {
648 let yml = r#"
649name: t
650version: "1"
651models_dir: models
652target_path: lake/gold
653materialize:
654 ref_strategy: memtable
655"#;
656 let cfg: RbtProjectConfig = serde_yaml::from_str(yml)?;
657 assert_eq!(cfg.materialize.ref_strategy, RefStrategy::Memtable);
658 assert_eq!(cfg.materialize.memtable_max_rows, DEFAULT_MEMTABLE_MAX_ROWS);
659 Ok(())
660 }
661
662 #[test]
663 fn absolute_layer_target_not_nested_under_project() -> Result<()> {
664 let yml = r#"
665name: multi_root_demo
666version: "1"
667models_dir: models
668target_path: /mnt/datalake/acme/nonprod/lake_us/lake/gold
669layers:
670 staging:
671 path: models/staging
672 target_path: /mnt/datalake/acme/nonprod/lake_us/lake/silver/stage
673 default_format: parquet
674"#;
675 let cfg: RbtProjectConfig = serde_yaml::from_str(yml)?;
676 let project = Path::new("/home/dev/rbt_projects/demo");
677 let stg =
678 cfg.resolve_model_target_path(project, "stg_events", ModelLayer::Staging, "parquet")?;
679 assert_eq!(
680 stg,
681 PathBuf::from(
682 "/mnt/datalake/acme/nonprod/lake_us/lake/silver/stage/stg_events.parquet"
683 )
684 );
685 assert!(!stg.starts_with(project));
686 Ok(())
687 }
688
689 #[test]
690 fn multi_root_template_in_layer_target() -> Result<()> {
691 let yml = r#"
692name: multi_root_demo
693version: "1"
694models_dir: models
695target_path: $nonprod_lake/gold
696roots:
697 nonprod_lake: /mnt/datalake/acme/nonprod/lake_us/lake
698layers:
699 staging:
700 path: models/staging
701 target_path: $nonprod_lake/silver/stage
702 default_format: parquet
703"#;
704 let cfg: RbtProjectConfig = serde_yaml::from_str(yml)?;
705 let project = Path::new("/home/dev/proj");
706 let dir = cfg.resolve_layer_target_dir(project, ModelLayer::Staging)?;
707 assert_eq!(
708 dir,
709 PathBuf::from("/mnt/datalake/acme/nonprod/lake_us/lake/silver/stage")
710 );
711 Ok(())
712 }
713
714 #[test]
715 fn bad_root_in_layer_target_is_error() {
716 let yml = r#"
717name: t
718version: "1"
719models_dir: models
720target_path: lake/gold
721layers:
722 staging:
723 path: models/staging
724 target_path: $missing_root/silver
725 default_format: parquet
726"#;
727 let cfg: RbtProjectConfig = serde_yaml::from_str(yml).unwrap();
728 let err = cfg
729 .resolve_layer_target_dir(Path::new("/proj"), ModelLayer::Staging)
730 .unwrap_err()
731 .to_string();
732 assert!(err.contains("E_RBT_ROOT_UNKNOWN") || err.contains("E_RBT_LAYER_PATH"));
733 }
734
735 #[test]
736 fn scan_config_defaults_protobuf_cap() {
737 let cfg = ScanConfig::default();
738 assert_eq!(
739 cfg.protobuf_max_payload_bytes,
740 DEFAULT_PROTOBUF_MAX_PAYLOAD_BYTES
741 );
742 assert_eq!(DEFAULT_PROTOBUF_MAX_PAYLOAD_BYTES, 1024 * 1024 * 1024);
743 }
744
745 #[test]
746 fn scan_config_override_from_yml() -> Result<()> {
747 let yml = r#"
748name: t
749version: "1"
750models_dir: models
751target_path: lake/gold
752scan:
753 protobuf_max_payload_bytes: 4096
754"#;
755 let cfg: RbtProjectConfig = serde_yaml::from_str(yml)?;
756 assert_eq!(cfg.scan.protobuf_max_payload_bytes, 4096);
757 let yml2 = r#"
759name: t
760version: "1"
761models_dir: models
762target_path: lake/gold
763"#;
764 let cfg2: RbtProjectConfig = serde_yaml::from_str(yml2)?;
765 assert_eq!(
766 cfg2.scan.protobuf_max_payload_bytes,
767 DEFAULT_PROTOBUF_MAX_PAYLOAD_BYTES
768 );
769 Ok(())
770 }
771
772 #[test]
774 fn load_workspace_example_projects() -> Result<()> {
775 let manifest = PathBuf::from(env!("CARGO_MANIFEST_DIR"));
776 let repo = manifest.join("../..");
777 for (rel, name, expect_root) in [
778 ("examples/smoke_fixture", "smoke_fixture", "lake"),
779 ("examples/full_e2e_rbt_example", "market_bars", "lake"),
780 ] {
781 let dir = repo.join(rel);
782 if !dir.join("rbt_project.yml").is_file() {
783 continue;
785 }
786 let cfg = RbtProjectConfig::load(&dir)?;
787 assert_eq!(cfg.name, name, "example {rel}");
788 assert_eq!(
789 cfg.roots.get("lake").map(String::as_str),
790 Some(expect_root),
791 "example {rel} should declare roots.lake"
792 );
793 assert_eq!(cfg.materialize.ref_strategy, RefStrategy::Parquet);
794 assert_eq!(
795 cfg.scan.protobuf_max_payload_bytes,
796 DEFAULT_PROTOBUF_MAX_PAYLOAD_BYTES
797 );
798 let silver = cfg.resolve_layer_target_dir(&dir, ModelLayer::Staging)?;
799 assert!(
800 silver.ends_with("lake/silver") || silver.ends_with("lake\\silver"),
801 "staging target for {rel}: {}",
802 silver.display()
803 );
804 if dir.join("models").is_dir() {
806 let dag = cfg.build_dag(&dir, None)?;
807 assert!(
808 dag.graph.node_count() >= 3,
809 "example {rel} expected ≥3 models, got {}",
810 dag.graph.node_count()
811 );
812 }
813 }
814 Ok(())
815 }
816}