Skip to content
Open
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
53 changes: 53 additions & 0 deletions datafusion/physical-expr/src/expressions/dynamic_filters.rs
Original file line number Diff line number Diff line change
Expand Up @@ -272,6 +272,15 @@ impl DynamicFilterPhysicalExpr {
});
}

/// Returns whether this dynamic filter has been marked complete.
///
/// This is a non-blocking snapshot of the completion state. It does not
/// synchronize with concurrent updates; use [`Self::wait_complete`] when
/// callers need to wait until completion.
pub fn is_complete(&self) -> bool {
self.inner.read().is_complete
}

/// Wait asynchronously for any update to this filter.
///
/// This method will return when [`Self::update`] is called and the generation increases.
Expand Down Expand Up @@ -636,11 +645,55 @@ mod test {

// Mark as complete immediately
dynamic_filter.mark_complete();
assert!(dynamic_filter.is_complete());

// wait_complete should return immediately
dynamic_filter.wait_complete().await;
}

#[test]
fn test_is_complete_initially_false() {
let dynamic_filter =
DynamicFilterPhysicalExpr::new(vec![], lit(42) as Arc<dyn PhysicalExpr>);

assert!(!dynamic_filter.is_complete());
}

#[test]
fn test_is_complete_false_after_update() {
let dynamic_filter =
DynamicFilterPhysicalExpr::new(vec![], lit(42) as Arc<dyn PhysicalExpr>);

dynamic_filter
.update(lit(43) as Arc<dyn PhysicalExpr>)
.expect("Failed to update expression");

assert!(!dynamic_filter.is_complete());
}

#[test]
fn test_is_complete_true_after_mark_complete() {
let dynamic_filter =
DynamicFilterPhysicalExpr::new(vec![], lit(42) as Arc<dyn PhysicalExpr>);

dynamic_filter.mark_complete();

assert!(dynamic_filter.is_complete());
}

#[test]
fn test_is_complete_true_after_update_then_mark_complete() {
let dynamic_filter =
DynamicFilterPhysicalExpr::new(vec![], lit(42) as Arc<dyn PhysicalExpr>);

dynamic_filter
.update(lit(43) as Arc<dyn PhysicalExpr>)
.expect("Failed to update expression");
dynamic_filter.mark_complete();

assert!(dynamic_filter.is_complete());
}

#[test]
fn test_with_new_children_independence() {
// Create a schema with columns a, b, c, d
Expand Down