rudb-plan 0.4.28

The logical and physical plan representations, their textual forms and Substrait conversion.
Documentation
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
//! The shape a plan runs as: which operator each node becomes, which pipeline it runs in, and which
//! pipeline waits for which.
//!
//! A pipeline is a run of operators from a source to a sink, and a plan breaks into several of them
//! wherever an operator has to see all of its input before it produces anything. A sort is the
//! plain case: everything under it is one pipeline that ends in the sort, and what reads the sorted
//! rows back is the next one, which cannot start until the first has finished. A join is two below
//! the one above it, because the side that is gathered has to be complete before the side that
//! probes it can run a single row.
//!
//! The operators are numbered by the same walk, because the two answers are the same answer. An
//! operator's id is what a metrics document calls it, what `EXPLAIN` prints beside it and what the
//! builder tags its counters with, and a number that three pieces of code work out separately is a
//! number that three pieces of code can disagree about.
//!
//! # Why this is here
//!
//! Two crates need all of it and neither can see the other. `rudb-exec` builds the tree, and
//! `rudb-opt` prints what `EXPLAIN` shows without building anything. Written twice it would be
//! right twice on the day it was written and wrong once some time after that, and the version that
//! would be wrong is the printed one, which is the version somebody reads when they are trying to
//! understand why a query is slow.
//!
//! It is physical knowledge about a logical tree, which is worth saying out loud. Whether an
//! operator is a pipeline breaker is a fact about how it is executed rather than about what it
//! means, and the reason it can live here anyway is that at this milestone the physical plan is the
//! logical plan with different words on it, which `crates/rudb-exec/src/build.rs` says at the top.
//! The day there is a physical plan this moves onto it and every caller keeps its call.
//!
//! # The rule for the pipelines
//!
//! The root of the plan produces into pipeline 0. Walking down from there, a node inherits the
//! pipeline of its parent, except that
//!
//! - an aggregate, a window, a sort, a top n and a distinct are sinks, so the node and everything
//!   under it are a new pipeline that the parent's waits for,
//! - a join and a set operation are two, the side that is gathered first and the side that reads
//!   it, with the second waiting for the first and the parent's waiting for the second, and which
//!   side a join gathers is the node's own build side rather than always the right one,
//! - a cross product keeps its left side and itself in the parent's pipeline, because the product
//!   is produced a chunk at a time and never held, and puts its right side in a new one, because
//!   that side is kept whole to be replayed,
//! - a materialised `WITH` is the sink of a new pipeline that its definition fills, and its body
//!   stays in the parent's, with the pipeline holding each read of the name waiting for the one
//!   that fills it.
//!
//! There is no scheduler reading any of this yet. It is written down because it is known, and an
//! edge reconstructed later from a tree somebody has already flattened is an edge somebody has to
//! guess at.
//!
//! # The rule for the numbers
//!
//! A node takes the next id when the walk reaches it, so the root is operator 0 and a parent is
//! always numbered before everything under it. A node with two inputs takes a second id straight
//! after its own, for the operator that holds the side that has to finish first: the gather under a
//! join or a set operation, and the kept chunks under a cross product. Those are operators in their
//! own right, they have their own counters and their own row in a metrics document, and they exist
//! because the plan has two inputs there rather than because somebody chose to add one.
//!
//! Then the children, and for a node with two inputs the side that runs first is walked first, so
//! the ids go in the order the work happens rather than in the order the tree prints.

use crate::node::{BuildSide, Node};
use crate::plan::Plan;
use crate::{NodeRef, OperatorRef, PipelineRef};

/// What a plan runs as.
#[derive(Debug, Clone)]
pub struct Shape {
    /// Per node in the arena, the operator it becomes and the pipeline that runs it, or none for a
    /// node the root does not reach.
    of: Vec<Option<Placed>>,
    /// What each pipeline waits for, indexed by pipeline.
    waits: Vec<Vec<PipelineRef>>,
    /// How many operators there are.
    operators: OperatorRef,
    /// Per operator, the operator its rows go to, or none for the one that produces the answer.
    ///
    /// Indexed by operator id, which is why it is filled where the ids are handed out rather than
    /// by a second walk. See [`Shape::consumer`] for what a reader does with it.
    consumes: Vec<Option<OperatorRef>>,
    /// The pipeline each materialisation currently being walked over is filled by.
    ///
    /// A stack rather than a map, because a materialised `WITH` inside another one is a `WITH`
    /// inside the body of the first, and the inner name is the one a scan under it reads.
    holding: Vec<(u32, PipelineRef)>,
}

