Skip to main content

rbt/engine/
bronze.rs

1//! Frontmatter-driven bronze source registration.
2//!
3//! ## Architecture
4//!
5//! * **Path A (DataFusion listing / external tables)** — Parquet, CSV, JSON/JSONL (no
6//!   jshift projection), Arrow IPC file: register via DataFusion native readers, then
7//!   wrap the resulting provider in [`BronzeTableProvider`]. Listing predicate
8//!   pushdown is available on this path.
9//! * **Path B (scan → MemTable)** — used when rbt must apply its own filters or
10//!   inject path-derived columns. **Any non-empty `path_glob` forces Path B** (as do
11//!   `partition_by` / `require_partitions` / `inject_source_path` / `force_scan` and
12//!   formats that require scan: log, txt, toml, Arrow IPC stream, protobuf).
13//!   DataFusion directory listing pushdown is **disabled** for that source by design.
14//!
15//! [`BronzeTableProvider`] is intentionally thin: it delegates scan/schema to the
16//! inner provider and carries bronze metadata for lineage / debugging.
17
18use crate::core::dag::{ModelDag, ModelNode};
19use crate::core::frontmatter::{SourceFormat, StagingFrontmatter};
20use crate::core::receipt::apply_scope_to_frontmatter;
21use crate::core::run_scope::{OnMissing, RunScope};
22use crate::scan::{LakeScanner, ScanRequest};
23use anyhow::{bail, Context, Result};
24use arrow::datatypes::SchemaRef;
25use arrow::record_batch::RecordBatch;
26use async_trait::async_trait;
27use datafusion::catalog::Session;
28use datafusion::catalog::TableProvider;
29use datafusion::common::TableReference;
30use datafusion::datasource::MemTable;
31use datafusion::error::Result as DFResult;
32use datafusion::execution::context::SessionContext;
33use datafusion::execution::options::ArrowReadOptions;
34use datafusion::logical_expr::{Expr, TableType};
35use datafusion::physical_plan::ExecutionPlan;
36use datafusion::prelude::{CsvReadOptions, JsonReadOptions, ParquetReadOptions};
37use std::any::Any;
38use std::collections::HashSet;
39use std::path::{Path, PathBuf};
40use std::sync::Arc;
41
42/// Metadata retained on the bronze provider for debugging and future lineage.
43#[derive(Debug, Clone)]
44pub struct BronzeSourceMeta {
45    pub model_name: String,
46    pub source_schema: String,
47    pub source_table: String,
48    pub format: SourceFormat,
49    pub scan_path: PathBuf,
50    pub registration_mode: BronzeRegistrationMode,
51}
52
53#[derive(Debug, Clone, Copy, PartialEq, Eq)]
54pub enum BronzeRegistrationMode {
55    /// Inner provider is a DataFusion listing / external table.
56    DataFusionListing,
57    /// Inner provider is a MemTable filled by `rbt::scan` (small / non-spill formats).
58    ScanMemTable,
59    /// Arrow IPC (etc.) spilled file-by-file to Parquet then listed — bounded peak RAM.
60    ScanSpillParquet,
61}
62
63/// Thin `TableProvider` wrapper around a DataFusion listing table or MemTable.
64#[derive(Debug)]
65pub struct BronzeTableProvider {
66    pub meta: BronzeSourceMeta,
67    inner: Arc<dyn TableProvider>,
68}
69
70impl BronzeTableProvider {
71    pub fn wrap(inner: Arc<dyn TableProvider>, meta: BronzeSourceMeta) -> Self {
72        Self { meta, inner }
73    }
74
75    pub fn inner(&self) -> &Arc<dyn TableProvider> {
76        &self.inner
77    }
78}
79
80#[async_trait]
81impl TableProvider for BronzeTableProvider {
82    fn as_any(&self) -> &dyn Any {
83        self
84    }
85
86    fn schema(&self) -> SchemaRef {
87        self.inner.schema()
88    }
89
90    fn table_type(&self) -> TableType {
91        self.inner.table_type()
92    }
93
94    async fn scan(
95        &self,
96        state: &dyn Session,
97        projection: Option<&Vec<usize>>,
98        filters: &[Expr],
99        limit: Option<usize>,
100    ) -> DFResult<Arc<dyn ExecutionPlan>> {
101        self.inner.scan(state, projection, filters, limit).await
102    }
103
104    fn supports_filters_pushdown(
105        &self,
106        filters: &[&Expr],
107    ) -> DFResult<Vec<datafusion::logical_expr::TableProviderFilterPushDown>> {
108        self.inner.supports_filters_pushdown(filters)
109    }
110}
111
112/// Registers all bronze sources declared by model frontmatter into `ctx`.
113///
114/// Idempotent per `(schema, table)` within a single run (tracked by `registered`).
115/// Uses `config.roots` and `config.scan` (no re-read of yml per model).
116/// Optional [`RunScope`] expands `{var}` templates and partition binds (P5a).
117pub async fn register_bronze_sources_for_dag(
118    ctx: &SessionContext,
119    dag: &ModelDag,
120    project_dir: &Path,
121    registered: &mut HashSet<(String, String)>,
122    config: &crate::core::project::RbtProjectConfig,
123) -> Result<usize> {
124    register_bronze_sources_for_dag_scoped(ctx, dag, project_dir, registered, config, None).await
125}
126
127/// Like [`register_bronze_sources_for_dag`] with an explicit run scope.
128pub async fn register_bronze_sources_for_dag_scoped(
129    ctx: &SessionContext,
130    dag: &ModelDag,
131    project_dir: &Path,
132    registered: &mut HashSet<(String, String)>,
133    config: &crate::core::project::RbtProjectConfig,
134    scope: Option<&RunScope>,
135) -> Result<usize> {
136    let mut count = 0;
137    for idx in dag.graph.node_indices() {
138        let node = &dag.graph[idx];
139        if let Some(n) =
140            register_bronze_for_model_scoped(ctx, node, project_dir, registered, config, scope)
141                .await?
142        {
143            count += n;
144        }
145    }
146    Ok(count)
147}
148
149/// Register bronze for a single model if it has a scan contract.
150pub async fn register_bronze_for_model(
151    ctx: &SessionContext,
152    node: &ModelNode,
153    project_dir: &Path,
154    registered: &mut HashSet<(String, String)>,
155    config: &crate::core::project::RbtProjectConfig,
156) -> Result<Option<usize>> {
157    register_bronze_for_model_scoped(ctx, node, project_dir, registered, config, None).await
158}
159
160/// Register bronze with optional run scope (partition binds + var templates + on_missing).
161pub async fn register_bronze_for_model_scoped(
162    ctx: &SessionContext,
163    node: &ModelNode,
164    project_dir: &Path,
165    registered: &mut HashSet<(String, String)>,
166    config: &crate::core::project::RbtProjectConfig,
167    scope: Option<&RunScope>,
168) -> Result<Option<usize>> {
169    let Some(fm_raw) = node.frontmatter.as_ref() else {
170        return Ok(None);
171    };
172    if !fm_raw.has_scan_contract() {
173        return Ok(None);
174    }
175
176    let default_scope = RunScope::default();
177    let scope = scope.unwrap_or(&default_scope);
178    let fm = apply_scope_to_frontmatter(fm_raw, scope);
179    let fm = &fm;
180
181    let (schema_name, table_name) = ModelDag::bronze_source_ident(node).with_context(|| {
182        format!(
183            "model '{}': frontmatter has scan_path but no source identity \
184             (add source() in SQL or source_name/source_table in frontmatter)",
185            node.name
186        )
187    })?;
188
189    let key = (schema_name.clone(), table_name.clone());
190    if registered.contains(&key) {
191        tracing::debug!(
192            "Bronze source {}.{} already registered; skipping model '{}'",
193            schema_name,
194            table_name,
195            node.name
196        );
197        return Ok(None);
198    }
199
200    ensure_schema(ctx, &schema_name).await?;
201
202    let format = fm
203        .resolve_format()
204        .with_context(|| format!("model '{}': cannot resolve source_format", node.name))?;
205
206    let raw_scan = fm.scan_path.as_deref().unwrap();
207    let resolved = crate::core::paths::resolve_project_path(project_dir, raw_scan, &config.roots)
208        .with_context(|| {
209        format!(
210            "E_RBT_BRONZE_PATH: model '{}': cannot resolve scan_path '{}'. \
211                     Check absolute paths and `roots:` templates in rbt_project.yml.",
212            node.name, raw_scan
213        )
214    })?;
215    let on_missing = fm.on_missing_policy();
216    if !resolved.exists() && !crate::core::frontmatter::is_remote_uri(raw_scan) {
217        if on_missing == OnMissing::Empty {
218            let provider = empty_memtable(fm, scope).with_context(|| {
219                format!(
220                    "model '{}': on_missing empty frame failed (scan_path missing)",
221                    node.name
222                )
223            })?;
224            return finish_register(
225                ctx,
226                registered,
227                key,
228                schema_name,
229                table_name,
230                node,
231                format,
232                resolved,
233                provider,
234                BronzeRegistrationMode::ScanMemTable,
235            )
236            .await;
237        }
238        bail!(
239            "E_RBT_BRONZE_SCAN_PATH_NOT_FOUND: model '{}': bronze scan_path does not exist: {} \
240             (resolved {}). Hint: verify the lake path and `$root` expansion, \
241             or set on_missing: empty for optional artifact families.",
242            node.name,
243            raw_scan,
244            resolved.display()
245        );
246    }
247
248    let path_str = resolved.to_string_lossy().to_string();
249    let use_scan = should_use_scan_path(fm, format);
250
251    let (inner, mode) = if use_scan {
252        if should_spill_to_parquet(format, config) {
253            match scan_spill_to_listing(
254                ctx,
255                project_dir,
256                fm,
257                format,
258                config,
259                &schema_name,
260                &table_name,
261            )
262            .await
263            {
264                Ok(provider) => (provider, BronzeRegistrationMode::ScanSpillParquet),
265                Err(e) if on_missing == OnMissing::Empty && is_empty_or_missing_err(&e) => {
266                    (
267                        empty_memtable(fm, scope)?,
268                        BronzeRegistrationMode::ScanMemTable,
269                    )
270                }
271                Err(e) => {
272                    return Err(e).with_context(|| {
273                        format!(
274                            "model '{}': bronze spill→parquet failed (format={})",
275                            node.name, format
276                        )
277                    })
278                }
279            }
280        } else {
281            match scan_to_memtable(project_dir, fm, format, config).await {
282                Ok(provider) => (provider, BronzeRegistrationMode::ScanMemTable),
283                Err(e) if on_missing == OnMissing::Empty && is_empty_or_missing_err(&e) => {
284                    (
285                        empty_memtable(fm, scope)?,
286                        BronzeRegistrationMode::ScanMemTable,
287                    )
288                }
289                Err(e) => {
290                    return Err(e).with_context(|| format!("model '{}': bronze scan failed", node.name))
291                }
292            }
293        }
294    } else {
295        let provider = listing_table_provider(ctx, &path_str, format)
296            .await
297            .with_context(|| {
298                format!(
299                    "model '{}': DataFusion listing registration failed for {}",
300                    node.name, path_str
301                )
302            })?;
303        (provider, BronzeRegistrationMode::DataFusionListing)
304    };
305
306    finish_register(
307        ctx,
308        registered,
309        key,
310        schema_name,
311        table_name,
312        node,
313        format,
314        resolved,
315        inner,
316        mode,
317    )
318    .await
319}
320
321async fn finish_register(
322    ctx: &SessionContext,
323    registered: &mut HashSet<(String, String)>,
324    key: (String, String),
325    schema_name: String,
326    table_name: String,
327    node: &ModelNode,
328    format: SourceFormat,
329    resolved: PathBuf,
330    inner: Arc<dyn TableProvider>,
331    mode: BronzeRegistrationMode,
332) -> Result<Option<usize>> {
333
334    let meta = BronzeSourceMeta {
335        model_name: node.name.clone(),
336        source_schema: schema_name.clone(),
337        source_table: table_name.clone(),
338        format,
339        scan_path: resolved,
340        registration_mode: mode,
341    };
342
343    let bronze = Arc::new(BronzeTableProvider::wrap(inner, meta));
344    let table_ref = TableReference::partial(schema_name.clone(), table_name.clone());
345
346    // Replace if present (re-runs / tests)
347    let _ = ctx.deregister_table(table_ref.clone());
348    ctx.register_table(table_ref, bronze)
349        .map_err(|e| anyhow::anyhow!("register {}.{}: {}", schema_name, table_name, e))?;
350
351    registered.insert(key);
352    tracing::info!(
353        "Registered bronze source {}.{} from model '{}' ({:?}, format={})",
354        schema_name,
355        table_name,
356        node.name,
357        mode,
358        format
359    );
360    Ok(Some(1))
361}
362
363fn is_empty_or_missing_err(e: &anyhow::Error) -> bool {
364    let s = format!("{e:#}");
365    s.contains("E_RBT_BRONZE_SCAN_EMPTY")
366        || s.contains("E_RBT_BRONZE_SCAN_PATH_NOT_FOUND")
367        || s.contains("bronze scan produced zero batches")
368}
369
370/// Zero-row MemTable with declared schema; partition keys filled from run scope when present.
371fn empty_memtable(fm: &StagingFrontmatter, scope: &RunScope) -> Result<Arc<dyn TableProvider>> {
372    let schema = fm.empty_frame_schema()?;
373    // Empty batch — partition constants are for SQL via require_partitions on non-empty scans;
374    // empty frame still exposes partition columns as nullable Utf8 for outer joins.
375    let batch = RecordBatch::new_empty(schema.clone());
376    let _ = scope; // reserved: future constant partition columns on empty frames
377    let mem = MemTable::try_new(schema, vec![vec![batch]])
378        .map_err(|e| anyhow::anyhow!("MemTable empty: {e}"))?;
379    Ok(Arc::new(mem))
380}
381
382fn should_use_scan_path(fm: &StagingFrontmatter, format: SourceFormat) -> bool {
383    if fm.force_scan.unwrap_or(false) {
384        return true;
385    }
386    // Hive partition injection / filters / source path / path_glob require the scan path
387    // (DataFusion listing does not inject path-derived columns or apply rbt globs).
388    if fm
389        .partition_by
390        .as_ref()
391        .map(|p| !p.is_empty())
392        .unwrap_or(false)
393        || fm
394            .require_partitions
395            .as_ref()
396            .map(|p| !p.is_empty())
397            .unwrap_or(false)
398        || fm
399            .path_glob
400            .as_ref()
401            .map(|p| !p.is_empty())
402            .unwrap_or(false)
403        || fm.inject_source_path.unwrap_or(false)
404    {
405        return true;
406    }
407    // jshift selective extract
408    if matches!(format, SourceFormat::Jsonl | SourceFormat::Json)
409        && fm.paths.as_ref().map(|p| !p.is_empty()).unwrap_or(false)
410    {
411        return true;
412    }
413    // Nested hive dirs, stream IPC, and opaque protobuf need the scan path.
414    matches!(
415        format,
416        SourceFormat::Log
417            | SourceFormat::Txt
418            | SourceFormat::Toml
419            | SourceFormat::ArrowIpc
420            | SourceFormat::ArrowIpcStream
421            | SourceFormat::Protobuf
422    )
423}
424
425async fn ensure_schema(ctx: &SessionContext, schema_name: &str) -> Result<()> {
426    // DataFusion accepts CREATE SCHEMA via SQL
427    let sql = format!(
428        "CREATE SCHEMA IF NOT EXISTS \"{}\"",
429        schema_name.replace('"', "")
430    );
431    ctx.sql(&sql)
432        .await
433        .with_context(|| format!("CREATE SCHEMA {}", schema_name))?
434        .collect()
435        .await
436        .with_context(|| format!("CREATE SCHEMA {} collect", schema_name))?;
437    Ok(())
438}
439
440/// Path A: materialize a DF listing provider, then return it for wrapping.
441async fn listing_table_provider(
442    ctx: &SessionContext,
443    path: &str,
444    format: SourceFormat,
445) -> Result<Arc<dyn TableProvider>> {
446    // Register under a private temp name, extract provider, deregister.
447    let tmp = format!(
448        "__rbt_bronze_tmp_{}",
449        std::time::SystemTime::now()
450            .duration_since(std::time::UNIX_EPOCH)
451            .map(|d| d.as_nanos())
452            .unwrap_or(0)
453    );
454
455    match format {
456        SourceFormat::Parquet => {
457            ctx.register_parquet(&tmp, path, ParquetReadOptions::default())
458                .await?;
459        }
460        SourceFormat::Csv => {
461            ctx.register_csv(&tmp, path, CsvReadOptions::default())
462                .await?;
463        }
464        SourceFormat::Jsonl => {
465            let opts = JsonReadOptions::default()
466                .file_extension(".jsonl")
467                .newline_delimited(true);
468            // DF register_type_check requires path to end with extension; if directory, ok
469            if let Err(e) = ctx.register_json(&tmp, path, opts).await {
470                // Fallback: .json extension / generic
471                tracing::debug!("jsonl register with .jsonl failed ({e}); retrying default");
472                ctx.register_json(&tmp, path, JsonReadOptions::default())
473                    .await?;
474            }
475        }
476        SourceFormat::Json => {
477            let opts = JsonReadOptions::default().newline_delimited(false);
478            ctx.register_json(&tmp, path, opts).await?;
479        }
480        SourceFormat::ArrowIpc => {
481            ctx.register_arrow(&tmp, path, ArrowReadOptions::default())
482                .await?;
483        }
484        other => bail!("listing_table_provider does not support format {}", other),
485    }
486
487    let provider = ctx
488        .table_provider(TableReference::bare(tmp.as_str()))
489        .await
490        .with_context(|| format!("lookup temp bronze table {}", tmp))?;
491    let _ = ctx.deregister_table(TableReference::bare(tmp.as_str()))?;
492    Ok(provider)
493}
494
495fn should_spill_to_parquet(
496    format: SourceFormat,
497    config: &crate::core::project::RbtProjectConfig,
498) -> bool {
499    config.scan.spill_arrow_ipc
500        && matches!(
501            format,
502            SourceFormat::ArrowIpc | SourceFormat::ArrowIpcStream
503        )
504}
505
506async fn scan_to_memtable(
507    project_dir: &Path,
508    fm: &StagingFrontmatter,
509    format: SourceFormat,
510    config: &crate::core::project::RbtProjectConfig,
511) -> Result<Arc<dyn TableProvider>> {
512    let mut req = ScanRequest::from_frontmatter_with_config(
513        project_dir,
514        fm,
515        config.roots.clone(),
516        &config.scan,
517    )?;
518    req.format = format;
519    let scanner = LakeScanner::from_request(&req);
520    let batches = scanner.scan(&req).await?;
521    if batches.is_empty() {
522        if req.allow_empty {
523            let schema = fm.empty_frame_schema().unwrap_or_else(|_| {
524                Arc::new(arrow::datatypes::Schema::new(vec![arrow::datatypes::Field::new(
525                    "_empty",
526                    arrow::datatypes::DataType::Utf8,
527                    true,
528                )]))
529            });
530            let batch = RecordBatch::new_empty(schema.clone());
531            let mem = MemTable::try_new(schema, vec![vec![batch]])
532                .map_err(|e| anyhow::anyhow!("MemTable::try_new empty: {e}"))?;
533            return Ok(Arc::new(mem));
534        }
535        bail!(
536            "E_RBT_BRONZE_SCAN_EMPTY: bronze scan produced zero batches for {}",
537            req.resolved_path()?.display()
538        );
539    }
540    let schema = batches[0].schema();
541    // MemTable expects Vec<Vec<RecordBatch>> partitions
542    let mem = MemTable::try_new(schema, vec![batches])
543        .map_err(|e| anyhow::anyhow!("MemTable::try_new: {}", e))?;
544    Ok(Arc::new(mem))
545}
546
547/// Stream Arrow IPC (etc.) file-by-file into a project spill Parquet, then DF-list it.
548async fn scan_spill_to_listing(
549    ctx: &SessionContext,
550    project_dir: &Path,
551    fm: &StagingFrontmatter,
552    format: SourceFormat,
553    config: &crate::core::project::RbtProjectConfig,
554    schema_name: &str,
555    table_name: &str,
556) -> Result<Arc<dyn TableProvider>> {
557    let mut req = ScanRequest::from_frontmatter_with_config(
558        project_dir,
559        fm,
560        config.roots.clone(),
561        &config.scan,
562    )?;
563    req.format = format;
564    let scanner = LakeScanner::from_request(&req);
565
566    let spill_root = crate::core::paths::resolve_project_path(
567        project_dir,
568        &config.scan.spill_dir,
569        &config.roots,
570    )
571    .with_context(|| {
572        format!(
573            "E_RBT_BRONZE_SPILL: resolve spill_dir '{}'",
574            config.scan.spill_dir
575        )
576    })?;
577    std::fs::create_dir_all(&spill_root).with_context(|| {
578        format!(
579            "E_RBT_BRONZE_SPILL: mkdir {}",
580            spill_root.display()
581        )
582    })?;
583    let safe = format!(
584        "{}__{}.parquet",
585        schema_name.replace('/', "_"),
586        table_name.replace('/', "_")
587    );
588    let spill_path = spill_root.join(safe);
589
590    let opts = crate::materializer::MaterializeWriteOptions::from_config(&config.materialize, true);
591    let stats = scanner
592        .scan_spill_to_parquet(&req, &spill_path, &opts)
593        .with_context(|| {
594            format!(
595                "E_RBT_BRONZE_SPILL: spill to {}",
596                spill_path.display()
597            )
598        })?;
599    tracing::info!(
600        "Bronze {}.{} spilled {} rows ({} batches) → {}",
601        schema_name,
602        table_name,
603        stats.rows,
604        stats.batches,
605        spill_path.display()
606    );
607
608    listing_table_provider(
609        ctx,
610        spill_path.to_str().unwrap_or_default(),
611        SourceFormat::Parquet,
612    )
613    .await
614}
615
616#[cfg(test)]
617mod tests {
618    use super::*;
619    use crate::core::dag::{Materialization, ModelDag, OutputFormat};
620
621    #[tokio::test]
622    async fn register_arrow_ipc_spills_to_parquet() -> Result<()> {
623        use arrow::array::Int64Array;
624        use arrow::datatypes::{DataType, Field, Schema};
625        use arrow::ipc::writer::FileWriter;
626        use arrow::record_batch::RecordBatch;
627        use std::sync::Arc;
628
629        let temp = tempfile::tempdir()?;
630        let bronze = temp.path().join("lake/bronze/symbol=X/timeframe=1m");
631        std::fs::create_dir_all(&bronze)?;
632        let schema = Arc::new(Schema::new(vec![
633            Field::new("symbol", DataType::Utf8, false),
634            Field::new("v", DataType::Int64, false),
635        ]));
636        let batch = RecordBatch::try_new(
637            schema.clone(),
638            vec![
639                Arc::new(arrow::array::StringArray::from(vec!["X", "X"])),
640                Arc::new(Int64Array::from(vec![1, 2])),
641            ],
642        )?;
643        let f = std::fs::File::create(bronze.join("chunk.arrow"))?;
644        let mut w = FileWriter::try_new(f, &schema)?;
645        w.write(&batch)?;
646        w.finish()?;
647
648        let sql = r#"---
649source_format: arrow_ipc
650scan_path: "lake/bronze"
651path_glob: "**/*.arrow"
652partition_by: [symbol, timeframe]
653require_partitions:
654  timeframe: "1m"
655inject_source_path: true
656---
657SELECT symbol, timeframe, v FROM {{ source('bronze', 'ohlcv') }}
658"#;
659        let mut dag = ModelDag::new();
660        dag.add_model_with_format(
661            "stg_ohlcv",
662            sql,
663            Materialization::Table,
664            OutputFormat::Parquet,
665            None,
666            "",
667        )?;
668        dag.build_graph()?;
669
670        let ctx = SessionContext::new();
671        let mut registered = HashSet::new();
672        let cfg = crate::core::project::RbtProjectConfig::default();
673        assert!(cfg.scan.spill_arrow_ipc);
674        let n = register_bronze_sources_for_dag(&ctx, &dag, temp.path(), &mut registered, &cfg)
675            .await?;
676        assert_eq!(n, 1);
677
678        let spill = temp
679            .path()
680            .join(".rbt/bronze_spill/bronze__ohlcv.parquet");
681        assert!(
682            spill.exists(),
683            "expected spill parquet at {}",
684            spill.display()
685        );
686
687        let df = ctx
688            .sql("SELECT COUNT(*) AS c FROM bronze.ohlcv")
689            .await?;
690        let batches = df.collect().await?;
691        // 2 data rows
692        let c = batches[0]
693            .column(0)
694            .as_any()
695            .downcast_ref::<Int64Array>()
696            .unwrap()
697            .value(0);
698        assert_eq!(c, 2);
699        Ok(())
700    }
701
702    #[tokio::test]
703    async fn register_jsonl_from_frontmatter() -> Result<()> {
704        let temp = tempfile::tempdir()?;
705        let bronze = temp.path().join("raw.jsonl");
706        std::fs::write(
707            &bronze,
708            r#"{"ticker":"NVDA","price":1.5}
709{"ticker":"AAPL","price":2.5}
710"#,
711        )?;
712
713        let sql = format!(
714            r#"---
715source_format: jsonl
716scan_path: "{}"
717---
718SELECT ticker, price FROM {{{{ source('bronze', 'raw_trades') }}}}
719"#,
720            bronze.file_name().unwrap().to_string_lossy()
721        );
722
723        let mut dag = ModelDag::new();
724        dag.add_model_with_format(
725            "stg_trades",
726            &sql,
727            Materialization::Table,
728            OutputFormat::Parquet,
729            None,
730            "",
731        )?;
732        dag.build_graph()?;
733
734        let engine_ctx = SessionContext::new();
735        let mut registered = HashSet::new();
736        let cfg = crate::core::project::RbtProjectConfig::default();
737        let n =
738            register_bronze_sources_for_dag(&engine_ctx, &dag, temp.path(), &mut registered, &cfg)
739                .await?;
740        assert_eq!(n, 1);
741
742        let df = engine_ctx
743            .sql("SELECT COUNT(*) AS c FROM bronze.raw_trades")
744            .await?;
745        let batches = df.collect().await?;
746        assert_eq!(batches[0].num_rows(), 1);
747        Ok(())
748    }
749}