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