Skip to main content

mz_timely_util/funded_spine/
spine_fueled.rs

1// Copyright (c) 2015 Frank McSherry
2// Copyright Materialize, Inc. and contributors. All rights reserved.
3//
4// Use of this software is governed by the Business Source License
5// included in the LICENSE file.
6//
7// As of the Change Date specified in that file, in accordance with
8// the Business Source License, use of this software will be governed
9// by the Apache License, Version 2.0.
10//
11// Portions of this file are derived from the fueled spine in
12// differential-dataflow 0.25.1. The original source code was retrieved from:
13//
14//     https://github.com/TimelyDataflow/differential-dataflow/blob/6e53ac674feacac83b0c3c5d123bdd5db488ef9d/differential-dataflow/src/trace/implementations/spine_fueled.rs
15//
16// The original source code is subject to the terms of the MIT license, a copy
17// of which can be found in the LICENSE file at the root of this repository.
18
19//! An append-only collection of update batches.
20//!
21//! The `Spine` is a general-purpose trace implementation based on collection and merging
22//! immutable batches of updates. It is generic with respect to the batch type, and can be
23//! instantiated for any implementor of `trace::Batch`.
24//!
25//! ## Design
26//!
27//! This spine is represented as a list of layers, where each element in the list is either
28//!
29//!   1. MergeState::Vacant  empty
30//!   2. MergeState::Single  a single batch
31//!   3. MergeState::Double  a pair of batches
32//!
33//! Each "batch" has the option to be `None`, indicating a non-batch that nonetheless acts
34//! as a number of updates proportionate to the level at which it exists (for bookkeeping).
35//!
36//! Each of the batches at layer i contains at most 2^i elements. The sequence of batches
37//! should have the upper bound of one match the lower bound of the next. Batches may be
38//! logically empty, with matching upper and lower bounds, as a bookkeeping mechanism.
39//!
40//! Each batch at layer i is treated as if it contains exactly 2^i elements, even though it
41//! may actually contain fewer elements. This allows us to decouple the physical representation
42//! from logical amounts of effort invested in each batch. It allows us to begin compaction and
43//! to reduce the number of updates, without compromising our ability to continue to move
44//! updates along the spine. We are explicitly making the trade-off that while some batches
45//! might compact at lower levels, we want to treat them as if they contained their full set of
46//! updates for accounting reasons (to apply work to higher levels).
47//!
48//! We maintain the invariant that for any in-progress merge at level k there should be fewer
49//! than 2^k records at levels lower than k. That is, even if we were to apply an unbounded
50//! amount of effort to those records, we would not have enough records to prompt a merge into
51//! the in-progress merge. Ideally, we maintain the extended invariant that for any in-progress
52//! merge at level k, the remaining effort required (number of records minus applied effort) is
53//! less than the number of records that would need to be added to reach 2^k records in layers
54//! below.
55//!
56//! ## Mathematics
57//!
58//! When a merge is initiated, there should be a non-negative *deficit* of updates before the layers
59//! below could plausibly produce a new batch for the currently merging layer. We must determine a
60//! factor of proportionality, so that newly arrived updates provide at least that amount of "fuel"
61//! towards the merging layer, so that the merge completes before lower levels invade.
62//!
63//! ### Deficit:
64//!
65//! A new merge is initiated only in response to the completion of a prior merge, or the introduction
66//! of new records from outside. The latter case is special, and will maintain our invariant trivially,
67//! so we will focus on the former case.
68//!
69//! When a merge at level k completes, assuming we have maintained our invariant then there should be
70//! fewer than 2^k records at lower levels. The newly created merge at level k+1 will require up to
71//! 2^k+2 units of work, and should not expect a new batch until strictly more than 2^k records are
72//! added. This means that a factor of proportionality of four should be sufficient to ensure that
73//! the merge completes before a new merge is initiated.
74//!
75//! When new records get introduced, we will need to roll up any batches at lower levels, which we
76//! treat as the introduction of records. Each of these virtual records introduced should either be
77//! accounted for the fuel it should contribute, as it results in the promotion of batches closer to
78//! in-progress merges.
79//!
80//! ### Fuel sharing
81//!
82//! We like the idea of applying fuel preferentially to merges at *lower* levels, under the idea that
83//! they are easier to complete, and we benefit from fewer total merges in progress. This does delay
84//! the completion of merges at higher levels, and may not obviously be a total win. If we choose to
85//! do this, we should make sure that we correctly account for completed merges at low layers: they
86//! should still extract fuel from new updates even though they have completed, at least until they
87//! have paid back any "debt" to higher layers by continuing to provide fuel as updates arrive.
88
89
90use differential_dataflow::logging::Logger;
91use differential_dataflow::trace::{Batch, ExertionLogic, Merger, Trace, TraceReader};
92
93use ::timely::dataflow::operators::generic::OperatorInfo;
94use ::timely::progress::{Antichain, frontier::AntichainRef};
95use ::timely::order::PartialOrder;
96
97/// An append-only collection of update tuples.
98///
99/// A spine maintains a small number of immutable collections of update tuples, merging the collections when
100/// two have similar sizes. In this way, it allows the addition of more tuples, which may then be merged with
101/// other immutable collections.
102pub struct Spine<B: Batch> {
103    operator: OperatorInfo,
104    logger: Option<Logger>,
105    logical_frontier: Antichain<B::Time>,   // Times after which the trace must accumulate correctly.
106    physical_frontier: Antichain<B::Time>,  // Times after which the trace must be able to subset its inputs.
107    merging: Vec<MergeState<B>>,            // Several possibly shared collections of updates.
108    pending: Vec<B>,                        // Batches at times in advance of `frontier`.
109    upper: Antichain<B::Time>,
110    effort: usize,
111    activator: Option<timely::scheduling::activate::Activator>,
112    /// Parameters to `exert_logic`, containing tuples of `(index, count, length)`.
113    exert_logic_param: Vec<(usize, usize, usize)>,
114    /// Logic to indicate whether and how many records we should introduce in the absence of actual updates.
115    exert_logic: Option<ExertionLogic>,
116    /// Fuel funded by inserted updates and not yet spent. It does not decay:
117    /// it accumulates in proportion to total input, and later turns spend it
118    /// down, so the total stays proportional to input.
119    consolidation_credit: usize,
120    /// Policy allowances funded by inserted batches and not yet spent.
121    progress_grants: usize,
122}
123
124/// Fuel each inserted update credits for policy-requested effort, scaled by
125/// the spine's effort multiplier, which is 1 unless set with `with_effort`.
126/// A granted request spends the effort the policy returns, so credit funds
127/// one grant per `effort / 8` inserted updates, about 125 at the cluster
128/// policy's effort of 1000. Small batches fund no grant from credit, and the
129/// banked allowances per batch are what bound them. The multiple matches the
130/// least `introduce_batch` already funds per update for active merges,
131/// which is `8 << level`.
132///
133/// NOTE: the credit bounds how many requests are granted, not the work each
134/// grant causes. A grant with no merge in progress introduces a virtual
135/// batch, which can start a merge up into the largest layer, and that merge
136/// then finishes unfunded. See the module documentation of `funded_spine`.
137const CONSOLIDATION_CREDIT_PER_UPDATE: usize = 8;
138
139/// Policy allowances each inserted batch funds regardless of its size.
140const CONSOLIDATION_GRANTS_PER_PROGRESS: usize = 8;
141
142/// Most progress-funded allowances held, eight batches' worth. The bank fills to
143/// this cap while the trace is reduced, so a long quiet spell funds at most this
144/// burst of consolidation when input returns.
145const MAX_BANKED_PROGRESS_GRANTS: usize = 8 * CONSOLIDATION_GRANTS_PER_PROGRESS;
146
147impl<B: Batch+Clone+'static> TraceReader for Spine<B> {
148
149    type Time = B::Time;
150    type Batch = B;
151
152    fn batches_through(&mut self, upper: AntichainRef<Self::Time>) -> Option<Vec<Self::Batch>> {
153
154        // If `upper` is the minimum frontier, we can return an empty cursor.
155        // This can happen with operators that are written to expect the ability to acquire cursors
156        // for their prior frontiers, and which start at `[T::minimum()]`, such as `Reduce`, sadly.
157        if upper.less_equal(&<Self::Time as timely::progress::Timestamp>::minimum()) {
158            return Some(Vec::new());
159        }
160
161        // The supplied `upper` should have the property that for each of our
162        // batch `lower` and `upper` frontiers, the supplied upper is comparable
163        // to the frontier; it should not be incomparable, because the frontiers
164        // that we created form a total order. If it is, there is a bug.
165        //
166        // We should acquire a cursor including all batches whose upper is less
167        // or equal to the supplied upper, excluding all batches whose lower is
168        // greater or equal to the supplied upper, and if a batch straddles the
169        // supplied upper it had better be empty.
170
171        // We shouldn't grab a cursor into a closed trace, right?
172        assert!(self.logical_frontier.borrow().len() > 0);
173
174        // Check that `upper` is greater or equal to `self.physical_frontier`.
175        // Otherwise, the cut could be in `self.merging` and it is user error anyhow.
176        // assert!(upper.iter().all(|t1| self.physical_frontier.iter().any(|t2| t2.less_equal(t1))));
177        assert!(PartialOrder::less_equal(&self.physical_frontier.borrow(), &upper));
178
179        let mut storage = Vec::new();
180
181        for merge_state in self.merging.iter().rev() {
182            match merge_state {
183                MergeState::Double(variant) => {
184                    match variant {
185                        MergeVariant::InProgress(batch1, batch2, _) => {
186                            if !batch1.is_empty() {
187                                storage.push(batch1.clone());
188                            }
189                            if !batch2.is_empty() {
190                                storage.push(batch2.clone());
191                            }
192                        },
193                        MergeVariant::Complete(Some((batch, _))) => {
194                            if !batch.is_empty() {
195                                storage.push(batch.clone());
196                            }
197                        }
198                        MergeVariant::Complete(None) => { },
199                    }
200                },
201                MergeState::Single(Some(batch)) => {
202                    if !batch.is_empty() {
203                        storage.push(batch.clone());
204                    }
205                },
206                MergeState::Single(None) => { },
207                MergeState::Vacant => { },
208            }
209        }
210
211        for batch in self.pending.iter() {
212
213            if !batch.is_empty() {
214
215                // For a non-empty `batch`, it is a catastrophic error if `upper`
216                // requires some-but-not-all of the updates in the batch. We can
217                // determine this from `upper` and the lower and upper bounds of
218                // the batch itself.
219                //
220                // TODO: It is not clear if this is the 100% correct logic, due
221                // to the possible non-total-orderedness of the frontiers.
222
223                let include_lower = PartialOrder::less_equal(&batch.lower().borrow(), &upper);
224                let include_upper = PartialOrder::less_equal(&batch.upper().borrow(), &upper);
225
226                if include_lower != include_upper && upper != batch.lower().borrow() {
227                    panic!("`cursor_through`: `upper` straddles batch");
228                }
229
230                // include pending batches
231                if include_upper {
232                    storage.push(batch.clone());
233                }
234            }
235        }
236
237        Some(storage)
238    }
239    #[inline]
240    fn set_logical_compaction(&mut self, frontier: AntichainRef<B::Time>) {
241        self.logical_frontier.clear();
242        self.logical_frontier.extend(frontier.iter().cloned());
243    }
244    #[inline]
245    fn get_logical_compaction(&mut self) -> AntichainRef<'_, B::Time> { self.logical_frontier.borrow() }
246    #[inline]
247    fn set_physical_compaction(&mut self, frontier: AntichainRef<'_, B::Time>) {
248        // We should never request to rewind the frontier.
249        debug_assert!(PartialOrder::less_equal(&self.physical_frontier.borrow(), &frontier), "FAIL\tthrough frontier !<= new frontier {:?} {:?}\n", self.physical_frontier, frontier);
250        self.physical_frontier.clear();
251        self.physical_frontier.extend(frontier.iter().cloned());
252        self.consider_merges();
253    }
254    #[inline]
255    fn get_physical_compaction(&mut self) -> AntichainRef<'_, B::Time> { self.physical_frontier.borrow() }
256
257    #[inline]
258    fn map_batches<F: FnMut(&Self::Batch)>(&self, mut f: F) {
259        for batch in self.merging.iter().rev() {
260            match batch {
261                MergeState::Double(MergeVariant::InProgress(batch1, batch2, _)) => { f(batch1); f(batch2); },
262                MergeState::Double(MergeVariant::Complete(Some((batch, _)))) => { f(batch) },
263                MergeState::Single(Some(batch)) => { f(batch) },
264                _ => { },
265            }
266        }
267        for batch in self.pending.iter() {
268            f(batch);
269        }
270    }
271}
272
273// A trace implementation for any key type that can be borrowed from or converted into `Key`.
274// TODO: Almost all this implementation seems to be generic with respect to the trace and batch types.
275impl<B: Batch+Clone+'static> Trace for Spine<B> {
276    fn new(
277        info: ::timely::dataflow::operators::generic::OperatorInfo,
278        logging: Option<differential_dataflow::logging::Logger>,
279        activator: Option<timely::scheduling::activate::Activator>,
280    ) -> Self {
281        Self::with_effort(1, info, logging, activator)
282    }
283
284    /// Apply some amount of effort to trace maintenance.
285    ///
286    /// Whether and how much effort to apply is determined by `self.exert_logic`, a closure the user can set.
287    fn exert(&mut self) {
288        // If there is work to be done, ...
289        self.tidy_layers();
290        // Determine whether we should apply effort independent of updates.
291        if let Some(effort) = self.exert_effort() {
292            // Optional effort is paid for by inserted updates and frontier advances
293            // while the input is open; only a closed input lifts that bound. A
294            // merge already in progress must complete before anything else lands
295            // at its level, so advancing it adds no work and proceeds unfunded.
296            // Funding is still spent first when available, so a funded spine does
297            // the same work either way.
298            let funded = self.upper.borrow().is_empty() || self.spend_funding(effort);
299            if funded || self.merging.iter().any(|b| b.is_double()) {
300                self.grant(effort);
301            }
302            // Unlike upstream, a declined request does not reactivate the
303            // operator. There is no funded work until the next insert, and
304            // the insert schedules the operator itself.
305        }
306    }
307
308    fn set_exert_logic(&mut self, logic: ExertionLogic) {
309        self.exert_logic = Some(logic);
310    }
311
312    // Ideally, this method acts as insertion of `batch`, even if we are not yet able to begin
313    // merging the batch. This means it is a good time to perform amortized work proportional
314    // to the size of batch.
315    fn insert(&mut self, batch: Self::Batch) {
316
317        // Log the introduction of a batch.
318        self.logger.as_ref().map(|l| l.log(differential_dataflow::logging::BatchEvent {
319            operator: self.operator.global_id,
320            length: batch.len()
321        }));
322
323        assert!(batch.lower() != batch.upper());
324        assert_eq!(batch.lower(), &self.upper);
325
326        self.upper.clone_from(batch.upper());
327        self.consolidation_credit = self.consolidation_credit.saturating_add(
328            batch.len().saturating_mul(CONSOLIDATION_CREDIT_PER_UPDATE).saturating_mul(self.effort),
329        );
330        self.progress_grants = (self.progress_grants + CONSOLIDATION_GRANTS_PER_PROGRESS)
331            .min(MAX_BANKED_PROGRESS_GRANTS);
332
333        // TODO: Consolidate or discard empty batches.
334        self.pending.push(batch);
335        self.consider_merges();
336    }
337
338    /// Completes the trace with a final empty batch.
339    fn close(&mut self) {
340        if !self.upper.borrow().is_empty() {
341            self.insert(B::empty(self.upper.clone(), Antichain::new()));
342        }
343    }
344}
345
346// Drop implementation allows us to log batch drops, to zero out maintained totals.
347impl<B: Batch> Drop for Spine<B> {
348    fn drop(&mut self) {
349        self.drop_batches();
350    }
351}
352
353
354impl<B: Batch> Spine<B> {
355    /// Drops and logs batches. Used in `set_logical_compaction` and drop.
356    fn drop_batches(&mut self) {
357        if let Some(logger) = &self.logger {
358            for batch in self.merging.drain(..) {
359                match batch {
360                    MergeState::Single(Some(batch)) => {
361                        logger.log(differential_dataflow::logging::DropEvent {
362                            operator: self.operator.global_id,
363                            length: batch.len(),
364                        });
365                    },
366                    MergeState::Double(MergeVariant::InProgress(batch1, batch2, _)) => {
367                        logger.log(differential_dataflow::logging::DropEvent {
368                            operator: self.operator.global_id,
369                            length: batch1.len(),
370                        });
371                        logger.log(differential_dataflow::logging::DropEvent {
372                            operator: self.operator.global_id,
373                            length: batch2.len(),
374                        });
375                    },
376                    MergeState::Double(MergeVariant::Complete(Some((batch, _)))) => {
377                        logger.log(differential_dataflow::logging::DropEvent {
378                            operator: self.operator.global_id,
379                            length: batch.len(),
380                        });
381                    }
382                    _ => { },
383                }
384            }
385            for batch in self.pending.drain(..) {
386                logger.log(differential_dataflow::logging::DropEvent {
387                    operator: self.operator.global_id,
388                    length: batch.len(),
389                });
390            }
391        }
392    }
393}
394
395impl<B: Batch> Spine<B> {
396    /// Determine the amount of effort we should exert in the absence of updates.
397    ///
398    /// This method prepares an iterator over batches, including the level, count, and length of each layer.
399    /// It supplies this to `self.exert_logic`, who produces the response of the amount of exertion to apply.
400    fn exert_effort(&mut self) -> Option<usize> {
401        self.exert_logic.as_ref().and_then(|exert_logic| {
402            self.exert_logic_param.clear();
403            self.exert_logic_param.extend(self.merging.iter().enumerate().rev().map(|(index, batch)| {
404                match batch {
405                    MergeState::Vacant => (index, 0, 0),
406                    MergeState::Single(_) => (index, 1, batch.len()),
407                    MergeState::Double(_) => (index, 2, batch.len()),
408                }
409            }));
410
411            (exert_logic)(&self.exert_logic_param[..])
412        })
413    }
414
415    /// Allocates a fueled `Spine` with a specified effort multiplier.
416    ///
417    /// This trace will merge batches progressively, with each inserted batch applying a multiple
418    /// of the batch's length in effort to each merge. The `effort` parameter is that multiplier.
419    /// This value should be at least one for the merging to happen; a value of zero is not helpful.
420    pub fn with_effort(
421        mut effort: usize,
422        operator: OperatorInfo,
423        logger: Option<differential_dataflow::logging::Logger>,
424        activator: Option<timely::scheduling::activate::Activator>,
425    ) -> Self {
426
427        // Zero effort is .. not smart.
428        if effort == 0 { effort = 1; }
429
430        Spine {
431            operator,
432            logger,
433            logical_frontier: Antichain::from_elem(<B::Time as timely::progress::Timestamp>::minimum()),
434            physical_frontier: Antichain::from_elem(<B::Time as timely::progress::Timestamp>::minimum()),
435            merging: Vec::new(),
436            pending: Vec::new(),
437            upper: Antichain::from_elem(<B::Time as timely::progress::Timestamp>::minimum()),
438            effort,
439            activator,
440            exert_logic_param: Vec::default(),
441            exert_logic: None,
442            consolidation_credit: 0,
443            progress_grants: 0,
444        }
445    }
446
447    /// Spend funding for one policy request, returning whether any was available.
448    fn spend_funding(&mut self, effort: usize) -> bool {
449        if self.consolidation_credit >= effort {
450            self.consolidation_credit -= effort;
451            true
452        } else if self.progress_grants > 0 {
453            self.progress_grants -= 1;
454            true
455        } else {
456            false
457        }
458    }
459
460    /// Apply one policy allowance: fuel for active merges, or a virtual introduction.
461    fn grant(&mut self, effort: usize) {
462        // If any merges exist, we can directly call `apply_fuel`.
463        if self.merging.iter().any(|b| b.is_double()) {
464            self.apply_fuel(&mut (effort as isize));
465        }
466        // Otherwise, we'll need to introduce fake updates to move merges along.
467        else {
468            // Introduce an empty batch with roughly *effort number of virtual updates.
469            let level = effort.next_power_of_two().trailing_zeros() as usize;
470            self.introduce_batch(None, level);
471        }
472        // We were not in reduced form, so let's check again in the future.
473        if let Some(activator) = &self.activator {
474            activator.activate();
475        }
476    }
477
478    /// Migrate data from `self.pending` into `self.merging`.
479    ///
480    /// This method reflects on the bookmarks held by others that may prevent merging, and in the
481    /// case that new batches can be introduced to the pile of mergeable batches, it gets on that.
482    #[inline(never)]
483    fn consider_merges(&mut self) {
484
485        // TODO: Consider merging pending batches before introducing them.
486        // TODO: We could use a `VecDeque` here to draw from the front and append to the back.
487        while !self.pending.is_empty() && PartialOrder::less_equal(self.pending[0].upper(), &self.physical_frontier)
488            //   self.physical_frontier.iter().all(|t1| self.pending[0].upper().iter().any(|t2| t2.less_equal(t1)))
489        {
490            // Batch can be taken in optimized insertion.
491            // Otherwise it is inserted normally at the end of the method.
492            let mut batch = Some(self.pending.remove(0));
493
494            // If `batch` and the most recently inserted batch are both empty, we can just fuse them.
495            // We can also replace a structurally empty batch with this empty batch, preserving the
496            // apparent record count but now with non-trivial lower and upper bounds.
497            if batch.as_ref().unwrap().len() == 0 {
498                if let Some(position) = self.merging.iter().position(|m| !m.is_vacant()) {
499                    if self.merging[position].is_single() && self.merging[position].len() == 0 {
500                        self.insert_at(batch.take(), position);
501                        let merged = self.complete_at(position);
502                        self.merging[position] = MergeState::Single(merged);
503                    }
504                }
505            }
506
507            // Normal insertion for the batch.
508            if let Some(batch) = batch {
509                let index = batch.len().next_power_of_two();
510                self.introduce_batch(Some(batch), index.trailing_zeros() as usize);
511            }
512        }
513
514        // Having performed all of our work, if we should perform more work reschedule ourselves.
515        if self.exert_effort().is_some() {
516            if let Some(activator) = &self.activator {
517                activator.activate();
518            }
519        }
520    }
521
522    /// Introduces a batch at an indicated level.
523    ///
524    /// The level indication is often related to the size of the batch, but
525    /// it can also be used to artificially fuel the computation by supplying
526    /// empty batches at non-trivial indices, to move merges along.
527    pub fn introduce_batch(&mut self, batch: Option<B>, batch_index: usize) {
528
529        // Step 0.  Determine an amount of fuel to use for the computation.
530        //
531        //          Fuel is used to drive maintenance of the data structure,
532        //          and in particular are used to make progress through merges
533        //          that are in progress. The amount of fuel to use should be
534        //          proportional to the number of records introduced, so that
535        //          we are guaranteed to complete all merges before they are
536        //          required as arguments to merges again.
537        //
538        //          The fuel use policy is negotiable, in that we might aim
539        //          to use relatively less when we can, so that we return
540        //          control promptly, or we might account more work to larger
541        //          batches. Not clear to me which are best, of if there
542        //          should be a configuration knob controlling this.
543
544        // The amount of fuel to use is proportional to 2^batch_index, scaled
545        // by a factor of self.effort which determines how eager we are in
546        // performing maintenance work. We need to ensure that each merge in
547        // progress receives fuel for each introduced batch, and so multiply
548        // by that as well.
549        if batch_index > 32 { println!("Large batch index: {}", batch_index); }
550
551        // We believe that eight units of fuel is sufficient for each introduced
552        // record, accounted as four for each record, and a potential four more
553        // for each virtual record associated with promoting existing smaller
554        // batches. We could try and make this be less, or be scaled to merges
555        // based on their deficit at time of instantiation. For now, we remain
556        // conservative.
557        let mut fuel = 8 << batch_index;
558        // Scale up by the effort parameter, which is calibrated to one as the
559        // minimum amount of effort.
560        fuel *= self.effort;
561        // Convert to an `isize` so we can observe any fuel shortfall.
562        let mut fuel = fuel as isize;
563
564        // Step 1.  Apply fuel to each in-progress merge.
565        //
566        //          Before we can introduce new updates, we must apply any
567        //          fuel to in-progress merges, as this fuel is what ensures
568        //          that the merges will be complete by the time we insert
569        //          the updates.
570        self.apply_fuel(&mut fuel);
571
572        // Step 2.  We must ensure the invariant that adjacent layers do not
573        //          contain two batches will be satisfied when we insert the
574        //          batch. We forcibly completing all merges at layers lower
575        //          than and including `batch_index`, so that the new batch
576        //          is inserted into an empty layer.
577        //
578        //          We could relax this to "strictly less than `batch_index`"
579        //          if the layer above has only a single batch in it, which
580        //          seems not implausible if it has been the focus of effort.
581        //
582        //          This should be interpreted as the introduction of some
583        //          volume of fake updates, and we will need to fuel merges
584        //          by a proportional amount to ensure that they are not
585        //          surprised later on. The number of fake updates should
586        //          correspond to the deficit for the layer, which perhaps
587        //          we should track explicitly.
588        self.roll_up(batch_index);
589
590        // Step 3. This insertion should be into an empty layer. It is a
591        //         logical error otherwise, as we may be violating our
592        //         invariant, from which all wonderment derives.
593        self.insert_at(batch, batch_index);
594
595        // Step 4. Tidy the largest layers.
596        //
597        //         It is important that we not tidy only smaller layers,
598        //         as their ascension is what ensures the merging and
599        //         eventual compaction of the largest layers.
600        self.tidy_layers();
601    }
602
603    /// Ensures that an insertion at layer `index` will succeed.
604    ///
605    /// This method is subject to the constraint that all existing batches
606    /// should occur at higher levels, which requires it to "roll up" batches
607    /// present at lower levels before the method is called. In doing this,
608    /// we should not introduce more virtual records than 2^index, as that
609    /// is the amount of excess fuel we have budgeted for completing merges.
610    fn roll_up(&mut self, index: usize) {
611
612        // Ensure entries sufficient for `index`.
613        while self.merging.len() <= index {
614            self.merging.push(MergeState::Vacant);
615        }
616
617        // We only need to roll up if there are non-vacant layers.
618        if self.merging[.. index].iter().any(|m| !m.is_vacant()) {
619
620            // Collect and merge all batches at layers up to but not including `index`.
621            let mut merged = None;
622            for i in 0 .. index {
623                self.insert_at(merged, i);
624                merged = self.complete_at(i);
625            }
626
627            // The merged results should be introduced at level `index`, which should
628            // be ready to absorb them (possibly creating a new merge at the time).
629            self.insert_at(merged, index);
630
631            // If the insertion results in a merge, we should complete it to ensure
632            // the upcoming insertion at `index` does not panic.
633            if self.merging[index].is_double() {
634                let merged = self.complete_at(index);
635                self.insert_at(merged, index + 1);
636            }
637        }
638    }
639
640    /// Applies an amount of fuel to merges in progress.
641    ///
642    /// The supplied `fuel` is for each in progress merge, and if we want to spend
643    /// the fuel non-uniformly (e.g. prioritizing merges at low layers) we could do
644    /// so in order to maintain fewer batches on average (at the risk of completing
645    /// merges of large batches later, but tbh probably not much later).
646    pub fn apply_fuel(&mut self, fuel: &mut isize) {
647        // For the moment our strategy is to apply fuel independently to each merge
648        // in progress, rather than prioritizing small merges. This sounds like a
649        // great idea, but we need better accounting in place to ensure that merges
650        // that borrow against later layers but then complete still "acquire" fuel
651        // to pay back their debts.
652        for index in 0 .. self.merging.len() {
653            // Give each level independent fuel, for now.
654            let mut fuel = *fuel;
655            // Pass along various logging stuffs, in case we need to report success.
656            self.merging[index].work(&mut fuel);
657            // `fuel` could have a deficit at this point, meaning we over-spent when
658            // we took a merge step. We could ignore this, or maintain the deficit
659            // and account future fuel against it before spending again. It isn't
660            // clear why that would be especially helpful to do; we might want to
661            // avoid overspends at multiple layers in the same invocation (to limit
662            // latencies), but there is probably a rich policy space here.
663
664            // If a merge completes, we can immediately merge it in to the next
665            // level, which is "guaranteed" to be complete at this point, by our
666            // fueling discipline.
667            if self.merging[index].is_complete() {
668                let complete = self.complete_at(index);
669                self.insert_at(complete, index+1);
670            }
671        }
672    }
673
674    /// Inserts a batch at a specific location.
675    ///
676    /// This is a non-public internal method that can panic if we try and insert into a
677    /// layer which already contains two batches (and is still in the process of merging).
678    fn insert_at(&mut self, batch: Option<B>, index: usize) {
679        // Ensure the spine is large enough.
680        while self.merging.len() <= index {
681            self.merging.push(MergeState::Vacant);
682        }
683
684        // Insert the batch at the location.
685        match self.merging[index].take() {
686            MergeState::Vacant => {
687                self.merging[index] = MergeState::Single(batch);
688            }
689            MergeState::Single(old) => {
690                // Log the initiation of a merge.
691                self.logger.as_ref().map(|l| l.log(
692                    differential_dataflow::logging::MergeEvent {
693                        operator: self.operator.global_id,
694                        scale: index,
695                        length1: old.as_ref().map(|b| b.len()).unwrap_or(0),
696                        length2: batch.as_ref().map(|b| b.len()).unwrap_or(0),
697                        complete: None,
698                    }
699                ));
700                let compaction_frontier = self.logical_frontier.borrow();
701                self.merging[index] = MergeState::begin_merge(old, batch, compaction_frontier);
702            }
703            MergeState::Double(_) => {
704                panic!("Attempted to insert batch into incomplete merge!")
705            }
706        };
707    }
708
709    /// Completes and extracts what ever is at layer `index`.
710    fn complete_at(&mut self, index: usize) -> Option<B> {
711        if let Some((merged, inputs)) = self.merging[index].complete() {
712            if let Some((input1, input2)) = inputs {
713                // Log the completion of a merge from existing parts.
714                self.logger.as_ref().map(|l| l.log(
715                    differential_dataflow::logging::MergeEvent {
716                        operator: self.operator.global_id,
717                        scale: index,
718                        length1: input1.len(),
719                        length2: input2.len(),
720                        complete: Some(merged.len()),
721                    }
722                ));
723            }
724            Some(merged)
725        }
726        else {
727            None
728        }
729    }
730
731    /// Attempts to draw down large layers to size appropriate layers.
732    fn tidy_layers(&mut self) {
733
734        // If the largest layer is complete (not merging), we can attempt
735        // to draw it down to the next layer. This is permitted if we can
736        // maintain our invariant that below each merge there are at most
737        // half the records that would be required to invade the merge.
738        if !self.merging.is_empty() {
739            let mut length = self.merging.len();
740            if self.merging[length-1].is_single() {
741
742                // To move a batch down, we require that it contain few
743                // enough records that the lower level is appropriate,
744                // and that moving the batch would not create a merge
745                // violating our invariant.
746
747                let appropriate_level = self.merging[length-1].len().next_power_of_two().trailing_zeros() as usize;
748
749                // Continue only as far as is appropriate
750                while appropriate_level < length-1 {
751
752                    match self.merging[length-2].take() {
753                        // Vacant or structurally empty batches can be absorbed.
754                        MergeState::Vacant | MergeState::Single(None) => {
755                            self.merging.remove(length-2);
756                            length = self.merging.len();
757                        }
758                        // Single batches may initiate a merge, if sizes are
759                        // within bounds, but terminate the loop either way.
760                        MergeState::Single(Some(batch)) => {
761
762                            // Determine the number of records that might lead
763                            // to a merge. Importantly, this is not the number
764                            // of actual records, but the sum of upper bounds
765                            // based on indices.
766                            let mut smaller = 0;
767                            for (index, batch) in self.merging[..(length-2)].iter().enumerate() {
768                                match batch {
769                                    MergeState::Vacant => { },
770                                    MergeState::Single(_) => { smaller += 1 << index; },
771                                    MergeState::Double(_) => { smaller += 2 << index; },
772                                }
773                            }
774
775                            if smaller <= (1 << length) / 8 {
776                                self.merging.remove(length-2);
777                                self.insert_at(Some(batch), length-2);
778                            }
779                            else {
780                                self.merging[length-2] = MergeState::Single(Some(batch));
781                            }
782                            return;
783                        }
784                        // If a merge is in progress there is nothing to do.
785                        MergeState::Double(state) => {
786                            self.merging[length-2] = MergeState::Double(state);
787                            return;
788                        }
789                    }
790                }
791            }
792        }
793    }
794}
795
796
797/// Describes the state of a layer.
798///
799/// A layer can be empty, contain a single batch, or contain a pair of batches
800/// that are in the process of merging into a batch for the next layer.
801enum MergeState<B: Batch> {
802    /// An empty layer, containing no updates.
803    Vacant,
804    /// A layer containing a single batch.
805    ///
806    /// The `None` variant is used to represent a structurally empty batch present
807    /// to ensure the progress of maintenance work.
808    Single(Option<B>),
809    /// A layer containing two batches, in the process of merging.
810    Double(MergeVariant<B>),
811}
812
813impl<B: Batch<Time: Eq>> MergeState<B> {
814
815    /// The number of actual updates contained in the level.
816    fn len(&self) -> usize {
817        match self {
818            MergeState::Single(Some(b)) => b.len(),
819            MergeState::Double(MergeVariant::InProgress(b1,b2,_)) => b1.len() + b2.len(),
820            MergeState::Double(MergeVariant::Complete(Some((b, _)))) => b.len(),
821            _ => 0,
822        }
823    }
824
825    /// True only for the MergeState::Vacant variant.
826    fn is_vacant(&self) -> bool {
827        if let MergeState::Vacant = self { true } else { false }
828    }
829
830    /// True only for the MergeState::Single variant.
831    fn is_single(&self) -> bool {
832        if let MergeState::Single(_) = self { true } else { false }
833    }
834
835    /// True only for the MergeState::Double variant.
836    fn is_double(&self) -> bool {
837        if let MergeState::Double(_) = self { true } else { false }
838    }
839
840    /// Immediately complete any merge.
841    ///
842    /// The result is either a batch, if there is a non-trivial batch to return
843    /// or `None` if there is no meaningful batch to return. This does not distinguish
844    /// between Vacant entries and structurally empty batches, which should be done
845    /// with the `is_complete()` method.
846    ///
847    /// There is the additional option of input batches.
848    fn complete(&mut self) -> Option<(B, Option<(B, B)>)>  {
849        match std::mem::replace(self, MergeState::Vacant) {
850            MergeState::Vacant => None,
851            MergeState::Single(batch) => batch.map(|b| (b, None)),
852            MergeState::Double(variant) => variant.complete(),
853        }
854    }
855
856    /// True iff the layer is a complete merge, ready for extraction.
857    fn is_complete(&mut self) -> bool {
858        if let MergeState::Double(MergeVariant::Complete(_)) = self {
859            true
860        }
861        else {
862            false
863        }
864    }
865
866    /// Performs a bounded amount of work towards a merge.
867    ///
868    /// If the merge completes, the resulting batch is returned.
869    /// If a batch is returned, it is the obligation of the caller
870    /// to correctly install the result.
871    fn work(&mut self, fuel: &mut isize) {
872        // We only perform work for merges in progress.
873        if let MergeState::Double(layer) = self {
874            layer.work(fuel)
875        }
876    }
877
878    /// Extract the merge state, typically temporarily.
879    fn take(&mut self) -> Self {
880        std::mem::replace(self, MergeState::Vacant)
881    }
882
883    /// Initiates the merge of an "old" batch with a "new" batch.
884    ///
885    /// The upper frontier of the old batch should match the lower
886    /// frontier of the new batch, with the resulting batch describing
887    /// their composed interval, from the lower frontier of the old
888    /// batch to the upper frontier of the new batch.
889    ///
890    /// Either batch may be `None` which corresponds to a structurally
891    /// empty batch whose upper and lower froniers are equal. This
892    /// option exists purely for bookkeeping purposes, and no computation
893    /// is performed to merge the two batches.
894    fn begin_merge(batch1: Option<B>, batch2: Option<B>, compaction_frontier: AntichainRef<B::Time>) -> MergeState<B> {
895        let variant =
896        match (batch1, batch2) {
897            (Some(batch1), Some(batch2)) => {
898                assert!(batch1.upper() == batch2.lower());
899                let begin_merge = <B as Batch>::begin_merge(&batch1, &batch2, compaction_frontier);
900                MergeVariant::InProgress(batch1, batch2, begin_merge)
901            }
902            (None, Some(x)) => MergeVariant::Complete(Some((x, None))),
903            (Some(x), None) => MergeVariant::Complete(Some((x, None))),
904            (None, None) => MergeVariant::Complete(None),
905        };
906
907        MergeState::Double(variant)
908    }
909}
910
911enum MergeVariant<B: Batch> {
912    /// Describes an actual in-progress merge between two non-trivial batches.
913    InProgress(B, B, <B as Batch>::Merger),
914    /// A merge that requires no further work. May or may not represent a non-trivial batch.
915    Complete(Option<(B, Option<(B, B)>)>),
916}
917
918impl<B: Batch> MergeVariant<B> {
919
920    /// Completes and extracts the batch, unless structurally empty.
921    ///
922    /// The result is either `None`, for structurally empty batches,
923    /// or a batch and optionally input batches from which it derived.
924    fn complete(mut self) -> Option<(B, Option<(B, B)>)> {
925        let mut fuel = isize::MAX;
926        self.work(&mut fuel);
927        if let MergeVariant::Complete(batch) = self { batch }
928        else { panic!("Failed to complete a merge!"); }
929    }
930
931    /// Applies some amount of work, potentially completing the merge.
932    ///
933    /// In case the work completes, the source batches are returned.
934    /// This allows the caller to manage the released resources.
935    fn work(&mut self, fuel: &mut isize) {
936        let variant = std::mem::replace(self, MergeVariant::Complete(None));
937        if let MergeVariant::InProgress(b1,b2,mut merge) = variant {
938            merge.work(&b1,&b2,fuel);
939            if *fuel > 0 {
940                *self = MergeVariant::Complete(Some((merge.done(), Some((b1,b2)))));
941            }
942            else {
943                *self = MergeVariant::InProgress(b1,b2,merge);
944            }
945        }
946        else {
947            *self = variant;
948        }
949    }
950}