/// One node's place in the shape.
#[derive(Debug, Clone, Copy)]
struct Placed {
    operator: OperatorRef,
    /// The operator that holds the side which has to finish first, for a node with two inputs.
    gathered: Option<OperatorRef>,
    pipeline: PipelineRef,
}

impl Shape {
    /// Works out the shape of a plan.
    #[must_use]
    pub fn of(plan: &Plan) -> Self {
        let mut shape = Self {
            of: vec![None; plan.node_count()],
            waits: vec![Vec::new()],
            operators: 0,
            consumes: Vec::new(),
            holding: Vec::new(),
        };
        shape.walk(plan, plan.root(), ROOT, None);
        shape
    }

    /// How many pipelines there are, which is at least one.
    #[must_use]
    pub fn pipelines(&self) -> usize {
        self.waits.len()
    }

    /// How many operators the tree has, which is at least one and is more than the plan has nodes
    /// whenever the plan has a node with two inputs in it.
    #[must_use]
    pub fn operators(&self) -> OperatorRef {
        self.operators
    }

    /// The operator this node becomes.
    ///
    /// # Panics
    ///
    /// If the node is not reachable from the plan's root, which is a node the arena is still
    /// holding after a rewrite replaced it.
    #[must_use]
    pub fn operator(&self, node: NodeRef) -> OperatorRef {
        self.placed(node).operator
    }

    /// The operator this node becomes, or none for a node the root does not reach.
    ///
    /// The tolerant form of [`Shape::operator`], for a caller walking the whole arena rather than
    /// the tree, which is what somebody filling one fact in per operator ends up doing.
    #[must_use]
    pub fn operator_of(&self, node: NodeRef) -> Option<OperatorRef> {
        self.of.get(node as usize).copied().flatten().map(|placed| placed.operator)
    }

    /// The operator holding the side of this node that has to finish first, if it has two inputs.
    ///
    /// # Panics
    ///
    /// The same as [`Shape::operator`].
    #[must_use]
    pub fn gathered(&self, node: NodeRef) -> Option<OperatorRef> {
        self.placed(node).gathered
    }

    /// The operator this one's rows go into, or none for the operator that produces the answer.
    ///
    /// Keyed by operator rather than by node because the operator tree is not quite the plan tree:
    /// a node with two inputs is two operators, and the side that is gathered feeds the one that
    /// holds it rather than the join above it. A reader of a metrics document that wants to check
    /// an operator's input against what its children produced needs the edges the way the rows
    /// actually moved, which is this.
    ///
    /// # Panics
    ///
    /// If there is no such operator.
    #[must_use]
    pub fn consumer(&self, operator: OperatorRef) -> Option<OperatorRef> {
        self.consumes[operator as usize]
    }

    /// The pipeline this node runs in.
    ///
    /// For a sink that is the pipeline it ends rather than the one above it, so a sort is in the
    /// pipeline that feeds it and the operator that reads the sorted rows is in the one above.
    ///
    /// # Panics
    ///
    /// The same as [`Shape::operator`].
    #[must_use]
    pub fn pipeline(&self, node: NodeRef) -> PipelineRef {
        self.placed(node).pipeline
    }

    /// What this pipeline has to wait for, in ascending order.
    ///
    /// # Panics
    ///
    /// If there is no such pipeline.
    #[must_use]
    pub fn waits_for(&self, pipeline: PipelineRef) -> &[PipelineRef] {
        &self.waits[pipeline as usize]
    }

    /// Every pipeline, from the root's outwards.
    pub fn all(&self) -> impl Iterator<Item = PipelineRef> {
        0..u32::try_from(self.waits.len()).unwrap_or(u32::MAX)
    }

    /// Where a node ended up.
    ///
    /// # Panics
    ///
    /// If the node is not reachable from the plan's root.
    fn placed(&self, node: NodeRef) -> Placed {
        self.of[node as usize].expect("a node under the root of the plan it was walked from")
    }

