See Implementable::implement, as PullLevel v2 can't implement Implementable directly.
(
&self,
nested: &mut Iterative<'b, S, u64>,
local_arrangements: &VariableMap<Iterative<'b, S, u64>>,
context: &mut I,
)
| 52 | /// See Implementable::implement, as PullLevel v2 can't implement |
| 53 | /// Implementable directly. |
| 54 | fn implement<'b, T, I, S>( |
| 55 | &self, |
| 56 | nested: &mut Iterative<'b, S, u64>, |
| 57 | local_arrangements: &VariableMap<Iterative<'b, S, u64>>, |
| 58 | context: &mut I, |
| 59 | ) -> ( |
| 60 | HashMap<PathId, Stream<S, (Vec<Value>, S::Timestamp, isize)>>, |
| 61 | ShutdownHandle, |
| 62 | ) |
| 63 | where |
| 64 | T: Timestamp + Lattice, |
| 65 | I: ImplContext<T>, |
| 66 | S: Scope<Timestamp = T>, |
| 67 | { |
| 68 | use differential_dataflow::operators::arrange::{Arrange, Arranged, TraceAgent}; |
| 69 | use differential_dataflow::operators::JoinCore; |
| 70 | use differential_dataflow::trace::implementations::ord::OrdValSpine; |
| 71 | use differential_dataflow::trace::TraceReader; |
| 72 | |
| 73 | assert_eq!(self.pull_attributes.is_empty(), false); |
| 74 | |
| 75 | let (input, mut shutdown_handle) = self.plan.implement(nested, local_arrangements, context); |
| 76 | |
| 77 | // Arrange input entities by eid. |
| 78 | let e_offset = input |
| 79 | .binds(self.pull_variable) |
| 80 | .expect("input relation doesn't bind pull_variable"); |
| 81 | |
| 82 | let paths = { |
| 83 | let (tuples, shutdown) = input.tuples(nested, context); |
| 84 | shutdown_handle.merge_with(shutdown); |
| 85 | tuples |
| 86 | }; |
| 87 | |
| 88 | let e_path: Arranged< |
| 89 | Iterative<S, u64>, |
| 90 | TraceAgent<OrdValSpine<Value, Vec<Value>, Product<T, u64>, isize>>, |
| 91 | > = paths.map(move |t| (t[e_offset].clone(), t)).arrange(); |
| 92 | |
| 93 | let mut shutdown_handle = shutdown_handle; |
| 94 | let path_streams = self |
| 95 | .pull_attributes |
| 96 | .iter() |
| 97 | .map(|a| { |
| 98 | let e_v = match context.forward_propose(a) { |
| 99 | None => panic!("attribute {:?} does not exist", a), |
| 100 | Some(propose_trace) => { |
| 101 | let frontier: Vec<T> = propose_trace.advance_frontier().to_vec(); |
| 102 | let (arranged, shutdown_propose) = |
| 103 | propose_trace.import_core(&nested.parent, a); |
| 104 | |
| 105 | let e_v = arranged.enter_at(nested, move |_, _, time| { |
| 106 | let mut forwarded = time.clone(); |
| 107 | forwarded.advance_by(&frontier); |
| 108 | Product::new(forwarded, 0) |
| 109 | }); |
| 110 | |
| 111 | shutdown_handle.add_button(shutdown_propose); |
nothing calls this directly
no test coverage detected