1use std::any::Any;
16use std::sync::Arc;
17
18use common_query::request::{
19 DynFilterPayload, INITIAL_REMOTE_DYN_FILTER_REGISTRATIONS_EXTENSION_KEY, InitialDynFilterReg,
20 InitialDynFilterRegs, InitialDynFilterSnapshot, REMOTE_DYN_FILTER_PAYLOAD_MAX_BYTES,
21};
22use datafusion_common::Result;
23use datafusion_physical_expr::PhysicalExpr;
24use datafusion_physical_expr::expressions::DynamicFilterPhysicalExpr;
25use session::context::{QueryContext, QueryContextRef};
26use store_api::storage::RegionId;
27
28use crate::dist_plan::filter_id::build_remote_dyn_filter_id;
29use crate::dist_plan::{
30 FilterId, QueryDynFilterRegistry, RemoteDynFilterProducerId, Subscriber, SubscriberRegistration,
31};
32use crate::region_query::RegionQueryTarget;
33
34#[derive(Debug, Clone)]
35pub(crate) struct CapturedDynFilter {
36 filter_id: FilterId,
37 initial_registration: InitialDynFilterReg,
38 pub(crate) alive_dyn_filter: Arc<DynamicFilterPhysicalExpr>,
39}
40
41#[derive(Debug, Clone)]
42pub(crate) struct RemoteDynFilterPushdown {
43 pub(crate) captured_dyn_filters: Vec<CapturedDynFilter>,
44 pub(crate) pushed_down: Vec<bool>,
46}
47
48pub(crate) fn capture_remote_dyn_filters_for_pushdown(
49 remote_dyn_filter_producer_id: RemoteDynFilterProducerId,
50 parent_filters: Vec<Arc<dyn datafusion::physical_plan::PhysicalExpr>>,
51) -> RemoteDynFilterPushdown {
52 let mut pushed_down = Vec::with_capacity(parent_filters.len());
53 let mut captured_dyn_filters = Vec::new();
54
55 for (producer_local_ordinal, filter) in parent_filters.into_iter().enumerate() {
56 let Some(alive_dyn_filter) = downcast_dynamic_filter(filter) else {
57 pushed_down.push(false);
58 continue;
59 };
60
61 match build_captured_dyn_filter(
62 remote_dyn_filter_producer_id,
63 producer_local_ordinal,
64 alive_dyn_filter,
65 ) {
66 Ok(captured_dyn_filter) => {
67 pushed_down.push(true);
68 captured_dyn_filters.push(captured_dyn_filter);
69 }
70 Err(error) => {
71 common_telemetry::warn!(error; "Remote dyn filter is not pushed down because initial registration cannot be built");
72 pushed_down.push(false);
73 }
74 }
75 }
76
77 if let Err(error) = validate_initial_registrations_for_pushdown(&captured_dyn_filters) {
78 common_telemetry::warn!(error; "Remote dyn filters are not pushed down because initial registrations are invalid");
79 return RemoteDynFilterPushdown {
80 captured_dyn_filters: Vec::new(),
81 pushed_down: vec![false; pushed_down.len()],
82 };
83 }
84
85 RemoteDynFilterPushdown {
86 captured_dyn_filters,
87 pushed_down,
88 }
89}
90
91fn downcast_dynamic_filter(
92 expr: Arc<dyn datafusion::physical_plan::PhysicalExpr>,
93) -> Option<Arc<DynamicFilterPhysicalExpr>> {
94 (expr as Arc<dyn Any + Send + Sync + 'static>)
95 .downcast::<DynamicFilterPhysicalExpr>()
96 .ok()
97}
98
99pub(crate) fn register_remote_dyn_filters(
100 registry: &QueryDynFilterRegistry,
101 captured_dyn_filters: &[CapturedDynFilter],
102) {
103 for captured_dyn_filter in captured_dyn_filters {
104 let _ = registry.register_remote_dyn_filter(
105 captured_dyn_filter.filter_id.clone(),
106 captured_dyn_filter.alive_dyn_filter.clone(),
107 );
108 }
109}
110
111pub(crate) fn register_dyn_filter_subscribers_for_region(
112 registry: &QueryDynFilterRegistry,
113 region_id: RegionId,
114 target: RegionQueryTarget,
115 captured_dyn_filters: &[CapturedDynFilter],
116) -> Vec<(FilterId, Subscriber)> {
117 let mut added = Vec::new();
118 for captured_dyn_filter in captured_dyn_filters {
119 let subscriber = Subscriber::new(region_id, target.clone());
120 let registration =
121 registry.register_subscriber(&captured_dyn_filter.filter_id, subscriber.clone());
122 match registration {
123 SubscriberRegistration::Added => {
124 added.push((captured_dyn_filter.filter_id.clone(), subscriber))
125 }
126 SubscriberRegistration::Duplicate => {}
127 SubscriberRegistration::MissingFilter => {
128 common_telemetry::warn!(
129 "Remote dynamic filter {} missing when registering subscriber for region {} at target {:?}",
130 captured_dyn_filter.filter_id,
131 region_id,
132 target
133 );
134 }
135 }
136 }
137 added
138}
139
140fn build_captured_dyn_filter(
141 remote_dyn_filter_producer_id: RemoteDynFilterProducerId,
142 producer_local_ordinal: usize,
143 alive_dyn_filter: Arc<DynamicFilterPhysicalExpr>,
144) -> Result<CapturedDynFilter> {
145 let children = alive_dyn_filter
146 .children()
147 .into_iter()
148 .cloned()
149 .collect::<Vec<_>>();
150 let filter_id = build_remote_dyn_filter_id(
151 remote_dyn_filter_producer_id,
152 producer_local_ordinal,
153 &children,
154 )?;
155 let initial_registration =
156 InitialDynFilterReg::from_filter_id_and_children(filter_id.to_string(), &children)?;
157
158 Ok(CapturedDynFilter {
159 filter_id,
160 initial_registration: attach_initial_snapshot(initial_registration, &alive_dyn_filter),
161 alive_dyn_filter,
162 })
163}
164
165fn validate_initial_registrations_for_pushdown(
166 captured_dyn_filters: &[CapturedDynFilter],
167) -> std::result::Result<(), String> {
168 let regs = build_initial_dyn_filter_regs_for_region(captured_dyn_filters);
169 regs.validate_default_bounds()?;
170 regs.to_extension_value()
171 .map_err(|error| error.to_string())?;
172 Ok(())
173}
174
175fn attach_initial_snapshot(
176 initial_registration: InitialDynFilterReg,
177 alive_dyn_filter: &DynamicFilterPhysicalExpr,
178) -> InitialDynFilterReg {
179 let Some(initial_snapshot) = initial_snapshot(alive_dyn_filter) else {
180 return initial_registration;
181 };
182
183 initial_registration.with_initial_snapshot(initial_snapshot)
184}
185
186fn initial_snapshot(
187 alive_dyn_filter: &DynamicFilterPhysicalExpr,
188) -> Option<InitialDynFilterSnapshot> {
189 let generation = alive_dyn_filter.snapshot_generation();
190 let current = match alive_dyn_filter.current() {
191 Ok(current) => current,
192 Err(error) => {
193 common_telemetry::warn!(error; "Failed to read remote dyn filter initial snapshot");
194 return None;
195 }
196 };
197
198 let payload = match DynFilterPayload::from_datafusion_expr(
199 ¤t,
200 REMOTE_DYN_FILTER_PAYLOAD_MAX_BYTES,
201 ) {
202 Ok(payload) => payload,
203 Err(error) => {
204 common_telemetry::warn!(error; "Failed to encode remote dyn filter initial snapshot");
205 return None;
206 }
207 };
208
209 let is_complete = false;
211 Some(InitialDynFilterSnapshot::new(
212 payload,
213 generation,
214 is_complete,
215 ))
216}
217
218fn build_initial_dyn_filter_regs_for_region(
219 captured_dyn_filters: &[CapturedDynFilter],
220) -> InitialDynFilterRegs {
221 InitialDynFilterRegs::new(
222 captured_dyn_filters
223 .iter()
224 .map(|captured| captured.initial_registration.clone())
225 .collect(),
226 )
227}
228
229pub(crate) fn query_context_with_initial_dyn_filter_regs(
230 query_ctx: &QueryContextRef,
231 region_id: RegionId,
232 captured_dyn_filters: &[CapturedDynFilter],
233) -> QueryContext {
234 let regs = build_initial_dyn_filter_regs_for_region(captured_dyn_filters);
235 query_context_with_initial_dyn_filter_regs_value(query_ctx, region_id, ®s)
236}
237
238pub(crate) fn query_context_with_refreshed_initial_dyn_filter_regs(
239 query_ctx: &QueryContextRef,
240 _region_id: RegionId,
241 captured_dyn_filters: &[CapturedDynFilter],
242) -> QueryContext {
243 let refreshed = InitialDynFilterRegs::new(
244 captured_dyn_filters
245 .iter()
246 .map(|captured| {
247 attach_initial_snapshot(
248 captured.initial_registration.clone(),
249 &captured.alive_dyn_filter,
250 )
251 })
252 .collect(),
253 );
254 let serialized = match serialize_initial_dyn_filter_regs(&refreshed) {
255 Ok(serialized) => serialized,
256 Err(error) => {
257 common_telemetry::warn!(error; "Failed to refresh remote dynamic filter initial registrations; using optimization-time snapshots");
258 return query_context_with_initial_dyn_filter_regs(
259 query_ctx,
260 _region_id,
261 captured_dyn_filters,
262 );
263 }
264 };
265 query_context_with_serialized_initial_dyn_filter_regs(query_ctx, serialized)
266}
267
268fn query_context_with_initial_dyn_filter_regs_value(
269 query_ctx: &QueryContextRef,
270 region_id: RegionId,
271 regs: &InitialDynFilterRegs,
272) -> QueryContext {
273 let region_query_ctx = query_ctx.as_ref().clone();
274 if regs.is_empty() {
275 return region_query_ctx;
276 }
277
278 match serialize_initial_dyn_filter_regs(regs) {
279 Ok(serialized) => {
280 return query_context_with_serialized_initial_dyn_filter_regs(query_ctx, serialized);
281 }
282 Err(error) => {
283 common_telemetry::warn!(error; "Failed to serialize initial remote dyn filter registrations for region {}", region_id)
284 }
285 }
286
287 region_query_ctx
288}
289
290fn query_context_with_serialized_initial_dyn_filter_regs(
291 query_ctx: &QueryContextRef,
292 serialized: String,
293) -> QueryContext {
294 let mut region_query_ctx = query_ctx.as_ref().clone();
295 region_query_ctx.set_extension(
296 INITIAL_REMOTE_DYN_FILTER_REGISTRATIONS_EXTENSION_KEY,
297 serialized,
298 );
299 region_query_ctx
300}
301
302fn serialize_initial_dyn_filter_regs(
303 regs: &InitialDynFilterRegs,
304) -> std::result::Result<String, String> {
305 regs.validate_default_bounds()?;
306 regs.to_extension_value().map_err(|error| error.to_string())
307}
308
309#[cfg(test)]
310mod tests {
311 use std::fmt;
312 use std::hash::{Hash, Hasher};
313
314 use common_meta::peer::Peer;
315 use datafusion::execution::TaskContext;
316 use datafusion_common::ScalarValue;
317 use datafusion_expr::ColumnarValue;
318 use datafusion_physical_expr::expressions::{Column, lit};
319 use session::query_id::QueryId;
320 use uuid::Uuid;
321
322 use super::*;
323 use crate::region_query::RegionQueryTarget;
324
325 #[derive(Debug)]
326 struct UnserializableExpr;
327
328 impl fmt::Display for UnserializableExpr {
329 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
330 write!(f, "unserializable_expr")
331 }
332 }
333
334 impl Hash for UnserializableExpr {
335 fn hash<H: Hasher>(&self, state: &mut H) {
336 "unserializable_expr".hash(state);
337 }
338 }
339
340 impl PartialEq for UnserializableExpr {
341 fn eq(&self, _other: &Self) -> bool {
342 true
343 }
344 }
345
346 impl Eq for UnserializableExpr {}
347
348 impl datafusion_physical_expr::PhysicalExpr for UnserializableExpr {
349 fn data_type(
350 &self,
351 _input_schema: &arrow_schema::Schema,
352 ) -> datafusion_common::Result<arrow_schema::DataType> {
353 Ok(arrow_schema::DataType::Boolean)
354 }
355
356 fn nullable(
357 &self,
358 _input_schema: &arrow_schema::Schema,
359 ) -> datafusion_common::Result<bool> {
360 Ok(false)
361 }
362
363 fn evaluate(
364 &self,
365 _batch: &common_recordbatch::DfRecordBatch,
366 ) -> datafusion_common::Result<ColumnarValue> {
367 Ok(ColumnarValue::Scalar(ScalarValue::Boolean(Some(true))))
368 }
369
370 fn children(&self) -> Vec<&Arc<dyn datafusion_physical_expr::PhysicalExpr>> {
371 Vec::new()
372 }
373
374 fn with_new_children(
375 self: Arc<Self>,
376 _children: Vec<Arc<dyn datafusion_physical_expr::PhysicalExpr>>,
377 ) -> datafusion_common::Result<Arc<dyn datafusion_physical_expr::PhysicalExpr>> {
378 Ok(self)
379 }
380
381 fn fmt_sql(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
382 write!(f, "{self}")
383 }
384 }
385
386 fn test_query_id(value: u128) -> QueryId {
387 QueryId::from(Uuid::from_u128(value))
388 }
389
390 fn test_remote_dyn_filter_producer_id(value: u64) -> RemoteDynFilterProducerId {
391 RemoteDynFilterProducerId::new(value)
392 }
393
394 fn test_target(id: u64) -> RegionQueryTarget {
395 RegionQueryTarget::new(Peer {
396 id,
397 addr: format!("127.0.0.1:{id}"),
398 })
399 }
400
401 fn test_captured_dyn_filter(
402 remote_dyn_filter_producer_id: RemoteDynFilterProducerId,
403 producer_local_ordinal: usize,
404 column_name: &str,
405 column_index: usize,
406 ) -> CapturedDynFilter {
407 build_captured_dyn_filter(
408 remote_dyn_filter_producer_id,
409 producer_local_ordinal,
410 Arc::new(DynamicFilterPhysicalExpr::new(
411 vec![Arc::new(Column::new(column_name, column_index)) as Arc<_>],
412 lit(true) as _,
413 )),
414 )
415 .unwrap()
416 }
417
418 fn test_dyn_filter_with_snapshot_payload(
419 column_name: &str,
420 column_index: usize,
421 payload_bytes: usize,
422 ) -> Arc<DynamicFilterPhysicalExpr> {
423 let dyn_filter = Arc::new(DynamicFilterPhysicalExpr::new(
424 vec![Arc::new(Column::new(column_name, column_index)) as Arc<_>],
425 lit(true) as _,
426 ));
427 dyn_filter
428 .update(lit(ScalarValue::Utf8(Some("x".repeat(payload_bytes)))) as _)
429 .unwrap();
430 dyn_filter
431 }
432
433 fn refreshed_regs(captured: &[CapturedDynFilter]) -> InitialDynFilterRegs {
434 let query_ctx = QueryContext::arc();
435 let region_query_ctx = query_context_with_refreshed_initial_dyn_filter_regs(
436 &query_ctx,
437 RegionId::new(1024, 7),
438 captured,
439 );
440 InitialDynFilterRegs::from_extension_value(
441 region_query_ctx
442 .extension(INITIAL_REMOTE_DYN_FILTER_REGISTRATIONS_EXTENSION_KEY)
443 .unwrap(),
444 )
445 .unwrap()
446 }
447
448 fn decode_snapshot(snapshot: &InitialDynFilterSnapshot, column_name: &str) -> String {
449 snapshot
450 .payload
451 .decode_datafusion_expr(
452 &TaskContext::default(),
453 &arrow_schema::Schema::new(vec![arrow_schema::Field::new(
454 column_name,
455 arrow_schema::DataType::Utf8,
456 false,
457 )]),
458 REMOTE_DYN_FILTER_PAYLOAD_MAX_BYTES,
459 )
460 .unwrap()
461 .to_string()
462 }
463
464 #[test]
465 fn capture_remote_dyn_filters_for_pushdown_preserves_parent_filter_ordinals() {
466 let parent_filters = vec![
467 Arc::new(Column::new("service", 0)) as Arc<dyn datafusion::physical_plan::PhysicalExpr>,
468 Arc::new(DynamicFilterPhysicalExpr::new(
469 vec![Arc::new(Column::new("host", 1)) as Arc<_>],
470 lit(true) as _,
471 )) as Arc<dyn datafusion::physical_plan::PhysicalExpr>,
472 Arc::new(Column::new("zone", 2)) as Arc<dyn datafusion::physical_plan::PhysicalExpr>,
473 Arc::new(DynamicFilterPhysicalExpr::new(
474 vec![Arc::new(Column::new("pod", 3)) as Arc<_>],
475 lit(true) as _,
476 )) as Arc<dyn datafusion::physical_plan::PhysicalExpr>,
477 ];
478
479 let remote_dyn_filter_producer_id = test_remote_dyn_filter_producer_id(42);
480 let captured =
481 capture_remote_dyn_filters_for_pushdown(remote_dyn_filter_producer_id, parent_filters)
482 .captured_dyn_filters;
483
484 assert_eq!(captured.len(), 2);
485 assert_eq!(
486 captured[0].filter_id.remote_dyn_filter_producer_id(),
487 remote_dyn_filter_producer_id
488 );
489 assert_eq!(
490 captured[1].filter_id.remote_dyn_filter_producer_id(),
491 remote_dyn_filter_producer_id
492 );
493 assert_eq!(captured[0].filter_id.producer_ordinal(), 1);
494 assert_eq!(captured[1].filter_id.producer_ordinal(), 3);
495 }
496
497 #[test]
498 fn capture_remote_dyn_filters_for_pushdown_marks_only_valid_initial_regs() {
499 let parent_filters = vec![
500 Arc::new(Column::new("service", 0)) as Arc<dyn datafusion::physical_plan::PhysicalExpr>,
501 Arc::new(DynamicFilterPhysicalExpr::new(
502 vec![Arc::new(Column::new("host", 1)) as Arc<_>],
503 lit(true) as _,
504 )) as Arc<dyn datafusion::physical_plan::PhysicalExpr>,
505 Arc::new(Column::new("zone", 2)) as Arc<dyn datafusion::physical_plan::PhysicalExpr>,
506 ];
507
508 let remote_dyn_filter_producer_id = test_remote_dyn_filter_producer_id(42);
509 let pushdown =
510 capture_remote_dyn_filters_for_pushdown(remote_dyn_filter_producer_id, parent_filters);
511
512 assert_eq!(pushdown.pushed_down, vec![false, true, false]);
513 assert_eq!(pushdown.captured_dyn_filters.len(), 1);
514 assert_eq!(
515 pushdown.captured_dyn_filters[0]
516 .filter_id
517 .remote_dyn_filter_producer_id(),
518 remote_dyn_filter_producer_id
519 );
520 assert_eq!(
521 pushdown.captured_dyn_filters[0]
522 .filter_id
523 .producer_ordinal(),
524 1
525 );
526 assert!(
527 pushdown.captured_dyn_filters[0]
528 .initial_registration
529 .initial_snapshot
530 .is_some()
531 );
532 }
533
534 #[test]
535 fn capture_remote_dyn_filters_for_pushdown_rejects_unencodable_registration() {
536 let parent_filters = vec![Arc::new(DynamicFilterPhysicalExpr::new(
537 vec![Arc::new(UnserializableExpr) as Arc<_>],
538 lit(true) as _,
539 ))
540 as Arc<dyn datafusion::physical_plan::PhysicalExpr>];
541
542 let pushdown = capture_remote_dyn_filters_for_pushdown(
543 test_remote_dyn_filter_producer_id(42),
544 parent_filters,
545 );
546
547 assert_eq!(pushdown.pushed_down, vec![false]);
548 assert!(pushdown.captured_dyn_filters.is_empty());
549 }
550
551 #[test]
552 fn capture_remote_dyn_filters_for_pushdown_attaches_initial_snapshot() {
553 let parent_filters = vec![Arc::new(DynamicFilterPhysicalExpr::new(
554 vec![Arc::new(Column::new("host", 1)) as Arc<_>],
555 lit(true) as _,
556 ))
557 as Arc<dyn datafusion::physical_plan::PhysicalExpr>];
558
559 let pushdown = capture_remote_dyn_filters_for_pushdown(
560 test_remote_dyn_filter_producer_id(42),
561 parent_filters,
562 );
563
564 assert_eq!(pushdown.pushed_down, vec![true]);
565 assert!(
566 pushdown.captured_dyn_filters[0]
567 .initial_registration
568 .initial_snapshot
569 .is_some()
570 );
571 }
572
573 #[test]
574 fn capture_remote_dyn_filters_for_pushdown_attaches_initial_snapshot_after_update() {
575 let dyn_filter = Arc::new(DynamicFilterPhysicalExpr::new(
576 vec![Arc::new(Column::new("host", 1)) as Arc<_>],
577 lit(true) as _,
578 ));
579 dyn_filter.update(lit(false) as _).unwrap();
580 let parent_filters = vec![dyn_filter as Arc<dyn datafusion::physical_plan::PhysicalExpr>];
581
582 let pushdown = capture_remote_dyn_filters_for_pushdown(
583 test_remote_dyn_filter_producer_id(42),
584 parent_filters,
585 );
586
587 assert_eq!(pushdown.pushed_down, vec![true]);
588 let snapshot = pushdown.captured_dyn_filters[0]
589 .initial_registration
590 .initial_snapshot
591 .as_ref()
592 .unwrap();
593 assert_eq!(snapshot.generation, 2);
594 assert!(!snapshot.is_complete);
595 assert!(matches!(
596 snapshot.payload,
597 DynFilterPayload::Datafusion(ref bytes) if !bytes.is_empty()
598 ));
599 }
600
601 #[test]
602 fn capture_remote_dyn_filters_for_pushdown_rejects_oversized_snapshots() {
603 let oversized_total_snapshot_bytes = REMOTE_DYN_FILTER_PAYLOAD_MAX_BYTES * 3 / 5;
604 let parent_filters = vec![
605 test_dyn_filter_with_snapshot_payload("host", 0, oversized_total_snapshot_bytes)
606 as Arc<dyn datafusion::physical_plan::PhysicalExpr>,
607 test_dyn_filter_with_snapshot_payload("pod", 1, oversized_total_snapshot_bytes)
608 as Arc<dyn datafusion::physical_plan::PhysicalExpr>,
609 ];
610
611 let pushdown = capture_remote_dyn_filters_for_pushdown(
612 test_remote_dyn_filter_producer_id(42),
613 parent_filters,
614 );
615
616 assert_eq!(pushdown.pushed_down, vec![false, false]);
617 assert!(pushdown.captured_dyn_filters.is_empty());
618 }
619
620 #[test]
621 fn capture_remote_dyn_filters_for_pushdown_rejects_too_many_regs_with_snapshots() {
622 const TOO_MANY_INITIAL_REGS: usize = 65;
623
624 let parent_filters = (0..TOO_MANY_INITIAL_REGS)
625 .map(|ordinal| {
626 test_dyn_filter_with_snapshot_payload(&format!("host_{ordinal}"), ordinal, 1)
627 as Arc<dyn datafusion::physical_plan::PhysicalExpr>
628 })
629 .collect::<Vec<_>>();
630
631 let pushdown = capture_remote_dyn_filters_for_pushdown(
632 test_remote_dyn_filter_producer_id(42),
633 parent_filters,
634 );
635
636 assert!(pushdown.captured_dyn_filters.is_empty());
637 assert_eq!(pushdown.pushed_down, vec![false; TOO_MANY_INITIAL_REGS]);
638 }
639
640 #[test]
641 fn capture_remote_dyn_filters_for_pushdown_rejects_regs_exceeding_bounds() {
642 const TOO_MANY_INITIAL_REGS: usize = 65;
643
644 let parent_filters = (0..TOO_MANY_INITIAL_REGS)
645 .map(|_| {
646 Arc::new(DynamicFilterPhysicalExpr::new(
647 vec![Arc::new(Column::new("host", 0)) as Arc<_>],
648 lit(true) as _,
649 )) as Arc<dyn datafusion::physical_plan::PhysicalExpr>
650 })
651 .collect::<Vec<_>>();
652
653 let pushdown = capture_remote_dyn_filters_for_pushdown(
654 test_remote_dyn_filter_producer_id(42),
655 parent_filters,
656 );
657
658 assert!(pushdown.captured_dyn_filters.is_empty());
659 assert_eq!(pushdown.pushed_down, vec![false; TOO_MANY_INITIAL_REGS]);
660 }
661
662 #[test]
663 fn register_dyn_filter_subscribers_for_region_reuses_existing_entry() {
664 let registry = QueryDynFilterRegistry::new(test_query_id(1));
665 let captured_dyn_filters = vec![test_captured_dyn_filter(
666 test_remote_dyn_filter_producer_id(42),
667 2,
668 "host",
669 0,
670 )];
671 let first_region_id = RegionId::new(1024, 7);
672 let second_region_id = RegionId::new(1024, 8);
673
674 register_remote_dyn_filters(®istry, &captured_dyn_filters);
675 register_dyn_filter_subscribers_for_region(
676 ®istry,
677 first_region_id,
678 test_target(1),
679 &captured_dyn_filters,
680 );
681 register_dyn_filter_subscribers_for_region(
682 ®istry,
683 second_region_id,
684 test_target(2),
685 &captured_dyn_filters,
686 );
687
688 assert_eq!(registry.entry_count(), 1);
689 let entry = registry.entries().pop().unwrap();
690 assert_eq!(
691 entry.filter_id().remote_dyn_filter_producer_id(),
692 test_remote_dyn_filter_producer_id(42)
693 );
694 assert_eq!(entry.filter_id().producer_ordinal(), 2);
695 let subscribers = entry.subscribers();
696 assert_eq!(subscribers.len(), 2);
697 assert!(
698 subscribers
699 .iter()
700 .any(|subscriber| subscriber.region_id() == first_region_id)
701 );
702 assert!(
703 subscribers
704 .iter()
705 .any(|subscriber| subscriber.region_id() == second_region_id)
706 );
707 }
708
709 #[test]
710 fn register_remote_dyn_filters_keeps_independent_producer_ids_distinct() {
711 let registry = QueryDynFilterRegistry::new(test_query_id(1));
712 let make_filter = |remote_dyn_filter_producer_id| {
713 test_captured_dyn_filter(remote_dyn_filter_producer_id, 2, "host", 0)
714 };
715
716 register_remote_dyn_filters(
717 ®istry,
718 &[make_filter(test_remote_dyn_filter_producer_id(42))],
719 );
720 register_remote_dyn_filters(
721 ®istry,
722 &[make_filter(test_remote_dyn_filter_producer_id(43))],
723 );
724
725 assert_eq!(registry.entry_count(), 2);
726 }
727
728 #[test]
729 fn query_context_includes_region_initial_dyn_filter_regs() {
730 let captured_dyn_filters = vec![test_captured_dyn_filter(
731 test_remote_dyn_filter_producer_id(42),
732 2,
733 "host",
734 0,
735 )];
736 let region_id = RegionId::new(1024, 7);
737 let query_ctx = QueryContext::arc();
738
739 let region_query_ctx = query_context_with_initial_dyn_filter_regs(
740 &query_ctx,
741 region_id,
742 &captured_dyn_filters,
743 );
744 let extension = region_query_ctx
745 .extension(INITIAL_REMOTE_DYN_FILTER_REGISTRATIONS_EXTENSION_KEY)
746 .unwrap();
747 let regs = InitialDynFilterRegs::from_extension_value(extension).unwrap();
748 let decoded_children = regs.regs[0]
749 .decode_children(
750 &TaskContext::default(),
751 &arrow_schema::Schema::new(vec![arrow_schema::Field::new(
752 "host",
753 arrow_schema::DataType::Utf8,
754 false,
755 )]),
756 1024,
757 )
758 .unwrap();
759 assert_eq!(regs.regs.len(), 1);
760 assert_eq!(
761 regs.regs[0].filter_id,
762 captured_dyn_filters[0].filter_id.to_string()
763 );
764 assert_eq!(decoded_children.len(), 1);
765 assert!(decoded_children[0].is::<Column>());
766 }
767
768 #[test]
769 fn refreshed_initial_regs_use_live_filter_snapshot_without_mutating_capture() {
770 let dyn_filter = Arc::new(DynamicFilterPhysicalExpr::new(
771 vec![Arc::new(Column::new("host", 0)) as Arc<_>],
772 lit(true) as _,
773 ));
774 let captured = vec![
775 build_captured_dyn_filter(
776 test_remote_dyn_filter_producer_id(42),
777 0,
778 dyn_filter.clone(),
779 )
780 .unwrap(),
781 ];
782 let original_generation = captured[0]
783 .initial_registration
784 .initial_snapshot
785 .as_ref()
786 .unwrap()
787 .generation;
788 dyn_filter.update(lit(false) as _).unwrap();
789 dyn_filter.mark_complete();
790
791 let regs = refreshed_regs(&captured);
792
793 let snapshot = regs.regs[0].initial_snapshot.as_ref().unwrap();
794 assert_eq!(snapshot.generation, 2);
795 assert!(!snapshot.is_complete);
796 assert_ne!(decode_snapshot(snapshot, "host"), "true");
797 assert_eq!(
798 captured[0]
799 .initial_registration
800 .initial_snapshot
801 .as_ref()
802 .unwrap()
803 .generation,
804 original_generation
805 );
806 }
807
808 #[test]
809 fn refresh_keeps_old_snapshot_when_one_live_filter_cannot_encode() {
810 let good = Arc::new(DynamicFilterPhysicalExpr::new(
811 vec![Arc::new(Column::new("good", 0)) as Arc<_>],
812 lit(true) as _,
813 ));
814 let bad = Arc::new(DynamicFilterPhysicalExpr::new(
815 vec![Arc::new(Column::new("bad", 0)) as Arc<_>],
816 lit(true) as _,
817 ));
818 let captured = vec![
819 build_captured_dyn_filter(test_remote_dyn_filter_producer_id(42), 0, good.clone())
820 .unwrap(),
821 build_captured_dyn_filter(test_remote_dyn_filter_producer_id(42), 1, bad.clone())
822 .unwrap(),
823 ];
824 let old_bad_generation = captured[1]
825 .initial_registration
826 .initial_snapshot
827 .as_ref()
828 .unwrap()
829 .generation;
830 good.update(lit(false) as _).unwrap();
831 bad.update(Arc::new(UnserializableExpr) as _).unwrap();
832
833 let regs = refreshed_regs(&captured);
834 assert_eq!(
835 regs.regs[0].initial_snapshot.as_ref().unwrap().generation,
836 2
837 );
838 assert_eq!(
839 regs.regs[1].initial_snapshot.as_ref().unwrap().generation,
840 old_bad_generation
841 );
842 }
843
844 #[test]
845 fn refresh_aggregate_overflow_falls_back_to_original_snapshots() {
846 let payload_bytes = REMOTE_DYN_FILTER_PAYLOAD_MAX_BYTES * 3 / 5;
847 let first = test_dyn_filter_with_snapshot_payload("first", 0, 1);
848 let second = test_dyn_filter_with_snapshot_payload("second", 0, 1);
849 let captured = vec![
850 build_captured_dyn_filter(test_remote_dyn_filter_producer_id(42), 0, first.clone())
851 .unwrap(),
852 build_captured_dyn_filter(test_remote_dyn_filter_producer_id(42), 1, second.clone())
853 .unwrap(),
854 ];
855 let original = captured
856 .iter()
857 .map(|captured| {
858 captured
859 .initial_registration
860 .initial_snapshot
861 .clone()
862 .unwrap()
863 })
864 .collect::<Vec<_>>();
865 first
866 .update(lit(ScalarValue::Utf8(Some("x".repeat(payload_bytes)))) as _)
867 .unwrap();
868 second
869 .update(lit(ScalarValue::Utf8(Some("y".repeat(payload_bytes)))) as _)
870 .unwrap();
871
872 let regs = refreshed_regs(&captured);
873 for (refreshed, original) in regs.regs.iter().zip(original) {
874 let refreshed = refreshed.initial_snapshot.as_ref().unwrap();
875 assert_eq!(refreshed.generation, original.generation);
876 assert_eq!(refreshed.payload, original.payload);
877 }
878 }
879
880 #[test]
881 fn query_context_drops_initial_regs_when_duplicate_filter_ids_exceed_bounds() {
882 let captured_dyn_filters = vec![
883 test_captured_dyn_filter(test_remote_dyn_filter_producer_id(42), 2, "host", 0),
884 test_captured_dyn_filter(test_remote_dyn_filter_producer_id(42), 2, "host", 0),
885 ];
886 let region_id = RegionId::new(1024, 7);
887 let query_ctx = QueryContext::arc();
888
889 let region_query_ctx = query_context_with_initial_dyn_filter_regs(
890 &query_ctx,
891 region_id,
892 &captured_dyn_filters,
893 );
894
895 assert!(
896 region_query_ctx
897 .extension(INITIAL_REMOTE_DYN_FILTER_REGISTRATIONS_EXTENSION_KEY)
898 .is_none()
899 );
900 }
901}