    /// A new pipeline that nothing waits for yet.
    fn fresh(&mut self) -> PipelineRef {
        self.waits.push(Vec::new());
        u32::try_from(self.waits.len() - 1).unwrap_or(u32::MAX)
    }

    /// Records that `pipeline` cannot start until `on` has finished.
    fn waits_on(&mut self, pipeline: PipelineRef, on: PipelineRef) {
        self.waits[pipeline as usize].push(on);
    }

    /// The next operator id, and the operator its rows go to.
    fn number(&mut self, into: Option<OperatorRef>) -> OperatorRef {
        let id = self.operators;
        self.operators += 1;
        self.consumes.push(into);
        id
    }

    /// A node with two inputs, where `held` has to be finished before `driving` can start a row.
    fn two(
        &mut self,
        plan: &Plan,
        node: NodeRef,
        operator: OperatorRef,
        pipeline: PipelineRef,
        held: NodeRef,
        driving: NodeRef,
    ) {
        let gathered = self.number(Some(operator));
        let first = self.fresh();
        let second = self.fresh();
        self.waits_on(second, first);
        self.waits_on(pipeline, second);
        self.of[node as usize] =
            Some(Placed { operator, gathered: Some(gathered), pipeline: second });
        // The gathered side's rows go into the operator that holds them rather than straight into
        // the join, which is the one place the operator tree has a shape the plan does not.
        self.walk(plan, held, first, Some(gathered));
        self.walk(plan, driving, second, Some(operator));
    }

    fn walk(
        &mut self,
        plan: &Plan,
        node: NodeRef,
        pipeline: PipelineRef,
        into: Option<OperatorRef>,
    ) {
        let operator = self.number(into);
        match *plan.node(node) {
            Node::Aggregate { input, .. }
            | Node::Window { input, .. }
            | Node::Sort { input, .. }
            | Node::TopN { input, .. }
            // A share of the input is not known until the input has ended, so this holds its rows
            // and is a sink, where a plain limit hands each chunk on and stays in the pipeline.
            | Node::LimitPercent { input, .. }
            | Node::Distinct { input, .. } => {
                let below = self.fresh();
                self.waits_on(pipeline, below);
                self.of[node as usize] = Some(Placed { operator, gathered: None, pipeline: below });
                self.walk(plan, input, below, Some(operator));
            }
            // Which side a join holds is the plan's own flag, written by `rudb_opt`'s `sides` pass
            // from an estimate of how many rows each input produces, and the builder reads that
            // same flag. Assuming the right side here would put the pipelines and the edges on the
            // wrong sides of every join the pass turned around, and the edges are the thing a
            // reader checks an operator's input against.
            Node::Join { left, right, build, .. } => {
                let (held, driving) =
                    if build == BuildSide::Left { (left, right) } else { (right, left) };
                self.two(plan, node, operator, pipeline, held, driving);
            }
            Node::DependentJoin { left, right, .. } | Node::SetOp { left, right, .. } => {
                self.two(plan, node, operator, pipeline, right, left);
            }
            // The definition is held whole and the body reads it, so the node is the sink of the
            // pipeline that fills it, the same way a sort is the sink of the pipeline under it. The
            // body stays where the parent is, because it produces into the parent a chunk at a time
            // and is never held.
            //
            // The edge is recorded at every scan rather than here, because that is where the
            // waiting actually is: a scan under a sort in the body is in a pipeline of its own, and
            // it is that pipeline which cannot start until the rows exist. A body that reads the
            // name nowhere waits for nothing, which is the shape of a materialisation the optimizer
            // is about to drop.
            Node::MaterializedCte { definition, body, cte, .. } => {
                // Its rows go nowhere. It is filled by its definition and read back by the scans
                // of its name, so nothing above it ever sees a row of it, and naming the node it
                // hangs under would claim rows that the body produced and this never handed on.
                // That makes it a second operator with nothing above it, which is what it is.
                self.consumes[operator as usize] = None;
                let filling = self.fresh();
                self.of[node as usize] =
                    Some(Placed { operator, gathered: None, pipeline: filling });
                self.walk(plan, definition, filling, Some(operator));
                self.holding.push((cte, filling));
                // The body produces the answer a chunk at a time and the node it hangs under never
                // sees those rows, so what consumes them is whatever consumes this node.
                self.walk(plan, body, pipeline, into);
                self.holding.pop();
            }
            Node::CteScan { cte, .. } => {
                self.of[node as usize] = Some(Placed { operator, gathered: None, pipeline });
                if let Some(&(_, filled)) =
                    self.holding.iter().rev().find(|&&(held, _)| held == cte)
                {
                    if !self.waits[pipeline as usize].contains(&filled) {
                        self.waits_on(pipeline, filled);
                    }
                }
            }
            Node::CrossProduct { left, right } => {
                let gathered = self.number(Some(operator));
                let aside = self.fresh();
                self.waits_on(pipeline, aside);
                self.of[node as usize] =
                    Some(Placed { operator, gathered: Some(gathered), pipeline });
                self.walk(plan, right, aside, Some(gathered));
                self.walk(plan, left, pipeline, Some(operator));
            }
            ref other => {
                self.of[node as usize] = Some(Placed { operator, gathered: None, pipeline });
                for child in other.children().into_iter().flatten() {
                    self.walk(plan, child, pipeline, Some(operator));
                }
            }
        }
    }
}

