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

Method implement

src/plan/pull_v2.rs:54–140  ·  view source on GitHub ↗

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,
    )

Source from the content-addressed store, hash-verified

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);

Callers

nothing calls this directly

Calls 5

tuplesMethod · 0.80
merge_withMethod · 0.80
forward_proposeMethod · 0.80
add_buttonMethod · 0.80
bindsMethod · 0.45

Tested by

no test coverage detected