mirror of
https://github.com/lancedb/lancedb.git
synced 2026-09-30 00:45:37 +00:00
feat: retire a Function binding when its output columns are dropped (#4231)
Dropping a column a Function binding writes was refused outright, which left a bad declaration unrecoverable: the binding is immutable, there is no rebind, and so the column could never be filled again. A drop that names every output of a binding now retires it instead. A binding lives in three places -- the envelope in schema metadata, a `computed_column.*` marker on each output field, and the columns themselves -- and readers cross-check the first two, so removing one and leaving the others is a table that refuses every write. The retirement therefore commits in two steps, each landing a state that stands on its own: one `UpdateConfig` rewrites the envelope and clears the markers together, leaving the outputs as ordinary columns holding their last values, and the drop follows. Interrupted between them, the columns are still there to be dropped again. Folding the metadata into the drop's own `Project` was the obvious alternative and does not work: the transaction proto records `Project` as fields alone, so the edit would survive only in the writer's memory. Naming one output of a multi-output binding is refused and names the missing siblings, since one refresh writes them in one commit. A multi-output binding's hidden `__function_assignment_*` column goes with it. Inputs stay protected while a surviving binding reads them, and drop in the request that retires the last one.
This commit is contained in:
@@ -418,23 +418,18 @@ pub(crate) fn ensure_not_function_bound<S: AsRef<str>>(
|
||||
touched: impl IntoIterator<Item = S>,
|
||||
) -> Result<()> {
|
||||
ensure_supported_function_metadata(schema)?;
|
||||
let mut protected = BTreeSet::new();
|
||||
for binding in function_bindings(schema)? {
|
||||
for input in binding.inputs() {
|
||||
protected.insert(field_root(&input.field_path)?);
|
||||
}
|
||||
protected.extend(
|
||||
binding
|
||||
.outputs()
|
||||
.iter()
|
||||
.map(|output| output.output_name.clone()),
|
||||
);
|
||||
protected.extend(
|
||||
binding
|
||||
.assignment()
|
||||
.map(|assignment| assignment.output_name.clone()),
|
||||
);
|
||||
}
|
||||
ensure_not_bound_by(&function_bindings(schema)?, operation, touched)
|
||||
}
|
||||
|
||||
/// [`ensure_not_function_bound`] against an explicit set of bindings, for a
|
||||
/// caller that is retiring some of the schema's own in the same operation and
|
||||
/// must be checked against what survives it.
|
||||
pub(crate) fn ensure_not_bound_by<S: AsRef<str>>(
|
||||
bindings: &[FunctionBinding],
|
||||
operation: &str,
|
||||
touched: impl IntoIterator<Item = S>,
|
||||
) -> Result<()> {
|
||||
let protected = protected_roots(bindings)?;
|
||||
if protected.is_empty() {
|
||||
return Ok(());
|
||||
}
|
||||
@@ -452,6 +447,129 @@ pub(crate) fn ensure_not_function_bound<S: AsRef<str>>(
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Every column `bindings` depend on, inputs taken at their root.
|
||||
fn protected_roots(bindings: &[FunctionBinding]) -> Result<BTreeSet<String>> {
|
||||
let mut protected = BTreeSet::new();
|
||||
for binding in bindings {
|
||||
for input in binding.inputs() {
|
||||
protected.insert(field_root(&input.field_path)?);
|
||||
}
|
||||
protected.extend(
|
||||
binding
|
||||
.outputs()
|
||||
.iter()
|
||||
.map(|output| output.output_name.clone()),
|
||||
);
|
||||
protected.extend(
|
||||
binding
|
||||
.assignment()
|
||||
.map(|assignment| assignment.output_name.clone()),
|
||||
);
|
||||
}
|
||||
Ok(protected)
|
||||
}
|
||||
|
||||
/// What a drop must do besides removing columns, so that the Function
|
||||
/// bindings it covers are retired rather than stranded.
|
||||
#[derive(Debug, Default)]
|
||||
pub struct FunctionUnbinding {
|
||||
/// The bindings the drop leaves in place. Later columns in the same
|
||||
/// operation are checked against these, not against the schema's.
|
||||
pub retained: Vec<FunctionBinding>,
|
||||
/// Assignment columns of the retired bindings that the caller did not
|
||||
/// name. Internal bookkeeping columns, so the drop adds them itself.
|
||||
pub assignment_columns: Vec<String>,
|
||||
/// Columns whose `computed_column.*` declaration metadata the unbind
|
||||
/// commit clears.
|
||||
pub cleared_columns: Vec<String>,
|
||||
/// The new value of [`FUNCTION_BINDINGS_META_KEY`]; `None` deletes it.
|
||||
pub bindings_metadata: Option<String>,
|
||||
}
|
||||
|
||||
impl FunctionUnbinding {
|
||||
/// True when the drop retires nothing and is an ordinary column drop.
|
||||
pub fn is_noop(&self) -> bool {
|
||||
self.cleared_columns.is_empty()
|
||||
}
|
||||
|
||||
/// Refuse `operation` on a path that a binding surviving this drop still
|
||||
/// reads or writes. The retired ones no longer protect anything.
|
||||
pub fn ensure_retained_unaffected<S: AsRef<str>>(
|
||||
&self,
|
||||
operation: &str,
|
||||
touched: impl IntoIterator<Item = S>,
|
||||
) -> Result<()> {
|
||||
ensure_not_bound_by(&self.retained, operation, touched)
|
||||
}
|
||||
}
|
||||
|
||||
/// Plan the retirement of every Function binding whose outputs `columns`
|
||||
/// covers.
|
||||
///
|
||||
/// Naming one output of a binding means naming all of them: the outputs of
|
||||
/// one Function are written by one refresh in one commit, so a binding that
|
||||
/// kept some of them would have no coherent shape to write. A partial drop is
|
||||
/// refused and names the siblings that are missing.
|
||||
pub fn plan_function_unbinding(
|
||||
schema: &ArrowSchema,
|
||||
columns: &[&str],
|
||||
) -> Result<FunctionUnbinding> {
|
||||
ensure_supported_function_metadata(schema)?;
|
||||
let named = columns
|
||||
.iter()
|
||||
.map(|column| whole_column(column))
|
||||
.collect::<Result<Vec<Option<String>>>>()?
|
||||
.into_iter()
|
||||
.flatten()
|
||||
.collect::<BTreeSet<String>>();
|
||||
|
||||
let mut plan = FunctionUnbinding::default();
|
||||
for binding in function_bindings(schema)? {
|
||||
let outputs = binding
|
||||
.outputs()
|
||||
.iter()
|
||||
.map(|output| output.output_name.clone())
|
||||
.collect::<Vec<_>>();
|
||||
let Some(dropped) = outputs.iter().find(|name| named.contains(name.as_str())) else {
|
||||
plan.retained.push(binding);
|
||||
continue;
|
||||
};
|
||||
let missing = outputs
|
||||
.iter()
|
||||
.filter(|name| !named.contains(name.as_str()))
|
||||
.cloned()
|
||||
.collect::<Vec<_>>();
|
||||
if !missing.is_empty() {
|
||||
return Err(Error::InvalidInput {
|
||||
message: format!(
|
||||
"dropping Function output '{dropped}' must drop every output of its \
|
||||
binding; the same request must also name {}",
|
||||
missing.join(", ")
|
||||
),
|
||||
});
|
||||
}
|
||||
plan.cleared_columns.extend(outputs);
|
||||
if let Some(assignment) = binding.assignment() {
|
||||
plan.cleared_columns.push(assignment.output_name.clone());
|
||||
if !named.contains(assignment.output_name.as_str()) {
|
||||
plan.assignment_columns.push(assignment.output_name.clone());
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Re-encoding an untouched set would rewrite bytes a stricter reader
|
||||
// compares against its own encoding, so leave the key alone when the drop
|
||||
// retires nothing.
|
||||
if !plan.is_noop() {
|
||||
plan.bindings_metadata = if plan.retained.is_empty() {
|
||||
None
|
||||
} else {
|
||||
Some(function_bindings_metadata(&plan.retained)?)
|
||||
};
|
||||
}
|
||||
Ok(plan)
|
||||
}
|
||||
|
||||
/// Read a field's computed-column declaration, if it carries one.
|
||||
///
|
||||
/// A field flagged computed but carrying no kind, or a SQL one missing its
|
||||
@@ -1320,7 +1438,30 @@ pub(crate) fn plan_function_application(
|
||||
/// Paths are compared at their root: a declaration reading `metadata` is
|
||||
/// invalidated by a change to `metadata.age` just as surely.
|
||||
pub(crate) fn ensure_not_an_input(schema: &SchemaRef, paths: &[&str]) -> Result<()> {
|
||||
ensure_not_an_input_of(schema, paths, &[])
|
||||
}
|
||||
|
||||
/// [`ensure_not_an_input`] where `removed` names columns whose own
|
||||
/// declarations the same operation deletes. A reader that is going away
|
||||
/// cannot be stranded by dropping what it reads, so a group drops together.
|
||||
pub(crate) fn ensure_not_an_input_of(
|
||||
schema: &SchemaRef,
|
||||
paths: &[&str],
|
||||
removed: &[&str],
|
||||
) -> Result<()> {
|
||||
// By identity, not spelling: a quoted path names the same column, while
|
||||
// a path *into* one removes no declaration at all.
|
||||
let removed = removed
|
||||
.iter()
|
||||
.map(|path| whole_column(path))
|
||||
.collect::<Result<Vec<_>>>()?
|
||||
.into_iter()
|
||||
.flatten()
|
||||
.collect::<BTreeSet<String>>();
|
||||
for declaration in computed_columns(schema) {
|
||||
if removed.contains(&declaration.name) {
|
||||
continue;
|
||||
}
|
||||
// The expression, not stored inputs, is the source of truth; an
|
||||
// expression that no longer parses proves nothing, so refuse.
|
||||
let inputs = match &declaration.kind {
|
||||
@@ -1583,6 +1724,21 @@ pub(crate) fn field_root(path: &str) -> Result<String> {
|
||||
})
|
||||
}
|
||||
|
||||
/// The column a path names when it names the whole column, and `None` when it
|
||||
/// addresses a field inside one. Retiring a binding takes the whole output;
|
||||
/// a path into one is left to the ordinary guard to refuse.
|
||||
fn whole_column(path: &str) -> Result<Option<String>> {
|
||||
let mut segments = parse_field_path(path)
|
||||
.map_err(|e| Error::InvalidInput {
|
||||
message: format!("invalid column path '{path}': {e}"),
|
||||
})?
|
||||
.into_iter();
|
||||
let column = segments.next().ok_or_else(|| Error::InvalidInput {
|
||||
message: format!("column path '{path}' is empty"),
|
||||
})?;
|
||||
Ok(segments.next().is_none().then_some(column))
|
||||
}
|
||||
|
||||
/// [`field_root`], falling back to the text before the first dot for a
|
||||
/// spelling lance would not resolve.
|
||||
pub(crate) fn root(path: &str) -> String {
|
||||
|
||||
@@ -9,7 +9,8 @@
|
||||
//! - [`drop_columns`](execute_drop_columns): Remove columns from the table
|
||||
|
||||
use arrow_schema::Schema as ArrowSchema;
|
||||
use lance::dataset::{ColumnAlteration, NewColumnTransform};
|
||||
use lance::dataset::transaction::{Operation, Transaction, UpdateMap, UpdateMapEntry};
|
||||
use lance::dataset::{ColumnAlteration, CommitBuilder, Dataset, NewColumnTransform};
|
||||
use serde::{Deserialize, Serialize};
|
||||
use std::collections::HashMap;
|
||||
|
||||
@@ -240,14 +241,96 @@ pub(crate) async fn execute_drop_columns(
|
||||
table.dataset.ensure_mutable()?;
|
||||
let mut dataset = (*table.dataset.get().await?).clone();
|
||||
let schema = std::sync::Arc::new(ArrowSchema::from(dataset.schema()));
|
||||
computed_columns::ensure_not_function_bound(schema.as_ref(), "schema evolution", columns)?;
|
||||
computed_columns::ensure_not_an_input(&schema, columns)?;
|
||||
dataset.drop_columns(columns).await?;
|
||||
let unbinding = computed_columns::plan_function_unbinding(schema.as_ref(), columns)?;
|
||||
|
||||
let mut names = columns.iter().map(|c| c.to_string()).collect::<Vec<_>>();
|
||||
names.extend(unbinding.assignment_columns.iter().cloned());
|
||||
let dropped = names.iter().map(String::as_str).collect::<Vec<_>>();
|
||||
computed_columns::ensure_not_bound_by(&unbinding.retained, "schema evolution", &dropped)?;
|
||||
computed_columns::ensure_not_an_input_of(&schema, &dropped, &dropped)?;
|
||||
|
||||
if !unbinding.is_noop() {
|
||||
// The retirement commits before the projection, so the whole
|
||||
// projection has to be known valid against this revision first --
|
||||
// otherwise the binding is gone and the columns are not. Lance plans
|
||||
// it rather than this repeating the rules: dropping a struct's last
|
||||
// child removes the struct, and data files that disagree on metadata
|
||||
// semantics are refused, neither of which an approximation here would
|
||||
// get right.
|
||||
dataset.plan_drop_columns(&dropped)?;
|
||||
dataset = commit_function_unbinding(dataset, &unbinding).await?;
|
||||
}
|
||||
dataset.drop_columns(&dropped).await?;
|
||||
let version = dataset.version().version;
|
||||
table.dataset.update(dataset);
|
||||
Ok(DropColumnsResult { version })
|
||||
}
|
||||
|
||||
/// Retire the bindings a drop covers, in the commit before it.
|
||||
///
|
||||
/// A binding lives in three places -- the schema-metadata envelope, each
|
||||
/// output's `computed_column.*` field metadata, and the columns themselves --
|
||||
/// and readers cross-check the first two, so a state with one of them removed
|
||||
/// is a table that refuses every write. Clearing both metadata halves in one
|
||||
/// `UpdateConfig` leaves the outputs as ordinary columns holding their last
|
||||
/// values: valid on its own, and the drop that follows is then an ordinary
|
||||
/// drop. Interrupted in between, the columns are still there to be dropped
|
||||
/// again.
|
||||
///
|
||||
/// Not folded into the drop's own `Operation::Project`, though it carries a
|
||||
/// schema: the transaction proto records `Project` as fields alone, so the
|
||||
/// metadata edit would survive only in this writer's memory.
|
||||
async fn commit_function_unbinding(
|
||||
dataset: Dataset,
|
||||
unbinding: &computed_columns::FunctionUnbinding,
|
||||
) -> Result<Dataset> {
|
||||
let bindings = UpdateMap {
|
||||
update_entries: vec![UpdateMapEntry {
|
||||
key: computed_columns::FUNCTION_BINDINGS_META_KEY.to_string(),
|
||||
value: unbinding.bindings_metadata.clone(),
|
||||
}],
|
||||
replace: false,
|
||||
};
|
||||
let mut field_metadata_updates = HashMap::new();
|
||||
for column in &unbinding.cleared_columns {
|
||||
let field = dataset
|
||||
.schema()
|
||||
.field(column)
|
||||
.ok_or_else(|| Error::InvalidInput {
|
||||
message: format!("Function output '{column}' does not exist in the table"),
|
||||
})?;
|
||||
let cleared = field
|
||||
.metadata
|
||||
.keys()
|
||||
.filter(|key| computed_columns::is_declaration_key(key))
|
||||
.map(|key| UpdateMapEntry {
|
||||
key: key.clone(),
|
||||
value: None,
|
||||
})
|
||||
.collect::<Vec<_>>();
|
||||
field_metadata_updates.insert(
|
||||
field.id,
|
||||
UpdateMap {
|
||||
update_entries: cleared,
|
||||
replace: false,
|
||||
},
|
||||
);
|
||||
}
|
||||
let transaction = Transaction::new(
|
||||
dataset.manifest.version,
|
||||
Operation::UpdateConfig {
|
||||
config_updates: None,
|
||||
table_metadata_updates: None,
|
||||
schema_metadata_updates: Some(bindings),
|
||||
field_metadata_updates,
|
||||
},
|
||||
None,
|
||||
);
|
||||
Ok(CommitBuilder::new(std::sync::Arc::new(dataset))
|
||||
.execute(transaction)
|
||||
.await?)
|
||||
}
|
||||
|
||||
/// Internal implementation of the update field metadata logic.
|
||||
///
|
||||
/// Merges or replaces per-field metadata, addressing fields by dot-path.
|
||||
@@ -330,8 +413,9 @@ mod tests {
|
||||
use crate::query::{ExecutableQuery, QueryBase, Select};
|
||||
use crate::table::NewColumnTransform;
|
||||
use crate::table::computed_columns::{
|
||||
FUNCTION_BINDINGS_META_KEY, ensure_supported_function_metadata, function_bindings,
|
||||
function_bindings_metadata, function_computed_column_metadata,
|
||||
FUNCTION_ASSIGNMENT_OUTPUT_ORDINAL, FUNCTION_BINDINGS_META_KEY,
|
||||
ensure_supported_function_metadata, function_bindings, function_bindings_metadata,
|
||||
function_computed_column_metadata, is_declaration_key,
|
||||
};
|
||||
use crate::{Error, Table};
|
||||
use std::collections::HashMap;
|
||||
@@ -351,38 +435,71 @@ mod tests {
|
||||
)
|
||||
.unwrap();
|
||||
let table = conn.create_table("bound", batch).execute().await.unwrap();
|
||||
let binding = FunctionBinding::from_json(include_str!(
|
||||
stamp_bindings(&table, &[text_features_binding("")]).await;
|
||||
table
|
||||
}
|
||||
|
||||
/// The fixture binding, its id and output names suffixed so a table can
|
||||
/// carry more than one.
|
||||
fn text_features_binding(suffix: &str) -> FunctionBinding {
|
||||
let raw = include_str!(
|
||||
"../../tests/fixtures/first_class_functions/v1/remote_function_binding.json"
|
||||
))
|
||||
.unwrap();
|
||||
);
|
||||
if suffix.is_empty() {
|
||||
return FunctionBinding::from_json(raw).unwrap();
|
||||
}
|
||||
FunctionBinding::from_json(
|
||||
&raw.replace("fb_01K3TEXT", &format!("fb_01K3TEXT{suffix}"))
|
||||
.replace("search_text", &format!("search_text{suffix}"))
|
||||
.replace("search_token_count", &format!("search_token_count{suffix}")),
|
||||
)
|
||||
.unwrap()
|
||||
}
|
||||
|
||||
/// Stamp bindings the way the server does it, since no local path
|
||||
/// declares one: the envelope on the schema, a declaration marker on
|
||||
/// every output field.
|
||||
async fn stamp_bindings(table: &Table, bindings: &[FunctionBinding]) {
|
||||
let native = table.as_native().unwrap();
|
||||
let mut dataset = native.dataset.get().await.unwrap().as_ref().clone();
|
||||
dataset
|
||||
.update_schema_metadata(vec![(
|
||||
FUNCTION_BINDINGS_META_KEY.to_string(),
|
||||
Some(function_bindings_metadata(std::slice::from_ref(&binding)).unwrap()),
|
||||
Some(function_bindings_metadata(bindings).unwrap()),
|
||||
)])
|
||||
.await
|
||||
.unwrap();
|
||||
let inputs = ["title".to_string(), "body".to_string()];
|
||||
let outputs = binding
|
||||
.outputs()
|
||||
let outputs = bindings
|
||||
.iter()
|
||||
.map(|output| {
|
||||
(
|
||||
dataset.schema().field(&output.output_name).unwrap().id as u32,
|
||||
function_computed_column_metadata(
|
||||
binding.binding_id(),
|
||||
output.output_ordinal,
|
||||
&inputs,
|
||||
),
|
||||
)
|
||||
.flat_map(|binding| {
|
||||
let assignment = binding.assignment().map(|assignment| {
|
||||
(
|
||||
assignment.output_name.clone(),
|
||||
FUNCTION_ASSIGNMENT_OUTPUT_ORDINAL,
|
||||
)
|
||||
});
|
||||
binding
|
||||
.outputs()
|
||||
.iter()
|
||||
.map(|output| (output.output_name.clone(), output.output_ordinal))
|
||||
.chain(assignment)
|
||||
.map(|(name, ordinal)| {
|
||||
(
|
||||
dataset.schema().field(&name).unwrap().id as u32,
|
||||
function_computed_column_metadata(
|
||||
binding.binding_id(),
|
||||
ordinal,
|
||||
&inputs,
|
||||
),
|
||||
)
|
||||
})
|
||||
.collect::<Vec<_>>()
|
||||
})
|
||||
.collect::<Vec<_>>();
|
||||
dataset.replace_field_metadata(outputs).await.unwrap();
|
||||
native.dataset.update(dataset);
|
||||
ensure_supported_function_metadata(&table.schema().await.unwrap()).unwrap();
|
||||
table
|
||||
}
|
||||
|
||||
fn metadata_update(path: &str) -> FieldMetadataUpdate {
|
||||
@@ -438,6 +555,18 @@ mod tests {
|
||||
);
|
||||
}
|
||||
|
||||
fn incomplete_group(err: Error) {
|
||||
assert!(
|
||||
matches!(&err, Error::InvalidInput { message }
|
||||
if message.contains("must drop every output of its binding")),
|
||||
"{err:?}"
|
||||
);
|
||||
}
|
||||
|
||||
async fn bindings_of(table: &Table) -> Vec<FunctionBinding> {
|
||||
function_bindings(table.schema().await.unwrap().as_ref()).unwrap()
|
||||
}
|
||||
|
||||
/// Every schema-evolution door refuses a column a binding reads or
|
||||
/// writes, including a rename onto one.
|
||||
#[tokio::test]
|
||||
@@ -445,7 +574,14 @@ mod tests {
|
||||
let table = bound_table().await;
|
||||
let version = table.version().await.unwrap();
|
||||
for column in ["title", "body", "search_text", "search_token_count"] {
|
||||
bound(table.drop_columns(&[column]).await.unwrap_err());
|
||||
// A lone drop of an output is refused too, but for its siblings
|
||||
// rather than for the binding; see the retirement tests below.
|
||||
let err = table.drop_columns(&[column]).await.unwrap_err();
|
||||
if column.starts_with("search_") {
|
||||
incomplete_group(err);
|
||||
} else {
|
||||
bound(err);
|
||||
}
|
||||
bound(
|
||||
table
|
||||
.alter_columns(&[ColumnAlteration::new(column.into()).rename("moved".into())])
|
||||
@@ -500,7 +636,12 @@ mod tests {
|
||||
let table = bound_table().await;
|
||||
let version = table.version().await.unwrap();
|
||||
for path in ["`title`", "`search_text`", "`title`.nested"] {
|
||||
bound(table.drop_columns(&[path]).await.unwrap_err());
|
||||
let err = table.drop_columns(&[path]).await.unwrap_err();
|
||||
if path == "`search_text`" {
|
||||
incomplete_group(err);
|
||||
} else {
|
||||
bound(err);
|
||||
}
|
||||
bound(
|
||||
table
|
||||
.alter_columns(&[ColumnAlteration::new(path.into()).set_nullable(false)])
|
||||
@@ -1038,6 +1179,264 @@ mod tests {
|
||||
assert_eq!(drop_result.version, v4);
|
||||
}
|
||||
|
||||
// Retiring a Function binding by dropping its outputs (ENT-2669).
|
||||
|
||||
/// The whole output set goes, and the binding with it: no envelope, no
|
||||
/// declaration markers, and the inputs it held are ordinary again.
|
||||
#[tokio::test]
|
||||
async fn dropping_every_output_retires_the_binding() {
|
||||
let table = bound_table().await;
|
||||
table
|
||||
.drop_columns(&["search_text", "search_token_count"])
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
let schema = table.schema().await.unwrap();
|
||||
assert!(schema.field_with_name("search_text").is_err());
|
||||
assert!(schema.field_with_name("search_token_count").is_err());
|
||||
assert!(!schema.metadata().contains_key(FUNCTION_BINDINGS_META_KEY));
|
||||
assert!(bindings_of(&table).await.is_empty());
|
||||
|
||||
// What the binding used to protect is now ordinary schema.
|
||||
table.drop_columns(&["title"]).await.unwrap();
|
||||
table
|
||||
.update_field_metadata(&[metadata_update("body")])
|
||||
.await
|
||||
.unwrap();
|
||||
ensure_supported_function_metadata(&table.schema().await.unwrap()).unwrap();
|
||||
}
|
||||
|
||||
/// One output of a multi-output binding cannot go alone: its siblings are
|
||||
/// written by the same refresh in the same commit.
|
||||
#[tokio::test]
|
||||
async fn dropping_part_of_an_output_group_names_the_missing_siblings() {
|
||||
let table = bound_table().await;
|
||||
let version = table.version().await.unwrap();
|
||||
for (named, missing) in [
|
||||
("search_text", "search_token_count"),
|
||||
("search_token_count", "search_text"),
|
||||
] {
|
||||
let err = table.drop_columns(&[named]).await.unwrap_err();
|
||||
let Error::InvalidInput { message } = &err else {
|
||||
panic!("{err:?}");
|
||||
};
|
||||
assert!(
|
||||
message.contains("must drop every output of its binding")
|
||||
&& message.contains(missing),
|
||||
"{message}"
|
||||
);
|
||||
}
|
||||
assert_eq!(table.version().await.unwrap(), version);
|
||||
assert_eq!(bindings_of(&table).await.len(), 1);
|
||||
}
|
||||
|
||||
/// An input is protected by the binding, not by its own name, so it drops
|
||||
/// in the request that retires the binding reading it.
|
||||
#[tokio::test]
|
||||
async fn an_input_drops_with_the_binding_that_reads_it() {
|
||||
let table = bound_table().await;
|
||||
bound(table.drop_columns(&["title"]).await.unwrap_err());
|
||||
|
||||
table
|
||||
.drop_columns(&["title", "body", "search_text", "search_token_count"])
|
||||
.await
|
||||
.unwrap();
|
||||
let schema = table.schema().await.unwrap();
|
||||
assert_eq!(schema.fields().len(), 1);
|
||||
assert!(schema.field_with_name("spare").is_ok());
|
||||
assert!(bindings_of(&table).await.is_empty());
|
||||
}
|
||||
|
||||
/// Lance resolves a quoted spelling to the same field, so it retires the
|
||||
/// same binding rather than slipping past the group rule.
|
||||
#[tokio::test]
|
||||
async fn a_quoted_output_spelling_retires_the_same_binding() {
|
||||
let table = bound_table().await;
|
||||
incomplete_group(table.drop_columns(&["`search_text`"]).await.unwrap_err());
|
||||
table
|
||||
.drop_columns(&["`search_text`", "`search_token_count`"])
|
||||
.await
|
||||
.unwrap();
|
||||
assert!(bindings_of(&table).await.is_empty());
|
||||
}
|
||||
|
||||
/// The unbind lands before the drop, and on its own: at that version the
|
||||
/// outputs are still there as plain columns, with no binding and no
|
||||
/// declaration metadata left pointing at one. A crash in between leaves
|
||||
/// that state, which every later write must accept.
|
||||
#[tokio::test]
|
||||
async fn the_unbind_commit_leaves_a_valid_table_before_the_drop() {
|
||||
let table = bound_table().await;
|
||||
let before = table.version().await.unwrap();
|
||||
table
|
||||
.drop_columns(&["search_text", "search_token_count"])
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(table.version().await.unwrap(), before + 2);
|
||||
|
||||
table.checkout(before + 1).await.unwrap();
|
||||
let schema = table.schema().await.unwrap();
|
||||
assert!(!schema.metadata().contains_key(FUNCTION_BINDINGS_META_KEY));
|
||||
for column in ["search_text", "search_token_count"] {
|
||||
let field = schema.field_with_name(column).unwrap();
|
||||
assert!(
|
||||
!field.metadata().keys().any(|key| is_declaration_key(key)),
|
||||
"{column}: {:?}",
|
||||
field.metadata()
|
||||
);
|
||||
}
|
||||
ensure_supported_function_metadata(&schema).unwrap();
|
||||
assert!(function_bindings(&schema).unwrap().is_empty());
|
||||
}
|
||||
|
||||
/// A table carrying two bindings retires only the one whose outputs the
|
||||
/// drop names; the other keeps its envelope entry and its protection.
|
||||
#[tokio::test]
|
||||
async fn retiring_one_binding_leaves_the_others_standing() {
|
||||
let conn = connect("memory://").execute().await.unwrap();
|
||||
let batch = record_batch!(
|
||||
("title", Utf8, ["a"]),
|
||||
("body", Utf8, ["b"]),
|
||||
("search_text", Utf8, ["a b"]),
|
||||
("search_token_count", Int64, [2]),
|
||||
("search_text_2", Utf8, ["a b"]),
|
||||
("search_token_count_2", Int64, [2])
|
||||
)
|
||||
.unwrap();
|
||||
let table = conn.create_table("two", batch).execute().await.unwrap();
|
||||
stamp_bindings(
|
||||
&table,
|
||||
&[text_features_binding(""), text_features_binding("_2")],
|
||||
)
|
||||
.await;
|
||||
|
||||
table
|
||||
.drop_columns(&["search_text", "search_token_count"])
|
||||
.await
|
||||
.unwrap();
|
||||
let remaining = bindings_of(&table).await;
|
||||
assert_eq!(remaining.len(), 1);
|
||||
assert_eq!(remaining[0].binding_id(), "fb_01K3TEXT_2");
|
||||
// The survivor still protects the inputs it shared with the retired one.
|
||||
bound(table.drop_columns(&["title"]).await.unwrap_err());
|
||||
// And its own outputs are still a group.
|
||||
incomplete_group(table.drop_columns(&["search_text_2"]).await.unwrap_err());
|
||||
}
|
||||
|
||||
/// A multi-output binding also owns a hidden `__function_assignment_*`
|
||||
/// column, which the caller never names. Retiring the binding has to take
|
||||
/// it: leaving it behind orphans a column nothing can fill or explain.
|
||||
#[tokio::test]
|
||||
async fn retiring_a_binding_drops_its_assignment_column() {
|
||||
const ASSIGNMENT: &str = "__function_assignment_fb_01K3TEXT";
|
||||
|
||||
let conn = connect("memory://").execute().await.unwrap();
|
||||
let batch = record_batch!(
|
||||
("title", Utf8, ["a"]),
|
||||
("body", Utf8, ["b"]),
|
||||
("search_text", Utf8, ["a b"]),
|
||||
("search_token_count", Int64, [2]),
|
||||
(ASSIGNMENT, Boolean, [Some(true)]),
|
||||
("spare", Int32, [1])
|
||||
)
|
||||
.unwrap();
|
||||
let table = conn
|
||||
.create_table("assigned", batch)
|
||||
.execute()
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
let raw = include_str!(
|
||||
"../../tests/fixtures/first_class_functions/v1/remote_function_binding.json"
|
||||
)
|
||||
.replace(
|
||||
r#""future_binding""#,
|
||||
&format!(
|
||||
r#""assignment": {{"output_name": "{ASSIGNMENT}", "output_field_id": -1}}, "future_binding""#
|
||||
),
|
||||
)
|
||||
// The assignment column is a physical sibling, so it belongs to the
|
||||
// binding's output schema too -- the server writes it that way.
|
||||
.replace(
|
||||
r#"{"name": "search_token_count", "nullable": true, "type": {"type": "int64"}}"#,
|
||||
&format!(
|
||||
r#"{{"name": "search_token_count", "nullable": true, "type": {{"type": "int64"}}}}, {{"name": "{ASSIGNMENT}", "nullable": true, "type": {{"type": "bool"}}}}"#
|
||||
),
|
||||
);
|
||||
let binding = FunctionBinding::from_json(&raw).unwrap();
|
||||
assert!(binding.assignment().is_some(), "fixture must carry one");
|
||||
stamp_bindings(&table, std::slice::from_ref(&binding)).await;
|
||||
|
||||
// The assignment column is the binding's, so it is protected too...
|
||||
bound(table.drop_columns(&[ASSIGNMENT]).await.unwrap_err());
|
||||
|
||||
// ...and goes with the outputs without the caller naming it.
|
||||
table
|
||||
.drop_columns(&["search_text", "search_token_count"])
|
||||
.await
|
||||
.unwrap();
|
||||
let schema = table.schema().await.unwrap();
|
||||
assert!(
|
||||
schema.field_with_name(ASSIGNMENT).is_err(),
|
||||
"the assignment column outlived its binding: {:?}",
|
||||
schema.fields().iter().map(|f| f.name()).collect::<Vec<_>>()
|
||||
);
|
||||
assert!(bindings_of(&table).await.is_empty());
|
||||
ensure_supported_function_metadata(&schema).unwrap();
|
||||
}
|
||||
|
||||
/// An input error must not cost the binding. The retirement commits
|
||||
/// first, so a projection that could never succeed is refused before it,
|
||||
/// leaving the table exactly as it was.
|
||||
#[rstest::rstest]
|
||||
#[case::unknown_column(
|
||||
&["search_text", "search_token_count", "nope"],
|
||||
"does not exist"
|
||||
)]
|
||||
#[case::every_column(
|
||||
&["title", "body", "search_text", "search_token_count", "spare"],
|
||||
"Cannot drop all columns"
|
||||
)]
|
||||
#[case::every_column_quoted(
|
||||
&["`title`", "`body`", "`search_text`", "`search_token_count`", "`spare`"],
|
||||
"Cannot drop all columns"
|
||||
)]
|
||||
#[tokio::test]
|
||||
async fn an_invalid_drop_leaves_the_binding_in_place(
|
||||
#[case] columns: &[&str],
|
||||
#[case] expected: &str,
|
||||
) {
|
||||
let table = bound_table().await;
|
||||
let version = table.version().await.unwrap();
|
||||
|
||||
let err = table.drop_columns(columns).await.unwrap_err();
|
||||
let Error::InvalidInput { message } = &err else {
|
||||
panic!("{err:?}");
|
||||
};
|
||||
assert!(message.contains(expected), "{message}");
|
||||
|
||||
// Read durable state, not this handle: the error path never calls
|
||||
// `dataset.update`, so the cached handle would report the old version
|
||||
// even if the unbind had committed.
|
||||
table.checkout_latest().await.unwrap();
|
||||
assert_eq!(table.version().await.unwrap(), version);
|
||||
let schema = table.schema().await.unwrap();
|
||||
assert!(schema.metadata().contains_key(FUNCTION_BINDINGS_META_KEY));
|
||||
assert_eq!(bindings_of(&table).await.len(), 1);
|
||||
for column in ["search_text", "search_token_count"] {
|
||||
assert!(
|
||||
schema
|
||||
.field_with_name(column)
|
||||
.unwrap()
|
||||
.metadata()
|
||||
.keys()
|
||||
.any(|key| is_declaration_key(key)),
|
||||
"{column} lost its declaration metadata"
|
||||
);
|
||||
}
|
||||
ensure_supported_function_metadata(&schema).unwrap();
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_update_field_metadata() {
|
||||
let conn = connect("memory://").execute().await.unwrap();
|
||||
|
||||
Reference in New Issue
Block a user