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/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 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 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 pub fn resolve_path(&self, project_dir: &Path, configured: &str) -> Result<PathBuf> {
343 resolve_project_path(project_dir, configured, &self.roots)
344 }
345
346 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 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 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 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 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 #[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 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 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}