/// The pipeline the root of a plan produces into.
///
/// Public because it is the one pipeline nothing drains. Every other pipeline ends in a sink and is
/// run by the loop that fills that sink, and this one is pulled from by whoever wanted the answer,
/// so whoever that is has to know which pipeline the loop they are writing belongs to.
pub const ROOT: PipelineRef = 0;

#[cfg(test)]
mod tests {
    use super::Shape;
    use crate::plan::Plan;

    fn shaped(text: &str) -> (Plan, Shape) {
        let plan =
            Plan::parse(text).unwrap_or_else(|error| panic!("{text} did not parse: {error}"));
        let shape = Shape::of(&plan);
        (plan, shape)
    }

    #[test]
    fn a_plan_with_nothing_that_buffers_is_one_pipeline() {
        let (plan, shape) = shaped(concat!(
            "Filter (#0.0::INTEGER > 1::INTEGER)::BOOLEAN\n",
            "  Get memory.main.t AS t #0 [a::INTEGER]\n",
        ));
        assert_eq!(shape.pipelines(), 1);
        assert_eq!(shape.pipeline(plan.root()), 0);
        assert!(shape.waits_for(0).is_empty());
    }

    #[test]
    fn a_sort_ends_the_pipeline_below_it_and_the_one_above_waits() {
        let (plan, shape) = shaped(concat!(
            "Sort [#0.0::INTEGER ASC NULLS LAST]\n",
            "  Get memory.main.t AS t #0 [a::INTEGER]\n",
        ));
        assert_eq!(shape.pipelines(), 2);
        assert_eq!(shape.pipeline(plan.root()), 1, "the sort is the sink of the one below");
        assert_eq!(shape.waits_for(0), [1]);
        assert!(shape.waits_for(1).is_empty());
    }

    #[test]
    fn a_join_is_two_pipelines_in_the_order_they_have_to_run() {
        let (plan, shape) = shaped(concat!(
            "Join INNER on=[(#0.0::INTEGER = #1.0::INTEGER)::BOOLEAN]\n",
            "  Get memory.main.l AS l #0 [a::INTEGER]\n",
            "  Get memory.main.r AS r #1 [a::INTEGER]\n",
        ));
        let [left, right] = plan.node(plan.root()).children();
        assert_eq!(shape.pipelines(), 3);
        assert_eq!(shape.pipeline(right.unwrap()), 1, "the gathered side runs first");
        assert_eq!(shape.pipeline(left.unwrap()), 2, "the probing side is the second");
        assert_eq!(shape.pipeline(plan.root()), 2, "and the join is its sink");
        assert_eq!(shape.waits_for(2), [1]);
        assert_eq!(shape.waits_for(0), [2]);
    }

    #[test]
    fn a_join_told_to_build_on_its_left_runs_that_side_first() {
        let (plan, shape) = shaped(concat!(
            "Join INNER on=[(#0.0::INTEGER = #1.0::INTEGER)::BOOLEAN] build=left\n",
            "  Get memory.main.l AS l #0 [a::INTEGER]\n",
            "  Get memory.main.r AS r #1 [a::INTEGER]\n",
        ));
        let [left, right] = plan.node(plan.root()).children();
        let (left, right) = (left.unwrap(), right.unwrap());
        assert_eq!(shape.pipeline(left), 1, "the side the flag names is the gathered one");
        assert_eq!(shape.pipeline(right), 2, "so the right side is the one that probes");
        assert_eq!(shape.pipeline(plan.root()), 2);
        let gathered = shape.gathered(plan.root()).expect("a join holds a side");
        assert_eq!(shape.consumer(shape.operator(left)), Some(gathered));
        assert_eq!(shape.consumer(shape.operator(right)), Some(shape.operator(plan.root())));
    }

