1use std::hash::{BuildHasher, Hash, Hasher};
19
20use columnation::Columnation;
21use differential_dataflow::consolidation::ConsolidatingContainerBuilder;
22use differential_dataflow::difference::{Multiply, Semigroup};
23use differential_dataflow::lattice::Lattice;
24use differential_dataflow::trace::Batcher;
25use differential_dataflow::{AsCollection, Collection, Hashable, VecCollection};
26use timely::container::{DrainContainer, PushInto};
27use timely::dataflow::channels::pact::{Exchange, ParallelizationContract, Pipeline};
28use timely::dataflow::operators::Capability;
29use timely::dataflow::operators::generic::builder_rc::{
30 OperatorBuilder as OperatorBuilderRc, OperatorBuilder,
31};
32use timely::dataflow::operators::generic::operator::{self, Operator};
33use timely::dataflow::operators::generic::{
34 InputHandleCore, OperatorInfo, OutputBuilder, OutputBuilderSession,
35};
36use timely::dataflow::{Scope, Stream, StreamVec};
37use timely::progress::operate::FrontierInterest;
38use timely::progress::{Antichain, Timestamp};
39use timely::{Container, ContainerBuilder, PartialOrder};
40
41use crate::columnation::{ColumnationChunker, ColumnationStack};
42
43pub trait StreamExt<'scope, T, C1>
45where
46 T: Timestamp,
47 C1: Container + DrainContainer + Clone + 'static,
48{
49 fn unary_fallible<DCB, ECB, B, P>(
60 self,
61 pact: P,
62 name: &str,
63 constructor: B,
64 ) -> (
65 Stream<'scope, T, DCB::Container>,
66 Stream<'scope, T, ECB::Container>,
67 )
68 where
69 DCB: ContainerBuilder,
70 ECB: ContainerBuilder,
71 B: FnOnce(
72 Capability<T>,
73 OperatorInfo,
74 ) -> Box<
75 dyn FnMut(
76 &mut InputHandleCore<T, C1, P::Puller>,
77 &mut OutputBuilderSession<'_, T, DCB>,
78 &mut OutputBuilderSession<'_, T, ECB>,
79 ) + 'static,
80 >,
81 P: ParallelizationContract<T, C1>;
82
83 fn flat_map_fallible<DCB, ECB, D2, E, I, L>(
88 self,
89 name: &str,
90 logic: L,
91 ) -> (
92 Stream<'scope, T, DCB::Container>,
93 Stream<'scope, T, ECB::Container>,
94 )
95 where
96 DCB: ContainerBuilder + PushInto<D2>,
97 ECB: ContainerBuilder + PushInto<E>,
98 I: IntoIterator<Item = Result<D2, E>>,
99 L: for<'a> FnMut(C1::Item<'a>) -> I + 'static;
100
101 fn partition_by<CB, L>(
109 self,
110 name: &str,
111 predicate: L,
112 ) -> (
113 Stream<'scope, T, CB::Container>,
114 Stream<'scope, T, CB::Container>,
115 )
116 where
117 CB: ContainerBuilder + for<'a> PushInto<C1::Item<'a>>,
118 L: for<'a> FnMut(&C1::Item<'a>) -> bool + 'static;
119
120 fn expire_stream_at(self, name: &str, expiration: T) -> Stream<'scope, T, C1>;
122}
123
124pub trait CollectionExt<'scope, T, D1, R>: Sized
126where
127 T: Timestamp,
128 R: Semigroup,
129{
130 fn empty(scope: Scope<'scope, T>) -> VecCollection<'scope, T, D1, R>;
132
133 fn map_fallible<DCB, ECB, D2, E, L>(
142 self,
143 name: &str,
144 mut logic: L,
145 ) -> (
146 VecCollection<'scope, T, D2, R>,
147 VecCollection<'scope, T, E, R>,
148 )
149 where
150 DCB: ContainerBuilder<Container = Vec<(D2, T, R)>> + PushInto<(D2, T, R)>,
151 ECB: ContainerBuilder<Container = Vec<(E, T, R)>> + PushInto<(E, T, R)>,
152 D2: Clone + 'static,
153 E: Clone + 'static,
154 L: FnMut(D1) -> Result<D2, E> + 'static,
155 {
156 self.flat_map_fallible::<DCB, ECB, _, _, _, _>(name, move |record| Some(logic(record)))
157 }
158
159 fn flat_map_fallible<DCB, ECB, D2, E, I, L>(
164 self,
165 name: &str,
166 logic: L,
167 ) -> (
168 Collection<'scope, T, DCB::Container>,
169 Collection<'scope, T, ECB::Container>,
170 )
171 where
172 DCB: ContainerBuilder + PushInto<(D2, T, R)>,
173 ECB: ContainerBuilder + PushInto<(E, T, R)>,
174 D2: Clone + 'static,
175 E: Clone + 'static,
176 I: IntoIterator<Item = Result<D2, E>>,
177 L: FnMut(D1) -> I + 'static;
178
179 fn expire_collection_at(self, name: &str, expiration: T) -> VecCollection<'scope, T, D1, R>;
181
182 fn explode_one<D2, R2, L>(
187 self,
188 logic: L,
189 ) -> VecCollection<'scope, T, D2, <R2 as Multiply<R>>::Output>
190 where
191 D2: differential_dataflow::Data,
192 R2: Semigroup + Multiply<R>,
193 <R2 as Multiply<R>>::Output: Clone + 'static + Semigroup,
194 L: FnMut(D1) -> (D2, R2) + 'static,
195 T: Lattice;
196
197 fn ensure_monotonic<E, IE>(
202 self,
203 into_err: IE,
204 ) -> (
205 VecCollection<'scope, T, D1, R>,
206 VecCollection<'scope, T, E, R>,
207 )
208 where
209 E: Clone + 'static,
210 IE: Fn(D1, R) -> (E, R) + 'static,
211 R: num_traits::sign::Signed;
212
213 fn consolidate_named_if<Ba>(self, must_consolidate: bool, name: &str) -> Self
216 where
217 D1: differential_dataflow::ExchangeData + Hash + Columnation,
218 R: Semigroup + differential_dataflow::ExchangeData + Columnation,
219 T: Lattice + Columnation,
220 Ba: Batcher<Time = T, Output = ColumnationStack<((D1, ()), T, R)>> + 'static;
221
222 fn consolidate_named<Ba>(self, name: &str) -> Self
224 where
225 D1: differential_dataflow::ExchangeData + Hash + Columnation,
226 R: Semigroup + differential_dataflow::ExchangeData + Columnation,
227 T: Lattice + Columnation,
228 Ba: Batcher<Time = T, Output = ColumnationStack<((D1, ()), T, R)>> + 'static;
229}
230
231impl<'scope, T, C1> StreamExt<'scope, T, C1> for Stream<'scope, T, C1>
232where
233 T: Timestamp,
234 C1: Container + DrainContainer + Clone + 'static,
235{
236 fn unary_fallible<DCB, ECB, B, P>(
237 self,
238 pact: P,
239 name: &str,
240 constructor: B,
241 ) -> (
242 Stream<'scope, T, DCB::Container>,
243 Stream<'scope, T, ECB::Container>,
244 )
245 where
246 DCB: ContainerBuilder,
247 ECB: ContainerBuilder,
248 B: FnOnce(
249 Capability<T>,
250 OperatorInfo,
251 ) -> Box<
252 dyn FnMut(
253 &mut InputHandleCore<T, C1, P::Puller>,
254 &mut OutputBuilderSession<'_, T, DCB>,
255 &mut OutputBuilderSession<'_, T, ECB>,
256 ) + 'static,
257 >,
258 P: ParallelizationContract<T, C1>,
259 {
260 let mut builder = OperatorBuilderRc::new(name.into(), self.scope());
261
262 let operator_info = builder.operator_info();
263
264 let mut input = builder.new_input(self.clone(), pact);
265 builder.set_notify_for(0, FrontierInterest::Never);
266 let (ok_output, ok_stream) = builder.new_output();
267 let mut ok_output = OutputBuilder::from(ok_output);
268 let (err_output, err_stream) = builder.new_output();
269 let mut err_output = OutputBuilder::from(err_output);
270
271 builder.build(move |mut capabilities| {
272 let capability = capabilities.pop().unwrap();
274 let mut logic = constructor(capability, operator_info);
275 move |_frontiers| {
276 let mut ok_output_handle = ok_output.activate();
277 let mut err_output_handle = err_output.activate();
278 logic(&mut input, &mut ok_output_handle, &mut err_output_handle);
279 }
280 });
281
282 (ok_stream, err_stream)
283 }
284
285 #[allow(clippy::redundant_closure)]
291 fn flat_map_fallible<DCB, ECB, D2, E, I, L>(
292 self,
293 name: &str,
294 mut logic: L,
295 ) -> (
296 Stream<'scope, T, DCB::Container>,
297 Stream<'scope, T, ECB::Container>,
298 )
299 where
300 DCB: ContainerBuilder + PushInto<D2>,
301 ECB: ContainerBuilder + PushInto<E>,
302 I: IntoIterator<Item = Result<D2, E>>,
303 L: for<'a> FnMut(C1::Item<'a>) -> I + 'static,
304 {
305 self.unary_fallible::<DCB, ECB, _, _>(Pipeline, name, move |_, _| {
306 Box::new(move |input, ok_output, err_output| {
307 input.for_each_time(|time, data| {
308 let mut ok_session = ok_output.session_with_builder(&time);
309 let mut err_session = err_output.session_with_builder(&time);
310 for r in data
311 .flat_map(DrainContainer::drain)
312 .flat_map(|d1| logic(d1))
313 {
314 match r {
315 Ok(d2) => ok_session.give(d2),
316 Err(e) => err_session.give(e),
317 }
318 }
319 })
320 })
321 })
322 }
323
324 fn partition_by<CB, L>(
325 self,
326 name: &str,
327 mut predicate: L,
328 ) -> (
329 Stream<'scope, T, CB::Container>,
330 Stream<'scope, T, CB::Container>,
331 )
332 where
333 CB: ContainerBuilder + for<'a> PushInto<C1::Item<'a>>,
334 L: for<'a> FnMut(&C1::Item<'a>) -> bool + 'static,
335 {
336 self.unary_fallible::<CB, CB, _, _>(Pipeline, name, move |_, _| {
337 Box::new(move |input, matching_output, rest_output| {
338 input.for_each_time(|time, data| {
339 let mut matching = matching_output.session_with_builder(&time);
340 let mut rest = rest_output.session_with_builder(&time);
341 for item in data.flat_map(DrainContainer::drain) {
342 if predicate(&item) {
343 matching.give(item);
344 } else {
345 rest.give(item);
346 }
347 }
348 })
349 })
350 })
351 }
352
353 fn expire_stream_at(self, name: &str, expiration: T) -> Stream<'scope, T, C1> {
354 let name = format!("expire_stream_at({name})");
355 self.unary_frontier(Pipeline, &name.clone(), move |cap, _| {
356 let cap = Some(cap.delayed(&expiration));
360 let mut warned = false;
361 move |(input, frontier), output| {
362 let _ = ∩
363 let frontier = frontier.frontier();
364 if !frontier.less_than(&expiration) && !warned {
365 tracing::warn!(
373 name = name,
374 frontier = ?frontier,
375 expiration = ?expiration,
376 "frontier not less than expiration"
377 );
378 warned = true;
379 }
380 input.for_each(|time, data| {
381 let mut session = output.session(&time);
382 session.give_container(data);
383 });
384 }
385 })
386 }
387}
388
389impl<'scope, T, D1, R> CollectionExt<'scope, T, D1, R> for VecCollection<'scope, T, D1, R>
390where
391 T: Timestamp + Clone + 'static,
392 D1: Clone + 'static,
393 R: Semigroup + 'static,
394{
395 fn empty(scope: Scope<'scope, T>) -> VecCollection<'scope, T, D1, R> {
396 operator::empty(scope).as_collection()
397 }
398
399 fn flat_map_fallible<DCB, ECB, D2, E, I, L>(
400 self,
401 name: &str,
402 mut logic: L,
403 ) -> (
404 Collection<'scope, T, DCB::Container>,
405 Collection<'scope, T, ECB::Container>,
406 )
407 where
408 DCB: ContainerBuilder + PushInto<(D2, T, R)>,
409 ECB: ContainerBuilder + PushInto<(E, T, R)>,
410 D2: Clone + 'static,
411 E: Clone + 'static,
412 I: IntoIterator<Item = Result<D2, E>>,
413 L: FnMut(D1) -> I + 'static,
414 {
415 let (ok_stream, err_stream) =
416 self.inner
417 .flat_map_fallible::<DCB, ECB, _, _, _, _>(name, move |(d1, t, r)| {
418 logic(d1).into_iter().map(move |res| match res {
419 Ok(d2) => Ok((d2, t.clone(), r.clone())),
420 Err(e) => Err((e, t.clone(), r.clone())),
421 })
422 });
423 (ok_stream.as_collection(), err_stream.as_collection())
424 }
425
426 fn expire_collection_at(self, name: &str, expiration: T) -> VecCollection<'scope, T, D1, R> {
427 self.inner
428 .expire_stream_at(name, expiration)
429 .as_collection()
430 }
431
432 fn explode_one<D2, R2, L>(
433 self,
434 mut logic: L,
435 ) -> VecCollection<'scope, T, D2, <R2 as Multiply<R>>::Output>
436 where
437 D2: differential_dataflow::Data,
438 R2: Semigroup + Multiply<R>,
439 <R2 as Multiply<R>>::Output: Clone + 'static + Semigroup,
440 L: FnMut(D1) -> (D2, R2) + 'static,
441 T: Lattice,
442 {
443 self.inner
444 .clone()
445 .unary::<ConsolidatingContainerBuilder<_>, _, _, _>(
446 Pipeline,
447 "ExplodeOne",
448 move |_, _| {
449 move |input, output| {
450 input.for_each(|time, data| {
451 output
452 .session_with_builder(&time)
453 .give_iterator(data.drain(..).map(|(x, t, d)| {
454 let (x, d2) = logic(x);
455 (x, t, d2.multiply(&d))
456 }));
457 });
458 }
459 },
460 )
461 .as_collection()
462 }
463
464 fn ensure_monotonic<E, IE>(
465 self,
466 into_err: IE,
467 ) -> (
468 VecCollection<'scope, T, D1, R>,
469 VecCollection<'scope, T, E, R>,
470 )
471 where
472 E: Clone + 'static,
473 IE: Fn(D1, R) -> (E, R) + 'static,
474 R: num_traits::sign::Signed,
475 {
476 let (oks, errs) = self
477 .inner
478 .unary_fallible(Pipeline, "EnsureMonotonic", move |_, _| {
479 Box::new(move |input, ok_output, err_output| {
480 input.for_each(|time, data| {
481 let mut ok_session = ok_output.session(&time);
482 let mut err_session = err_output.session(&time);
483 for (x, t, d) in data.drain(..) {
484 if d.is_positive() {
485 ok_session.give((x, t, d))
486 } else {
487 let (e, d2) = into_err(x, d);
488 err_session.give((e, t, d2))
489 }
490 }
491 })
492 })
493 });
494 (oks.as_collection(), errs.as_collection())
495 }
496
497 fn consolidate_named_if<Ba>(self, must_consolidate: bool, name: &str) -> Self
498 where
499 D1: differential_dataflow::ExchangeData + Hash + Columnation,
500 R: Semigroup + differential_dataflow::ExchangeData + Columnation,
501 T: Lattice + Ord + Columnation,
502 Ba: Batcher<Time = T, Output = ColumnationStack<((D1, ()), T, R)>> + 'static,
503 {
504 if must_consolidate {
505 let random_state = crate::hash::fixed_state();
518 let exchange = Exchange::new(move |update: &((D1, _), T, R)| {
519 let data = &(update.0).0;
520 let mut h = random_state.build_hasher();
521 data.hash(&mut h);
522 h.finish()
523 });
524 consolidate_pact::<ColumnationChunker<((D1, ()), T, R)>, Ba, _, _>(
525 self.map(|k| (k, ())).inner,
526 exchange,
527 name,
528 )
529 .unary(Pipeline, "unpack consolidated", |_, _| {
530 |input, output| {
531 input.for_each(|time, data| {
532 let mut session = output.session(&time);
533 for ((k, ()), t, d) in data.iter().flatten().flat_map(|chunk| chunk.iter())
534 {
535 session.give((k.clone(), t.clone(), d.clone()))
536 }
537 })
538 }
539 })
540 .as_collection()
541 } else {
542 self
543 }
544 }
545
546 fn consolidate_named<Ba>(self, name: &str) -> Self
547 where
548 D1: differential_dataflow::ExchangeData + Hash + Columnation,
549 R: Semigroup + differential_dataflow::ExchangeData + Columnation,
550 T: Lattice + Ord + Columnation,
551 Ba: Batcher<Time = T, Output = ColumnationStack<((D1, ()), T, R)>> + 'static,
552 {
553 let exchange = Exchange::new(move |update: &((D1, ()), T, R)| (update.0).0.hashed());
554
555 consolidate_pact::<ColumnationChunker<((D1, ()), T, R)>, Ba, _, _>(
556 self.map(|k| (k, ())).inner,
557 exchange,
558 name,
559 )
560 .unary(Pipeline, &format!("Unpack {name}"), |_, _| {
561 |input, output| {
562 input.for_each(|time, data| {
563 let mut session = output.session(&time);
564 for ((k, ()), t, d) in data.iter().flatten().flat_map(|chunk| chunk.iter()) {
565 session.give((k.clone(), t.clone(), d.clone()))
566 }
567 })
568 }
569 })
570 .as_collection()
571 }
572}
573
574pub fn consolidate_pact<'scope, Chu, Ba, C, P>(
581 stream: Stream<'scope, Ba::Time, C>,
582 pact: P,
583 name: &str,
584) -> StreamVec<'scope, Ba::Time, Vec<Ba::Output>>
585where
586 Ba: Batcher + 'static,
587 Chu: ContainerBuilder<Container = Ba::Output> + for<'a> PushInto<&'a mut C> + 'static,
588 C: Container + Clone + 'static,
589 Ba::Output: Clone,
590 P: ParallelizationContract<Ba::Time, C>,
591{
592 let logger = stream
593 .scope()
594 .worker()
595 .logger_for("differential/arrange")
596 .map(Into::into);
597 stream.unary_frontier(pact, name, |_cap, info| {
598 let mut batcher = Ba::new(logger, info.global_id);
601 let mut chunker = Chu::default();
604 let mut capabilities = Antichain::<Capability<Ba::Time>>::new();
606 let mut prev_frontier = Antichain::from_elem(Ba::Time::minimum());
607
608 move |(input, frontier), output| {
609 input.for_each(|cap, data| {
610 capabilities.insert(cap.retain(0));
611 chunker.push_into(data);
612 while let Some(chunk) = chunker.extract() {
613 batcher.push_into(std::mem::take(chunk));
614 }
615 });
616
617 if prev_frontier.borrow() != frontier.frontier() {
618 while let Some(chunk) = chunker.finish() {
621 batcher.push_into(std::mem::take(chunk));
622 }
623
624 if capabilities
625 .elements()
626 .iter()
627 .any(|c| !frontier.less_equal(c.time()))
628 {
629 let mut upper = Antichain::new(); for (index, capability) in capabilities.elements().iter().enumerate() {
633 if !frontier.less_equal(capability.time()) {
634 upper.clear();
638 for time in frontier.frontier().iter() {
639 upper.insert(time.clone());
640 }
641 for other_capability in &capabilities.elements()[(index + 1)..] {
642 upper.insert(other_capability.time().clone());
643 }
644
645 let mut session = output.session(&capabilities.elements()[index]);
647 let (chain, _description) = batcher.seal(upper.clone());
649 session.give(chain);
650 }
651 }
652
653 let mut new_capabilities = Antichain::new();
659 for time in batcher.frontier().iter() {
660 if let Some(capability) = capabilities
661 .elements()
662 .iter()
663 .find(|c| c.time().less_equal(time))
664 {
665 new_capabilities.insert(capability.delayed(time));
666 } else {
667 panic!("failed to find capability");
668 }
669 }
670
671 capabilities = new_capabilities;
672 }
673
674 prev_frontier.clear();
675 prev_frontier.extend(frontier.frontier().iter().cloned());
676 }
677 }
678 })
679}
680
681pub trait ConcatenateFlatten<'scope, T: Timestamp, C: Container + DrainContainer> {
683 fn concatenate_flatten<I, CB>(&self, sources: I) -> Stream<'scope, T, CB::Container>
704 where
705 I: IntoIterator<Item = Stream<'scope, T, C>>,
706 CB: ContainerBuilder + for<'a> PushInto<C::Item<'a>>;
707}
708
709impl<'scope, T, C> ConcatenateFlatten<'scope, T, C> for Stream<'scope, T, C>
710where
711 T: Timestamp,
712 C: Container + DrainContainer + Clone + 'static,
713{
714 fn concatenate_flatten<I, CB>(&self, sources: I) -> Stream<'scope, T, CB::Container>
715 where
716 I: IntoIterator<Item = Stream<'scope, T, C>>,
717 CB: ContainerBuilder + for<'a> PushInto<C::Item<'a>>,
718 {
719 self.scope()
720 .concatenate_flatten::<_, CB>(Some(Clone::clone(self)).into_iter().chain(sources))
721 }
722}
723
724impl<'scope, T, C> ConcatenateFlatten<'scope, T, C> for Scope<'scope, T>
725where
726 T: Timestamp,
727 C: Container + DrainContainer,
728{
729 fn concatenate_flatten<I, CB>(&self, sources: I) -> Stream<'scope, T, CB::Container>
730 where
731 I: IntoIterator<Item = Stream<'scope, T, C>>,
732 CB: ContainerBuilder + for<'a> PushInto<C::Item<'a>>,
733 {
734 let mut builder = OperatorBuilder::new("ConcatenateFlatten".to_string(), self.clone());
735
736 let mut handles = sources
738 .into_iter()
739 .map(|s| builder.new_input(s, Pipeline))
740 .collect::<Vec<_>>();
741 for i in 0..handles.len() {
742 builder.set_notify_for(i, FrontierInterest::Never);
743 }
744
745 let (output, result) = builder.new_output::<CB::Container>();
747 let mut output = OutputBuilder::<_, CB>::from(output);
748
749 builder.build(move |_capability| {
750 move |_frontier| {
751 let mut output = output.activate();
752 for handle in handles.iter_mut() {
753 handle.for_each_time(|time, data| {
754 output
755 .session_with_builder(&time)
756 .give_iterator(data.flat_map(DrainContainer::drain));
757 })
758 }
759 }
760 });
761
762 result
763 }
764}
765
766pub trait ClearContainer {
768 fn clear(&mut self);
770}
771
772impl<T> ClearContainer for Vec<T> {
773 fn clear(&mut self) {
774 Vec::clear(self)
775 }
776}
777
778#[cfg(test)]
779mod tests {
780 use timely::container::CapacityContainerBuilder;
781 use timely::dataflow::operators::Capture;
782 use timely::dataflow::operators::capture::Extract;
783 use timely::dataflow::operators::core::to_stream::ToStreamBuilder;
784 use timely::dataflow::operators::vec::ToStream;
785
786 use crate::columnar::Column;
787 use crate::columnar::builder::ColumnBuilder;
788
789 use super::*;
790
791 const UPDATES: [(u64, u64, i64); 4] = [(0, 0, 1), (1, 3, 1), (2, 5, -1), (3, 7, 1)];
794 const SPLIT: u64 = 5;
795 const MATCHING: [(u64, u64, i64); 2] = [(2, 5, -1), (3, 7, 1)];
796 const REST: [(u64, u64, i64); 2] = [(0, 0, 1), (1, 3, 1)];
797
798 #[mz_ore::test]
799 fn partition_by_routes_vec_records_by_their_own_time() {
800 let (matching, rest) = timely::execute_directly(|worker| {
801 worker.dataflow::<u64, _, _>(|scope| {
802 let (matching, rest) = UPDATES
803 .to_vec()
804 .to_stream(scope)
805 .partition_by::<CapacityContainerBuilder<Vec<_>>, _>("Test", |(_, time, _)| {
806 *time >= SPLIT
807 });
808 (matching.capture(), rest.capture())
809 })
810 });
811 let flatten = |captured: Vec<(u64, Vec<(u64, u64, i64)>)>| {
812 captured
813 .into_iter()
814 .flat_map(|(capability, updates)| {
815 assert_eq!(capability, 0, "outputs keep the input capability");
816 updates
817 })
818 .collect::<Vec<_>>()
819 };
820 assert_eq!(flatten(matching.extract()), MATCHING);
821 assert_eq!(flatten(rest.extract()), REST);
822 }
823
824 #[mz_ore::test]
825 fn partition_by_routes_columnar_records_by_their_own_time() {
826 let (matching, rest) = timely::execute_directly(|worker| {
827 worker.dataflow::<u64, _, _>(|scope| {
828 let (matching, rest) = UPDATES
829 .to_vec()
830 .to_stream_with_builder::<_, ColumnBuilder<(u64, u64, i64)>>(scope)
831 .partition_by::<ColumnBuilder<(u64, u64, i64)>, _>("Test", |(_, time, _)| {
832 **time >= SPLIT
833 });
834 (matching.capture(), rest.capture())
835 })
836 });
837 let flatten = |captured: Vec<(u64, Column<(u64, u64, i64)>)>| {
838 captured
839 .into_iter()
840 .flat_map(|(capability, mut column)| {
841 assert_eq!(capability, 0, "outputs keep the input capability");
842 column
843 .drain()
844 .map(|(data, time, diff)| (*data, *time, *diff))
845 .collect::<Vec<_>>()
846 })
847 .collect::<Vec<_>>()
848 };
849 assert_eq!(flatten(matching.extract()), MATCHING);
850 assert_eq!(flatten(rest.extract()), REST);
851 }
852}