From fb57550a80db5d80dfb48bb63c6542a17a9e5625 Mon Sep 17 00:00:00 2001 From: Chase Wilson Date: Wed, 14 Apr 2021 15:38:35 -0500 Subject: [PATCH 1/6] Fixed antichain bug --- timely/src/progress/frontier.rs | 5 +++++ timely/src/progress/subgraph.rs | 34 ++++++++++++++++----------------- 2 files changed, 21 insertions(+), 18 deletions(-) diff --git a/timely/src/progress/frontier.rs b/timely/src/progress/frontier.rs index 4d370b6aa..d8e23d5d3 100644 --- a/timely/src/progress/frontier.rs +++ b/timely/src/progress/frontier.rs @@ -43,6 +43,11 @@ impl Antichain { } } + /// Reserves capacity for at least additional more elements to be inserted in the given `Antichain` + pub fn reserve(&mut self, additional: usize) { + self.elements.reserve(additional); + } + /// Performs a sequence of insertion and return true iff any insertion does. /// /// # Examples diff --git a/timely/src/progress/subgraph.rs b/timely/src/progress/subgraph.rs index b178532cb..d306b7a18 100644 --- a/timely/src/progress/subgraph.rs +++ b/timely/src/progress/subgraph.rs @@ -530,24 +530,22 @@ where assert_eq!(self.children[0].outputs, self.inputs()); assert_eq!(self.children[0].inputs, self.outputs()); - let mut internal_summary = Vec::with_capacity(self.inputs()); - for summary in self.scope_summary.iter() { - let scope_summary = summary.iter() - .map(|output| { - output.elements() - .iter() - .cloned() - .fold( - Antichain::with_capacity(output.elements().len()), - |mut antichain, path_summary| { - antichain.insert(TInner::summarize(path_summary)); - antichain - }, - ) - }) - .collect(); - - internal_summary.push(scope_summary); + let mut internal_summary: Vec> = (0..self.inputs()) + .map(|_| (0..self.outputs()).map(|_| Antichain::new()).collect()) + .collect(); + + for (summary, internal_summary) in + self.scope_summary.iter().zip(internal_summary.iter_mut()) + { + internal_summary.reserve(summary.len()); + + for (output, antichain) in summary.iter().zip(internal_summary) { + antichain.reserve(output.elements().len()); + + for path_summary in output.elements() { + antichain.insert(TInner::summarize(path_summary.clone())); + } + } } // Each child has expressed initial capabilities (their `shared_progress.internals`). From 347510c3f1b4432fb7bb246b3d710048d6f4bd23 Mon Sep 17 00:00:00 2001 From: Chase Wilson Date: Wed, 14 Apr 2021 15:40:49 -0500 Subject: [PATCH 2/6] Added assertions after summary construction --- timely/src/progress/subgraph.rs | 10 ++++++++++ 1 file changed, 10 insertions(+) diff --git a/timely/src/progress/subgraph.rs b/timely/src/progress/subgraph.rs index d306b7a18..c8cd48230 100644 --- a/timely/src/progress/subgraph.rs +++ b/timely/src/progress/subgraph.rs @@ -548,6 +548,16 @@ where } } + debug_assert_eq!( + internal_summary.len(), + self.inputs(), + "the internal summary should have as many elements as there are inputs", + ); + debug_assert!( + internal_summary.iter().all(|summary| summary.len() == self.outputs()), + "each element of the internal summary should have as many elements as there are outputs", + ); + // Each child has expressed initial capabilities (their `shared_progress.internals`). // We introduce these into the progress tracker to determine the scope's initial // internal capabilities. From 25514f77ecb15a4245c507a3561f01e31e7fe501 Mon Sep 17 00:00:00 2001 From: Chase Wilson Date: Wed, 14 Apr 2021 15:45:11 -0500 Subject: [PATCH 3/6] Added comment about summary length --- timely/src/progress/subgraph.rs | 3 +++ 1 file changed, 3 insertions(+) diff --git a/timely/src/progress/subgraph.rs b/timely/src/progress/subgraph.rs index c8cd48230..909b3a093 100644 --- a/timely/src/progress/subgraph.rs +++ b/timely/src/progress/subgraph.rs @@ -530,6 +530,9 @@ where assert_eq!(self.children[0].outputs, self.inputs()); assert_eq!(self.children[0].inputs, self.outputs()); + // Note that we need to have `self.inputs()` elements in the summary + // with each element containing `self.outputs()` antichains regardless + // of how long `self.scope_summary` is let mut internal_summary: Vec> = (0..self.inputs()) .map(|_| (0..self.outputs()).map(|_| Antichain::new()).collect()) .collect(); From 7b792baf005a50fe7e0215f6cf4892450b66056f Mon Sep 17 00:00:00 2001 From: Chase Wilson Date: Wed, 14 Apr 2021 18:54:41 -0500 Subject: [PATCH 4/6] Simplified summary loop --- timely/src/progress/subgraph.rs | 14 ++++---------- 1 file changed, 4 insertions(+), 10 deletions(-) diff --git a/timely/src/progress/subgraph.rs b/timely/src/progress/subgraph.rs index 909b3a093..50d7a1cb4 100644 --- a/timely/src/progress/subgraph.rs +++ b/timely/src/progress/subgraph.rs @@ -537,17 +537,11 @@ where .map(|_| (0..self.outputs()).map(|_| Antichain::new()).collect()) .collect(); - for (summary, internal_summary) in - self.scope_summary.iter().zip(internal_summary.iter_mut()) - { - internal_summary.reserve(summary.len()); - - for (output, antichain) in summary.iter().zip(internal_summary) { + for (input_idx, input) in self.scope_summary.iter().enumerate() { + for (output_idx, output) in input.iter().enumerate() { + let antichain = &mut internal_summary[input_idx][output_idx]; antichain.reserve(output.elements().len()); - - for path_summary in output.elements() { - antichain.insert(TInner::summarize(path_summary.clone())); - } + antichain.extend(output.elements().iter().cloned().map(TInner::summarize)); } } From c9b3778705ee5be4cf558749d6516274520c8bcc Mon Sep 17 00:00:00 2001 From: Chase Wilson Date: Wed, 14 Apr 2021 18:57:11 -0500 Subject: [PATCH 5/6] Fancy asserts --- timely/src/progress/subgraph.rs | 13 +++++++++++-- 1 file changed, 11 insertions(+), 2 deletions(-) diff --git a/timely/src/progress/subgraph.rs b/timely/src/progress/subgraph.rs index 50d7a1cb4..133fe2807 100644 --- a/timely/src/progress/subgraph.rs +++ b/timely/src/progress/subgraph.rs @@ -637,8 +637,17 @@ impl PerOperatorState { let (internal_summary, shared_progress) = scope.get_internal_summary(); - assert_eq!(internal_summary.len(), inputs); - assert!(!internal_summary.iter().any(|x| x.len() != outputs)); + assert_eq!( + internal_summary.len(), + inputs, + "operator summary has {} inputs when {} were expected", + internal_summary.len(), + inputs, + ); + assert!( + !internal_summary.iter().any(|x| x.len() != outputs), + "operator summary had too few outputs", + ); PerOperatorState { name: scope.name().to_owned(), From f148a20027006cf39748724b134cda301fe15930 Mon Sep 17 00:00:00 2001 From: Chase Wilson Date: Wed, 14 Apr 2021 18:58:05 -0500 Subject: [PATCH 6/6] Simplified summary init --- timely/src/progress/subgraph.rs | 5 +---- 1 file changed, 1 insertion(+), 4 deletions(-) diff --git a/timely/src/progress/subgraph.rs b/timely/src/progress/subgraph.rs index 133fe2807..309d9077c 100644 --- a/timely/src/progress/subgraph.rs +++ b/timely/src/progress/subgraph.rs @@ -533,10 +533,7 @@ where // Note that we need to have `self.inputs()` elements in the summary // with each element containing `self.outputs()` antichains regardless // of how long `self.scope_summary` is - let mut internal_summary: Vec> = (0..self.inputs()) - .map(|_| (0..self.outputs()).map(|_| Antichain::new()).collect()) - .collect(); - + let mut internal_summary = vec![vec![Antichain::new(); self.outputs()]; self.inputs()]; for (input_idx, input) in self.scope_summary.iter().enumerate() { for (output_idx, output) in input.iter().enumerate() { let antichain = &mut internal_summary[input_idx][output_idx];