    #[test]
    fn a_cross_product_keeps_its_left_side_where_it_was() {
        let (plan, shape) = shaped(concat!(
            "CrossProduct\n",
            "  Get memory.main.l AS l #0 [a::INTEGER]\n",
            "  Get memory.main.r AS r #1 [a::INTEGER]\n",
        ));
        let [left, right] = plan.node(plan.root()).children();
        assert_eq!(shape.pipelines(), 2);
        assert_eq!(shape.pipeline(plan.root()), 0, "the product streams");
        assert_eq!(shape.pipeline(left.unwrap()), 0, "and so does the side it streams");
        assert_eq!(shape.pipeline(right.unwrap()), 1, "the side that is kept is its own");
        assert_eq!(shape.waits_for(0), [1]);
    }

    #[test]
    fn a_materialisation_is_the_sink_of_the_pipeline_that_fills_it() {
        let (plan, shape) = shaped(concat!(
            "MaterializedCte c @0 [a::INTEGER]\n",
            "  Get memory.main.t AS t #0 [a::INTEGER]\n",
            "  Project #2 [#1.0::INTEGER AS a]\n",
            "    CteScan c @0 #1 [a::INTEGER]\n",
        ));
        let [definition, body] = plan.node(plan.root()).children();
        assert_eq!(shape.pipelines(), 2);
        assert_eq!(shape.pipeline(plan.root()), 1, "the node is what the definition fills");
        assert_eq!(shape.pipeline(definition.unwrap()), 1, "and the definition is under it");
        assert_eq!(shape.pipeline(body.unwrap()), 0, "the body produces the answer");
        assert_eq!(shape.waits_for(0), [1], "and cannot start before the rows are there");
    }

    #[test]
    fn the_pipeline_that_waits_is_the_one_the_read_is_in() {
        let (plan, shape) = shaped(concat!(
            "MaterializedCte c @0 [a::INTEGER]\n",
            "  Get memory.main.t AS t #0 [a::INTEGER]\n",
            "  Sort [#1.0::INTEGER ASC NULLS LAST]\n",
            "    CteScan c @0 #1 [a::INTEGER]\n",
        ));
        assert_eq!(shape.pipelines(), 3);
        assert_eq!(shape.waits_for(0), [2], "the answer waits for the sort");
        assert_eq!(shape.waits_for(2), [1], "and the sort waits for the rows it reads");
        assert!(shape.waits_for(1).is_empty(), "which wait for nothing");
        let [_, body] = plan.node(plan.root()).children();
        assert_eq!(shape.pipeline(body.unwrap()), 2);
    }

    /// A body that never names the materialisation waits for nothing, which is the shape the pass
    /// that drops an unread one is about to remove.
    #[test]
    fn a_body_that_reads_nothing_waits_for_nothing() {
        let (_, shape) = shaped(concat!(
            "MaterializedCte c @0 [a::INTEGER]\n",
            "  Get memory.main.t AS t #0 [a::INTEGER]\n",
            "  Project #2 [1::INTEGER AS one]\n",
            "    Dummy\n",
        ));
        assert_eq!(shape.pipelines(), 2);
        assert!(shape.waits_for(0).is_empty());
    }

    #[test]
    fn two_sorts_under_one_another_are_three_pipelines_in_a_line() {
        let (plan, shape) = shaped(concat!(
            "Sort [#0.0::INTEGER ASC NULLS LAST]\n",
            "  Limit 10 offset 0\n",
            "    Sort [#0.0::INTEGER DESC NULLS FIRST]\n",
            "      Get memory.main.t AS t #0 [a::INTEGER]\n",
        ));
        assert_eq!(shape.pipelines(), 3);
        assert_eq!(shape.pipeline(plan.root()), 1);
        assert_eq!(shape.waits_for(0), [1]);
        assert_eq!(shape.waits_for(1), [2]);
        assert!(shape.waits_for(2).is_empty());
    }

