diff --git a/rust/lancedb/src/materialized_view.rs b/rust/lancedb/src/materialized_view.rs index 08d6c921e..95480098d 100644 --- a/rust/lancedb/src/materialized_view.rs +++ b/rust/lancedb/src/materialized_view.rs @@ -95,6 +95,10 @@ pub struct ViewProjection { pub struct MaterializedViewDefinition { /// Name of the source table, in the same database as the view. pub source_table: String, + /// Namespace path holding the source table; empty is the root namespace. + /// A definition written before namespaced sources reads as root. + #[serde(default, skip_serializing_if = "Vec::is_empty")] + pub source_namespace: Vec, /// The projected output columns, in view schema order. pub projections: Vec, /// SQL predicate selecting the source rows the view holds. @@ -166,6 +170,7 @@ pub fn materialized_view_kind( pub(crate) fn plan( source_schema: SchemaRef, source_table: &str, + source_namespace: &[String], projections: &[(String, String)], filter: Option<&str>, limit: Option, @@ -319,6 +324,7 @@ pub(crate) fn plan( let definition = MaterializedViewDefinition { source_table: source_table.to_string(), + source_namespace: source_namespace.to_vec(), projections: projections .into_iter() .map(|(output, expression)| ViewProjection { output, expression }) @@ -602,7 +608,7 @@ pub struct PreparedDeclaration { definition: MaterializedViewDefinition, /// The source's own database: the only place /// [`PreparedDeclaration::create`] will put the view, because refresh - /// resolves the recorded source name through the view's database. + /// resolves the recorded source coordinate through the view's database. database: Arc, } @@ -622,10 +628,21 @@ impl PreparedDeclaration { /// Create the view table and verify it, consuming the declaration. /// - /// The view goes in the source's own database, where refresh resolves the - /// recorded source name. Stable row ids are requested at both levels and - /// verified rather than trusted; nothing is rolled back on failure. + /// The view goes at the root of the source's own database, where refresh + /// resolves the recorded source coordinate. Stable row ids are requested + /// at both levels and verified rather than trusted; nothing is rolled + /// back on failure. pub async fn create(self, name: &str) -> Result { + self.create_in(&[], name).await + } + + /// Create the view in `namespace_path`, empty for the root namespace. + /// Otherwise [`PreparedDeclaration::create`]. + pub async fn create_in( + self, + namespace_path: &[String], + name: &str, + ) -> Result { let empty: Vec> = vec![]; // Minted here, not at preparation: a declaration can be cloned and @@ -640,6 +657,7 @@ impl PreparedDeclaration { let reader: Box = Box::new(arrow_array::RecordBatchIterator::new(empty, schema)); let mut request = CreateTableRequest::new(name.to_string(), Box::new(reader)); + request.namespace_path = namespace_path.to_vec(); let write_params = request .write_options .lance_write_params @@ -680,8 +698,8 @@ impl PreparedDeclaration { /// Validate a view declaration against its live source and hold what its /// creation needs. The declaration is canonicalized through the coordinate a -/// refresh will resolve, so a handle that does not resolve back to itself is -/// rejected, as is a namespaced source. Same creation-time checks as +/// refresh will resolve -- name and namespace both -- so a handle that does +/// not resolve back to itself is rejected. Same creation-time checks as /// [`Connection::create_materialized_view`]. /// /// ```no_run @@ -710,17 +728,9 @@ pub async fn prepare_declaration( message: "materialized views are supported only on local databases".into(), }); }; - // The definition records the source by bare name; any other source - // form would be recorded as a name its refresh cannot resolve. - if !source.namespace().is_empty() { - return Err(Error::NotSupported { - message: format!( - "a namespaced source cannot be recorded in a view definition; \ - '{}' must be a root-namespace table", - source.name() - ), - }); - } + // Refresh resolves the source at exactly this coordinate, so the + // definition records the namespace alongside the name. + let source_namespace = source.namespace().to_vec(); let database = source .database_opt() .ok_or_else(|| Error::InvalidInput { @@ -734,7 +744,7 @@ pub async fn prepare_declaration( let resolved = database .open_table(OpenTableRequest { name: source.name().to_string(), - namespace_path: vec![], + namespace_path: source_namespace.clone(), index_cache_size: None, lance_read_params: None, location: None, @@ -780,6 +790,7 @@ pub async fn prepare_declaration( let (definition, mut fields, lineage) = plan( source_schema.clone(), resolved.name(), + &source_namespace, projections, filter, limit, @@ -839,7 +850,9 @@ fn ensure_local(connection: &Connection) -> Result<()> { pub struct CreateMaterializedViewBuilder { connection: Connection, name: String, + namespace: Vec, source: String, + source_namespace: Vec, projections: Vec<(String, String)>, filter: Option, limit: Option, @@ -850,13 +863,31 @@ impl CreateMaterializedViewBuilder { Self { connection, name, + namespace: Vec::new(), source, + source_namespace: Vec::new(), projections: Vec::new(), filter: None, limit: None, } } + /// The namespace to create the view in. Defaults to the root namespace. + pub fn namespace(mut self, namespace: impl IntoIterator>) -> Self { + self.namespace = namespace.into_iter().map(Into::into).collect(); + self + } + + /// The namespace holding the source table. Defaults to the root + /// namespace, and is recorded in the definition for refresh to resolve. + pub fn source_namespace( + mut self, + namespace: impl IntoIterator>, + ) -> Self { + self.source_namespace = namespace.into_iter().map(Into::into).collect(); + self + } + /// The view's columns, as `(name, SQL expression)` pairs. Not calling /// this selects every source column, expanded at creation time. pub fn select( @@ -887,7 +918,12 @@ impl CreateMaterializedViewBuilder { /// provenance across compaction, and cannot be enabled later. pub async fn execute(self) -> Result { ensure_local(&self.connection)?; - let source = self.connection.open_table(&self.source).execute().await?; + let source = self + .connection + .open_table(&self.source) + .namespace(self.source_namespace.clone()) + .execute() + .await?; let prepared = prepare_declaration( &source, &self.projections, @@ -895,7 +931,7 @@ impl CreateMaterializedViewBuilder { self.limit, ) .await?; - prepared.create(&self.name).await + prepared.create_in(&self.namespace, &self.name).await } } @@ -1152,6 +1188,7 @@ mod tests { view.definition(), &MaterializedViewDefinition { source_table: "people".into(), + source_namespace: Vec::new(), projections: vec![ ViewProjection { output: "name".into(), @@ -2083,33 +2120,75 @@ mod tests { .await .unwrap_err(); assert!(err.to_string().contains("custom_loc"), "{err}"); + } - // A namespaced source cannot be recorded in the definition: the - // bare name refresh resolves would reach a different table or none. - let namespaced = crate::table::NativeTable::create( - "memory://ns_src", - "ns_src", - vec!["ns".to_string()], - Box::new(arrow_array::RecordBatchIterator::new( - vec![], - std::sync::Arc::new(arrow_schema::Schema::new(vec![arrow_schema::Field::new( - "id", - arrow_schema::DataType::Int32, - true, - )])), - )) as Box, - None, - None, - None, - None, - std::collections::HashSet::new(), - ) + /// A view declared over a namespaced source records that namespace, and + /// refresh resolves the source through it -- the coordinate round-trips. + #[tokio::test] + async fn a_namespaced_source_round_trips_through_refresh() { + use lance_namespace::models::CreateNamespaceRequest; + + let tmp = tempfile::tempdir().unwrap(); + let mut properties = std::collections::HashMap::new(); + properties.insert("root".to_string(), tmp.path().to_str().unwrap().to_string()); + let conn = crate::connect_namespace("dir", properties) + .execute() + .await + .unwrap(); + conn.create_namespace(CreateNamespaceRequest { + id: Some(vec!["ns".into()]), + ..Default::default() + }) .await .unwrap(); - let namespaced = Table::new(std::sync::Arc::new(namespaced), conn.database().clone()); - let err = prepare_declaration(&namespaced, &[], None, None) + + let batch = record_batch!( + ("name", Utf8, ["ada", "grace", "alan"]), + ("age", Int32, [36, 85, 41]) + ) + .unwrap(); + conn.create_table("people", batch) + .namespace(vec!["ns".to_string()]) + .write_options(stable_row_ids()) + .execute() .await - .unwrap_err(); - assert!(err.to_string().contains("namespaced source"), "{err}"); + .unwrap(); + + // A decoy of the same name at the root: resolving the source at the + // wrong namespace materializes one row here instead of three. + let decoy = record_batch!(("name", Utf8, ["mallory"]), ("age", Int32, [42])).unwrap(); + conn.create_table("people", decoy) + .write_options(stable_row_ids()) + .execute() + .await + .unwrap(); + + let view = conn + .create_materialized_view("adults", "people") + .namespace(vec!["ns".to_string()]) + .source_namespace(vec!["ns".to_string()]) + .select([("name", "name")]) + .only_if("age >= 18") + .execute() + .await + .unwrap(); + + assert_eq!(view.definition().source_table, "people"); + assert_eq!(view.definition().source_namespace, vec!["ns".to_string()]); + assert_eq!(view.table().namespace(), &["ns"]); + + // Refresh resolves the source at the recorded namespace, not at root. + let result = view.refresh().execute().await.unwrap(); + assert_eq!(result.rows_written, 3); + } + + /// A definition stored before namespaced sources existed carries no + /// namespace key and must read as the root namespace. + #[test] + fn a_definition_without_a_namespace_reads_as_root() { + let stored = + r#"{"source_table":"people","projections":[{"output":"name","expression":"name"}]}"#; + let definition: MaterializedViewDefinition = serde_json::from_str(stored).unwrap(); + assert!(definition.source_namespace.is_empty()); } } diff --git a/rust/lancedb/src/materialized_view/refresh.rs b/rust/lancedb/src/materialized_view/refresh.rs index b967e81f8..efddd1c7b 100644 --- a/rust/lancedb/src/materialized_view/refresh.rs +++ b/rust/lancedb/src/materialized_view/refresh.rs @@ -170,6 +170,7 @@ pub(crate) async fn execute_refresh( let (replanned, mut planned_fields, _renames) = super::plan( source_schema, &definition.source_table, + &definition.source_namespace, &projections, definition.filter.as_deref(), definition.limit, @@ -590,7 +591,7 @@ async fn open_source(view: &Table, definition: &MaterializedViewDefinition) -> R let source = database .open_table(OpenTableRequest { name: definition.source_table.clone(), - namespace_path: Vec::new(), + namespace_path: definition.source_namespace.clone(), index_cache_size: None, lance_read_params: None, location: None, @@ -2919,6 +2920,7 @@ mod tests { let replacement = crate::materialized_view::MaterializedViewDefinition { source_table: "src".into(), + source_namespace: Vec::new(), projections: vec![ crate::materialized_view::ViewProjection { output: "x".into(), @@ -2958,6 +2960,7 @@ mod tests { let narrower = crate::materialized_view::MaterializedViewDefinition { source_table: "src".into(), + source_namespace: Vec::new(), projections: vec![crate::materialized_view::ViewProjection { output: "x".into(), expression: "x".into(),