1use 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#[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 DataFusionListing,
57 ScanMemTable,
59 ScanSpillParquet,
61}
62
63#[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
112pub 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
127pub 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
149pub 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
160pub 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 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
370fn empty_memtable(fm: &StagingFrontmatter, scope: &RunScope) -> Result<Arc<dyn TableProvider>> {
372 let schema = fm.empty_frame_schema()?;
373 let batch = RecordBatch::new_empty(schema.clone());
376 let _ = scope; 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 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 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 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 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
440async fn listing_table_provider(
442 ctx: &SessionContext,
443 path: &str,
444 format: SourceFormat,
445) -> Result<Arc<dyn TableProvider>> {
446 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 if let Err(e) = ctx.register_json(&tmp, path, opts).await {
470 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 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
547async 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 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}