    #[test]
    fn a_parent_is_numbered_before_everything_under_it() {
        let (plan, shape) = shaped(concat!(
            "Sort [#0.0::INTEGER ASC NULLS LAST]\n",
            "  Filter (#0.0::INTEGER > 1::INTEGER)::BOOLEAN\n",
            "    Get memory.main.t AS t #0 [a::INTEGER]\n",
        ));
        let filter = plan.node(plan.root()).children()[0].unwrap();
        let get = plan.node(filter).children()[0].unwrap();
        assert_eq!(shape.operator(plan.root()), 0);
        assert_eq!(shape.operator(filter), 1);
        assert_eq!(shape.operator(get), 2);
        assert_eq!(shape.operators(), 3);
        assert_eq!(shape.gathered(plan.root()), None, "one input, nothing to hold");
    }

    #[test]
    fn a_node_with_two_inputs_is_two_operators_and_the_first_side_is_numbered_first() {
        let (plan, shape) = shaped(concat!(
            "Join INNER on=[(#0.0::INTEGER = #1.0::INTEGER)::BOOLEAN]\n",
            "  Get memory.main.l AS l #0 [a::INTEGER]\n",
            "  Get memory.main.r AS r #1 [a::INTEGER]\n",
        ));
        let [left, right] = plan.node(plan.root()).children();
        assert_eq!(shape.operator(plan.root()), 0);
        assert_eq!(shape.gathered(plan.root()), Some(1), "the gather is an operator of its own");
        assert_eq!(shape.operator(right.unwrap()), 2, "the side that has to finish first");
        assert_eq!(shape.operator(left.unwrap()), 3);
        assert_eq!(shape.operators(), 4);
    }

    #[test]
    fn every_operator_but_the_one_that_answers_names_what_its_rows_go_into() {
        let (_, shape) = shaped(concat!(
            "Sort [#0.0::INTEGER ASC NULLS LAST]\n",
            "  Filter (#0.0::INTEGER > 1::INTEGER)::BOOLEAN\n",
            "    Get memory.main.t AS t #0 [a::INTEGER]\n",
        ));
        assert_eq!(shape.consumer(0), None, "the sort is what the answer is read from");
        assert_eq!(shape.consumer(1), Some(0));
        assert_eq!(shape.consumer(2), Some(1));
    }

    /// The one place the operator tree is a different shape than the plan. The gathered side feeds
    /// the operator that holds it, and that one feeds the join, so a check that read the plan tree
    /// instead would compare the join's input against rows it never saw.
    #[test]
    fn the_gathered_side_feeds_the_operator_that_holds_it_rather_than_the_join() {
        let (_, shape) = shaped(concat!(
            "Join INNER on=[(#0.0::INTEGER = #1.0::INTEGER)::BOOLEAN]\n",
            "  Get memory.main.l AS l #0 [a::INTEGER]\n",
            "  Get memory.main.r AS r #1 [a::INTEGER]\n",
        ));
        assert_eq!(shape.consumer(0), None, "the join produces the answer");
        assert_eq!(shape.consumer(1), Some(0), "the gather hands the held side to the join");
        assert_eq!(shape.consumer(2), Some(1), "the side that finishes first is what it holds");
        assert_eq!(shape.consumer(3), Some(0), "the driving side goes straight into the join");
    }

    /// A materialisation is filled by its definition and read back by the scans of its name, so no
    /// row of it is ever handed upwards and it is an operator with nothing above it. Its body is
    /// what streams into whatever the whole thing hangs under.
    #[test]
    fn a_materialisation_is_fed_by_its_definition_and_its_body_streams_past_it() {
        let (_, shape) = shaped(concat!(
            "Project #2 [#3.0::INTEGER AS a]\n",
            "  MaterializedCte c @0 [a::INTEGER]\n",
            "    Get memory.main.t AS t #0 [a::INTEGER]\n",
            "    CteScan c @0 #1 [a::INTEGER]\n",
        ));
        assert_eq!(shape.consumer(0), None);
        assert_eq!(shape.consumer(1), None, "the materialisation hands nothing upwards");
        assert_eq!(shape.consumer(2), Some(1), "the definition is what fills it");
        assert_eq!(shape.consumer(3), Some(0), "and the body produces into the project");
    }
}