MCPcopy Create free account
hub / github.com/comnik/declarative-dataflow / extend

Method extend

src/plan/hector.rs:976–1011  ·  view source on GitHub ↗
(
        &self,
        extenders: &mut [Extender<'a, S, P, E>],
    )

Source from the content-addressed store, hash-verified

974
975impl<'a, S: Scope, P: ExchangeData + Ord> ProposeExtensionMethod<'a, S, P> for Collection<S, P> {
976 fn extend<E: ExchangeData + Ord>(
977 &self,
978 extenders: &mut [Extender<'a, S, P, E>],
979 ) -> Collection<S, (P, E)> {
980 if extenders.is_empty() {
981 // @TODO don't panic
982 panic!("No extenders specified.");
983 } else if extenders.len() == 1 {
984 extenders[0].propose(&self.clone())
985 } else {
986 let mut counts = self.map(|p| (p, 1 << 31, 0));
987 for (index, extender) in extenders.iter_mut().enumerate() {
988 if let Some(new_counts) = extender.count(&counts, index) {
989 counts = new_counts;
990 }
991 }
992
993 let parts = counts
994 .inner
995 .partition(extenders.len() as u64, |((p, _, i), t, d)| {
996 (i as u64, (p, t, d))
997 });
998
999 let mut results = Vec::new();
1000 for (index, nominations) in parts.into_iter().enumerate() {
1001 let mut extensions = extenders[index].propose(&nominations.as_collection());
1002 for other in (0..extenders.len()).filter(|&x| x != index) {
1003 extensions = extenders[other].validate(&extensions);
1004 }
1005
1006 results.push(extensions.inner); // save extensions
1007 }
1008
1009 self.scope().concatenate(results).as_collection()
1010 }
1011 }
1012}
1013
1014struct ConstantExtender<P, V>

Callers 8

sinkMethod · 0.80
implementMethod · 0.80
countMethod · 0.80
proposeMethod · 0.80
validateMethod · 0.80
implementMethod · 0.80
implementMethod · 0.80
implementMethod · 0.80

Calls 3

proposeMethod · 0.80
countMethod · 0.80
validateMethod · 0.80

Tested by

no test coverage detected