Skip to main content

mz_timely_util/
reclock.rs

1// Copyright Materialize, Inc. and contributors. All rights reserved.
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License in the LICENSE file at the
6// root of this repository, or online at
7//
8//     http://www.apache.org/licenses/LICENSE-2.0
9//
10// Unless required by applicable law or agreed to in writing, software
11// distributed under the License is distributed on an "AS IS" BASIS,
12// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
13// See the License for the specific language governing permissions and
14// limitations under the License.
15
16//! ## Notation
17//!
18//! Collections are represented with capital letters (T, S, R), collection traces as bold letters
19//! (𝐓, 𝐒, 𝐑), and difference traces as δ𝐓.
20//!
21//! Indexing a collection trace 𝐓 to obtain its version at `t` is written as 𝐓(t). Indexing a
22//! collection to obtain the multiplicity of a record `x` is written as T\[x\]. These can be combined
23//! to obtain the multiplicity of a record `x` at some version `t` as 𝐓(t)\[x\].
24//!
25//! ## Overview
26//!
27//! Reclocking transforms a source collection `S` that evolves with some timestamp `FromTime` into
28//! a collection `T` that evolves with some other timestamp `IntoTime`. The reclocked collection T
29//! contains all updates `u ∈ S` that are not beyond some `FromTime` frontier R(t). The collection
30//! `R` is called the remap collection.
31//!
32//! More formally, for some arbitrary time `t` of `IntoTime` and some arbitrary record `x`, the
33//! reclocked collection `T(t)[x]` is defined to be the `sum{δ𝐒(s)[x]: !(𝐑(t) βͺ― s)}`. Since this
34//! holds for any record we can write the definition of Reclock(𝐒, 𝐑) as:
35//!
36//! > Reclock(𝐒, 𝐑) β‰œ 𝐓: βˆ€ t ∈ IntoTime : 𝐓(t) = sum{δ𝐒(s): !(𝐑(t) βͺ― s)}
37//!
38//! In order for the reclocked collection `T` to have a sensible definition of progress we require
39//! that `t1 ≀ t2 β‡’ 𝐑(t1) βͺ― 𝐑(t2)` where the first `≀` is the partial order of `IntoTime` and the
40//! second one the partial order of `FromTime` antichains.
41//!
42//! ## Total order simplification
43//!
44//! In order to simplify the implementation we will require that `IntoTime` is a total order. This
45//! limitation can be lifted in the future but further elaboration on the mechanics of reclocking
46//! is required to ensure a correct implementation.
47//!
48//! ## The difference trace
49//!
50//! By the definition of difference traces we have:
51//!
52//! ```text
53//!     δ𝐓(t) = T(t) - sum{δ𝐓(s): s < t}
54//! ```
55//!
56//! Due to the total order assumption we only need to consider two cases.
57//!
58//! **Case 1:** `t` is the minimum timestamp
59//!
60//! In this case `sum{δ𝐓(s): s < t}` is the empty set and so we obtain:
61//!
62//! ```text
63//!     δ𝐓(min) = T(min) = sum{δ𝐒(s): !(𝐑(min) ≀ s}
64//! ```
65//!
66//! **Case 2:** `t` is a timestamp with a predecessor `prev`
67//!
68//! In this case `sum{δ𝐓(s): s < t}` is equal to `T(prev)` because:
69//!
70//! ```text
71//!     sum{δ𝐓(s): s < t} = sum{δ𝐓(s): s ≀ prev} + sum{δ𝐓(s): prev < s < t}
72//!                       = T(prev) + βˆ…
73//!                       = T(prev)
74//! ```
75//!
76//! And therefore the difference trace of T is:
77//!
78//! ```text
79//!     δ𝐓(t) = 𝐓(t) - 𝐓(prev)
80//!           = sum{δ𝐒(s): !(𝐑(t) βͺ― s)} - sum{δ𝐒(s): !(𝐑(prev) βͺ― s)}
81//!           = sum{δ𝐒(s): (𝐑(prev) βͺ― s) ∧ !(𝐑(t) βͺ― s)}
82//! ```
83//!
84//! ## Unique mapping property
85//!
86//! Given the definition above we can derive the fact that for any source difference δ𝐒(s) there is
87//! at most one target timestamp t that it must be reclocked to. This property can be exploited by
88//! the implementation of the operator as it can safely discard source updates once a matching
89//! Ξ΄T(t) has been found, making it "stateless" with respect to the source trace. A formal proof of
90//! this property is [provided below](#unique-mapping-property-proof).
91//!
92//! ## Operational description
93//!
94//! The operator follows a run-to-completion model where on each scheduling it completes all
95//! outstanding work that can be completed.
96//!
97//! ### Unique mapping property proof
98//!
99//! This section contains the formal proof the unique mapping property. The proof follows the
100//! structure proof notation created by Leslie Lamport. Readers unfamiliar with structured proofs
101//! can read about them here <https://lamport.azurewebsites.net/pubs/proof.pdf>.
102//!
103//! #### Statement
104//!
105//! AtMostOne(X, Ο†(x)) β‰œ βˆ€ x1, x2 ∈ X : Ο†(x1) ∧ Ο†(x2) β‡’ x1 = x2
106//!
107//! * **THEOREM** UniqueMapping β‰œ
108//!     * **ASSUME**
109//!         * **NEW** (FromTime, βͺ―) ∈ PartiallyOrderedTimestamps
110//!         * **NEW** (IntoTime, ≀) ∈ TotallyOrderedTimestamps
111//!         * **NEW** 𝐒 ∈ SetOfCollectionTraces(FromTime)
112//!         * **NEW** 𝐑 ∈ SetOfCollectionTraces(IntoTime)
113//!         * βˆ€ t ∈ IntoTime: 𝐑(t) ∈ SetOfAntichains(FromTime)
114//!         * βˆ€ t1, t1 ∈ IntoTime: t1 ≀ t2 β‡’ 𝐑(t1) βͺ― 𝐑(t2)
115//!         * **NEW** 𝐓 = Reclock(𝐒, 𝐑)
116//!     * **PROVE**  βˆ€ s ∈ FromTime : AtMostOne(IntoTime, δ𝐒(s) ∈ δ𝐓(x))
117//!
118//! #### Proof
119//!
120//! 1. **SUFFICES ASSUME** βˆƒ s ∈ FromTime: Β¬AtMostOne(IntoTime, δ𝐒(s) ∈ δ𝐓(x))
121//!     * **PROVE FALSE**
122//!     * _By proof by contradiction._
123//! 2. **PICK** s ∈ FromTime : Β¬AtMostOne(IntoTime, δ𝐒(s) ∈ δ𝐓(x))
124//!    * _Proof: Such time exists by <1>1._
125//! 3. βˆƒ t1, t2 ∈ IntoTime : t1 β‰  t2 ∧ δ𝐒(s) ∈ δ𝐓(t1) ∧ δ𝐒(s) ∈ δ𝐓(t2)
126//!     1. Β¬(βˆ€ x1, x2 ∈ X : (δ𝐒(s) ∈ δ𝐓(x1)) ∧ (δ𝐒(s) ∈ δ𝐓(x2)) β‡’ x1 = x2)
127//!         * _Proof: By <1>2 and definition of AtMostOne._
128//!     2. Q.E.D
129//!         * _Proof: By <2>1, quantifier negation rules, and theorem of propositional logic Β¬(P β‡’ Q) ≑ P ∧ Β¬Q._
130//! 4. **PICK** t1, t2 ∈ IntoTime : t1 < t2 ∧ δ𝐒(s) ∈ δ𝐓(t1) ∧ δ𝐒(s) ∈ δ𝐓(t2)
131//!    * _Proof: By <1>3. Assume t1 < t2 without loss of generality._
132//! 5. Β¬(𝐑(t1) βͺ― s)
133//!     1. **CASE** t1 = min(IntoTime)
134//!         1. δ𝐓(t1) = sum{δ𝐒(s): !(𝐑(t1)) βͺ― s}
135//!             * _Proof: By definition of δ𝐓(min)._
136//!         2. δ𝐒(s) ∈ δ𝐓(t1)
137//!             * _Proof: By <1>4._
138//!         3. Q.E.D
139//!             * _Proof: By <3>1 and <3>2._
140//!     2. **CASE** t1 > min(IntoTime)
141//!         1. **PICK** t1_prev = Predecessor(t1)
142//!             * _Proof: Predecessor exists because the set {t: t < t1} is non-empty since it must contain at least min(IntoTime)._
143//!         2. δ𝐓(t1) = sum{δ𝐒(s): (𝐑(t1_prev) βͺ― s) ∧ !(𝐑(t1) βͺ― s)}
144//!             * _Proof: By definition of δ𝐓(t)._
145//!         3. δ𝐒(s) ∈ δ𝐓(t1)
146//!             * _Proof: By <1>4._
147//!         3. Q.E.D
148//!             * _Proof: By <3>2 and <3>3._
149//!     3. Q.E.D
150//!         * _Proof: From cases <2>1 and <2>2 which are exhaustive_
151//! 6. **PICK** t2_prev ∈ IntoTime : t2_prev = Predecessor(t2)
152//!    * _Proof: Predecessor exists because by <1>4 the set {t: t < t2} is non empty since it must contain at least t1._
153//! 7. t1 ≀ t2_prev
154//!    * _Proof: t1 ∈ {t: t < t2} and t2_prev is the maximum element of the set._
155//! 8. 𝐑(t2) βͺ― s
156//!     1. t2 > min(IntoTime)
157//!         * _Proof: By <1>5._
158//!     2. **PICK** t2_prev = Predecessor(t2)
159//!         * _Proof: Predecessor exists because the set {t: t < t2} is non-empty since it must contain at least min(IntoTime)._
160//!     3. δ𝐓(t) = sum{δ𝐒(s): (𝐑(t2_prev) βͺ― s) ∧ !(𝐑(t) βͺ― s)}
161//!         * _Proof: By definition of δ𝐓(t)_
162//!     4. δ𝐒(s) ∈ δ𝐓(t1)
163//!         * _Proof: By <1>4._
164//!     5. Q.E.D
165//!         * _Proof: By <2>3 and <2>4._
166//! 9. 𝐑(t1) βͺ― 𝐑(t2_prev)
167//!     * _Proof: By <1>.7 and hypothesis on R_
168//! 10. 𝐑(t1) βͺ― s
169//!     * _Proof: By <1>8 and <1>9._
170//! 11. Q.E.D
171//!     * _Proof: By <1>5 and <1>10_
172
173use std::cmp::{Ordering, Reverse};
174use std::collections::VecDeque;
175use std::collections::binary_heap::{BinaryHeap, PeekMut};
176use std::iter::FromIterator;
177
178use differential_dataflow::difference::Semigroup;
179use differential_dataflow::lattice::Lattice;
180use differential_dataflow::{AsCollection, ExchangeData, VecCollection, consolidation};
181use mz_ore::Overflowing;
182use mz_ore::collections::CollectionExt;
183use timely::communication::{Pull, Push};
184use timely::dataflow::channels::pact::Pipeline;
185use timely::dataflow::operators::CapabilitySet;
186use timely::dataflow::operators::capture::Event;
187use timely::dataflow::operators::generic::OutputBuilder;
188use timely::dataflow::operators::generic::builder_rc::OperatorBuilder;
189use timely::order::{PartialOrder, TotalOrder};
190use timely::progress::frontier::{AntichainRef, MutableAntichain};
191use timely::progress::{Antichain, Timestamp};
192
193/// Constructs an operator that reclocks a `source` collection varying with some time `FromTime`
194/// into the corresponding `reclocked` collection varying over some time `IntoTime` using the
195/// provided `remap` collection.
196///
197/// In order for the operator to read the `source` collection a `Pusher` is returned which can be
198/// used with timely's capture facilities to connect a collection from a foreign scope to this
199/// operator.
200pub fn reclock<'scope, D, FromTime, IntoTime, R>(
201    remap_collection: VecCollection<'scope, IntoTime, FromTime, Overflowing<i64>>,
202    as_of: Antichain<IntoTime>,
203) -> (
204    Box<dyn Push<Event<FromTime, Vec<(D, FromTime, R)>>>>,
205    VecCollection<'scope, IntoTime, D, R>,
206)
207where
208    D: ExchangeData,
209    FromTime: Timestamp,
210    IntoTime: Timestamp + Lattice + TotalOrder,
211    R: Semigroup + 'static,
212{
213    let scope = remap_collection.scope();
214    let mut builder = OperatorBuilder::new("Reclock".into(), scope.clone());
215    // Here we create a channel that can be used to send data from a foreign scope into this
216    // operator. The channel is associated with this operator's address so that it is activated
217    // every time events are available for consumption. This mechanism is similar to Timely's input
218    // handles where data can be introduced into a timely scope from an exogenous source.
219    let info = builder.operator_info();
220    let channel_id = scope.worker().new_identifier();
221    let (pusher, mut events) = scope
222        .worker()
223        .pipeline::<Event<FromTime, Vec<(D, FromTime, R)>>>(channel_id, info.address);
224
225    let mut remap_input = builder.new_input(remap_collection.inner, Pipeline);
226    // Keep the output disconnected from the remap collection input so that we can drop the output
227    // capability when we finish reclocking the source, which can happen before the
228    // remap_collection reaches the input frontier. See the `test_finalized_source` for an example.
229    let (output, reclocked) = builder.new_output_connection([]);
230    let mut output = OutputBuilder::from(output);
231
232    builder.build(move |caps| {
233        let mut capset = CapabilitySet::from_elem(caps.into_element());
234        capset.downgrade(&as_of.borrow());
235
236        // Received remap updates at times `into_time` greater or equal to `remap_input`'s input
237        // frontier. As the input frontier advances, we drop elements out of this priority queue
238        // and mint new associations.
239        let mut pending_remap: BinaryHeap<Reverse<(IntoTime, FromTime, i64)>> = BinaryHeap::new();
240        // A trace of `remap_input` that accumulates correctly for all times that are beyond
241        // `remap_since` and not beyond `remap_upper`. The updates in `remap_trace` are maintained
242        // in time order. An actual DD trace could be used here at the expense of a more
243        // complicated API to traverse it. This is left for future work if the naive trace
244        // maintenance implemented in this operator becomes problematic.
245        let mut remap_upper = Antichain::from_elem(IntoTime::minimum());
246        let mut remap_since = as_of.clone();
247        let mut remap_trace = VecDeque::new();
248
249        // A stash of source updates for which we don't know the corresponding binding yet.
250        let mut deferred_source_updates: Vec<ChainBatch<_, _, _>> = Vec::new();
251        // The frontier of the `events` input
252        let mut source_frontier = MutableAntichain::from_elem(FromTime::minimum());
253
254        let mut binding_buffer = Vec::new();
255
256        // Accumulation buffer for `remap_input` updates.
257        use timely::progress::ChangeBatch;
258        let mut remap_accum_buffer: ChangeBatch<(IntoTime, FromTime)> = ChangeBatch::new();
259
260        // The operator drains `remap_input` and organizes new bindings that are not beyond
261        // `remap_input`'s frontier into the time ordered `remap_trace`.
262        //
263        // All received data events can either be reclocked to a time included in the
264        // `remap_trace`, or deferred until new associations are minted. Each data event that
265        // happens at some `FromTime` is mapped to the first `IntoTime` whose associated antichain
266        // is not less or equal to the input `FromTime`.
267        //
268        // As progress events are received from the `events` input, we can advance our
269        // held capability to track the least `IntoTime` a newly received `FromTime` could possibly
270        // map to and also compact the maintained `remap_trace` to that time.
271        move |frontiers| {
272            let Some(cap) = capset.get(0).cloned() else {
273                remap_input.for_each(|_, _| {});
274                return;
275            };
276            let mut output = output.activate();
277            let mut session = output.session(&cap);
278
279            // STEP 1. Accept new bindings into `pending_remap`.
280            // Advance all `into` times by `as_of`, and consolidate all updates at that frontier.
281            remap_input.for_each(|_, data| {
282                for (from, mut into, diff) in data.drain(..) {
283                    into.advance_by(as_of.borrow());
284                    remap_accum_buffer.update((into, from), diff.into_inner());
285                }
286            });
287            // Drain consolidated bindings into the `pending_remap` heap.
288            // Only do this once any of the `remap_input` frontier has passed `as_of`.
289            // For as long as the input frontier is less-equal `as_of`, we have no finalized times.
290            if !PartialOrder::less_equal(&frontiers[0].frontier(), &as_of.borrow()) {
291                for ((into, from), diff) in remap_accum_buffer.drain() {
292                    pending_remap.push(Reverse((into, from, diff)));
293                }
294            }
295
296            // STEP 2. Extract bindings not beyond `remap_frontier` and commit them into `remap_trace`.
297            let prev_remap_upper =
298                std::mem::replace(&mut remap_upper, frontiers[0].frontier().to_owned());
299            while let Some(update) = pending_remap.peek_mut() {
300                if !remap_upper.less_equal(&update.0.0) {
301                    let Reverse((into, from, diff)) = PeekMut::pop(update);
302                    remap_trace.push_back((from, into, diff));
303                } else {
304                    break;
305                }
306            }
307
308            // STEP 3. Receive new data updates
309            //         The `events` input describes arbitrary progress and data over `FromTime`,
310            //         which must be translated to `IntoTime`. Each `FromTime` can be found as the
311            //         first `IntoTime` associated with a `[FromTime]` that is not less or equal to
312            //         the input `FromTime`. Received events that are not yet associated to an
313            //         `IntoTime` are collected, and formed into a "chain batch": a sequence of
314            //         chains that results from sorting the updates by `FromTime`, and then
315            //         segmenting the sequence at elements where the partial order on `FromTime` is
316            //         violated.
317            let mut stash = Vec::new();
318            // Consolidate progress updates before applying them to `source_frontier`, to avoid quadratic
319            // behavior in overload scenarios.
320            let mut change_batch = ChangeBatch::<FromTime, 2>::default();
321            while let Some(event) = events.pull() {
322                match event {
323                    Event::Progress(changes) => {
324                        change_batch.extend(changes.drain(..));
325                    }
326                    Event::Messages(_, data) => stash.append(data),
327                }
328            }
329            source_frontier.update_iter(change_batch.drain());
330            stash.sort_unstable_by(|(_, t1, _): &(D, FromTime, R), (_, t2, _)| t1.cmp(t2));
331            let mut new_source_updates = ChainBatch::from_iter(stash);
332
333            // STEP 4: Reclock new and deferred updates
334            //         We are now ready to step through the remap bindings in time order and
335            //         perform the following actions:
336            //         4.1. Match `new_source_updates` against the entirety of bindings contained
337            //              in the trace.
338            //         4.2. Match `deferred_source_updates` against the bindings that were just
339            //              added in the trace.
340            //         4.3. Reclock `source_frontier` to calculate the new since frontier of the
341            //              remap trace.
342            //
343            //         The steps above only make sense to perform if there are any times for which
344            //         we can correctly accumulate the remap trace, which is what we check here.
345            if remap_since.iter().all(|t| !remap_upper.less_equal(t)) {
346                let mut cur_binding = MutableAntichain::new();
347
348                let mut remap = remap_trace.iter().peekable();
349                let mut reclocked_source_frontier = remap_upper.clone();
350
351                // We go over all the times for which we might need to output data at. These times
352                // are restrticted to the times at which there exists an update in `remap_trace`
353                // and the minimum timestamp for the case where `remap_trace` is completely empty,
354                // in which case the minimum timestamp maps to the empty `FromTime` frontier and
355                // therefore all data events map to that minimum timestamp.
356                //
357                // The approach taken here will take time proportional to the number of elements in
358                // `remap_trace`. During development an alternative approach was considered where
359                // the updates in `remap_trace` are instead fully materialized into an ordered list
360                // of antichains in which every data update can be binary searched into. The are
361                // two concerns with this alternative approach that led to preferring this one:
362                // 1. Materializing very wide antichains with small differences between them
363                //    needs memory proportial to the number of bindings times the width of the
364                //    antichain.
365                // 2. It locks in the requirement of a totally ordered target timestamp since only
366                //    in that case can one binary search a binding.
367                // The linear scan is expected to be fine due to the run-to-completion nature of
368                // the operator since its cost is amortized among the number of outstanding
369                // updates.
370                let mut min_time = IntoTime::minimum();
371                min_time.advance_by(remap_since.borrow());
372                let mut prev_cur_time = None;
373                let mut interesting_times = std::iter::once(&min_time)
374                    .chain(remap_trace.iter().map(|(_, t, _)| t))
375                    .filter(|&v| {
376                        let prev = prev_cur_time.replace(v);
377                        prev != prev_cur_time
378                    });
379                let mut frontier_reclocked = false;
380                while !(new_source_updates.is_empty()
381                    && deferred_source_updates.is_empty()
382                    && frontier_reclocked)
383                    && let Some(cur_time) = interesting_times.next()
384                {
385                    // 4.0. Load updates of `cur_time` from the trace into `cur_binding` to
386                    //      construct the `[FromTime]` frontier that `cur_time` maps to.
387                    while let Some((t_from, _, diff)) = remap.next_if(|(_, t, _)| t == cur_time) {
388                        binding_buffer.push((t_from.clone(), *diff));
389                    }
390                    cur_binding.update_iter(binding_buffer.drain(..));
391                    let cur_binding = cur_binding.frontier();
392
393                    // 4.1. Extract updates from `new_source_updates`
394                    for (data, _, diff) in new_source_updates.extract(cur_binding) {
395                        session.give((data, cur_time.clone(), diff));
396                    }
397
398                    // 4.2. Extract updates from `deferred_source_updates`.
399                    //      The deferred updates contain all updates that were not able to be
400                    //      reclocked with the bindings until `prev_remap_upper`. For this reason
401                    //      we only need to reconsider these updates when we start looking at new
402                    //      bindings, i.e bindings that are beyond `prev_remap_upper`.
403                    if prev_remap_upper.less_equal(cur_time) {
404                        deferred_source_updates.retain_mut(|batch| {
405                            for (data, _, diff) in batch.extract(cur_binding) {
406                                session.give((data, cur_time.clone(), diff));
407                            }
408                            // Retain non-empty batches
409                            !batch.is_empty()
410                        })
411                    }
412
413                    // 4.3. Reclock `source_frontier`
414                    //      If any FromTime in source frontier could possibly be reclocked to this
415                    //      binding then we must maintain our capability to emit data at that time
416                    //      and not compact past it. Since we iterate over this loop in time order
417                    //      and IntoTime is a total order we only need to perform this step once.
418                    //      Once a `cur_time` is inserted into `reclocked_source_frontier` no more
419                    //      changes can be made to the frontier by inserting times later in the
420                    //      loop.
421                    if !frontier_reclocked
422                        && source_frontier
423                            .frontier()
424                            .iter()
425                            .any(|t| !cur_binding.less_equal(t))
426                    {
427                        reclocked_source_frontier.insert(cur_time.clone());
428                        frontier_reclocked = true;
429                    }
430                }
431
432                // STEP 5. Downgrade capability and compact remap trace if our since frontier
433                //         advanced
434                if PartialOrder::less_than(&remap_since, &reclocked_source_frontier) {
435                    capset.downgrade(&reclocked_source_frontier.borrow());
436                    remap_since = reclocked_source_frontier;
437
438                    // The remap trace is stored in time order and T is a total order. The updates
439                    // that can have their time advanced will form a prefix, which we extract here.
440                    let mut advanced = vec![];
441                    while !remap_trace.is_empty() && !remap_since.less_equal(&remap_trace[0].1) {
442                        let (d, mut t, r) = remap_trace.pop_front().unwrap();
443                        t.advance_by(remap_since.borrow());
444                        advanced.push((d, t, r));
445                    }
446                    if !advanced.is_empty() {
447                        // If we have updates whose time was advanced, further peel the prefix of
448                        // updates that sit *at* the time of the since frontier.
449                        while !remap_trace.is_empty() && remap_since.contains(&remap_trace[0].1) {
450                            advanced.push(remap_trace.pop_front().unwrap());
451                        }
452                        consolidation::consolidate_updates(&mut advanced);
453                        advanced.sort_unstable_by(|(_, t1, _): &(_, IntoTime, _), (_, t2, _)| {
454                            t1.cmp(t2)
455                        });
456                        for u in advanced.into_iter().rev() {
457                            remap_trace.push_front(u);
458                        }
459                        // If using less than a quarter of the capacity, shrink the container. To avoid having
460                        // to resize the container on a subsequent push, shrink to 2x the length, which is
461                        // what push would grow it to.
462                        if remap_trace.len() < remap_trace.capacity() / 4 {
463                            remap_trace.shrink_to(remap_trace.len() * 2);
464                        }
465                    }
466                }
467            }
468
469            // STEP 6. Tidy up deferred updates
470            //         Deferred updates are represented as a list of chain batches where each batch
471            //         contains two times the updates of the batch proceeding it. This organization
472            //         leads to a logarithmic number of batches with respect to the outstanding
473            //         number of updates.
474            deferred_source_updates.sort_unstable_by_key(|b| Reverse(b.len()));
475            if !new_source_updates.is_empty() {
476                deferred_source_updates.push(new_source_updates);
477            }
478            let dsu = &mut deferred_source_updates;
479            while dsu.len() > 1 && (dsu[dsu.len() - 1].len() >= dsu[dsu.len() - 2].len() / 2) {
480                let a = dsu.pop().unwrap();
481                let b = dsu.pop().unwrap();
482                dsu.push(a.merge_with(b));
483            }
484
485            // If using less than a quarter of the capacity, shrink the container. To avoid having
486            // to resize the container on a subsequent push, shrink to 2x the length, which is
487            // what push would grow it to.
488            if deferred_source_updates.len() < deferred_source_updates.capacity() / 4 {
489                deferred_source_updates.shrink_to(deferred_source_updates.len() * 2);
490            }
491
492            // If this condition holds then this operator has no more work to do, so we drop the capability.
493            if source_frontier.frontier().is_empty() && deferred_source_updates.is_empty() {
494                capset = CapabilitySet::new();
495            }
496        }
497    });
498
499    (Box::new(pusher), reclocked.as_collection())
500}
501
502/// A batch of differential updates that vary over some partial order. This type maintains the data
503/// as a set of chains that allows for efficient extraction of batches given a frontier.
504#[derive(Debug, PartialEq)]
505struct ChainBatch<D, T, R> {
506    /// A list of chains (sets of mutually comparable times) sorted by the partial order.
507    chains: Vec<VecDeque<(D, T, R)>>,
508}
509
510impl<D, T: Timestamp, R> ChainBatch<D, T, R> {
511    /// Extracts all updates with time not greater or equal to any time in `upper`.
512    fn extract<'a>(
513        &'a mut self,
514        upper: AntichainRef<'a, T>,
515    ) -> impl Iterator<Item = (D, T, R)> + 'a {
516        self.chains.retain(|chain| !chain.is_empty());
517        self.chains.iter_mut().flat_map(move |chain| {
518            // A chain is a sorted list of mutually comparable elements so we keep extracting
519            // elements that are not beyond upper.
520            std::iter::from_fn(move || {
521                let (_, into, _) = chain.front()?;
522                if !upper.less_equal(into) {
523                    chain.pop_front()
524                } else {
525                    None
526                }
527            })
528        })
529    }
530
531    fn merge_with(
532        mut self: ChainBatch<D, T, R>,
533        mut other: ChainBatch<D, T, R>,
534    ) -> ChainBatch<D, T, R>
535    where
536        D: ExchangeData,
537        T: Timestamp,
538        R: Semigroup,
539    {
540        let mut updates1 = self.chains.drain(..).flatten().peekable();
541        let mut updates2 = other.chains.drain(..).flatten().peekable();
542
543        let merged = std::iter::from_fn(|| {
544            match (updates1.peek(), updates2.peek()) {
545                (Some((d1, t1, _)), Some((d2, t2, _))) => {
546                    match (t1, d1).cmp(&(t2, d2)) {
547                        Ordering::Less => updates1.next(),
548                        Ordering::Greater => updates2.next(),
549                        // If the same (d, t) pair is found, consolidate their diffs
550                        Ordering::Equal => {
551                            let (d1, t1, mut r1) = updates1.next().unwrap();
552                            while let Some((_, _, r)) =
553                                updates1.next_if(|(d, t, _)| (d, t) == (&d1, &t1))
554                            {
555                                r1.plus_equals(&r);
556                            }
557                            while let Some((_, _, r)) =
558                                updates2.next_if(|(d, t, _)| (d, t) == (&d1, &t1))
559                            {
560                                r1.plus_equals(&r);
561                            }
562                            Some((d1, t1, r1))
563                        }
564                    }
565                }
566                (Some(_), None) => updates1.next(),
567                (None, Some(_)) => updates2.next(),
568                (None, None) => None,
569            }
570        });
571
572        ChainBatch::from_iter(merged.filter(|(_, _, r)| !r.is_zero()))
573    }
574
575    /// Returns the number of updates in the batch.
576    fn len(&self) -> usize {
577        self.chains.iter().map(|chain| chain.len()).sum()
578    }
579
580    /// Returns true if the batch contains no updates.
581    fn is_empty(&self) -> bool {
582        self.len() == 0
583    }
584}
585
586impl<D, T: Timestamp, R> FromIterator<(D, T, R)> for ChainBatch<D, T, R> {
587    /// Computes the chain decomposition of updates according to the partial order `T`.
588    fn from_iter<I: IntoIterator<Item = (D, T, R)>>(updates: I) -> Self {
589        let mut chains = vec![];
590        let mut updates = updates.into_iter();
591        if let Some((d, t, r)) = updates.next() {
592            let mut chain = VecDeque::new();
593            chain.push_back((d, t, r));
594            for (d, t, r) in updates {
595                let prev_t = &chain[chain.len() - 1].1;
596                if !PartialOrder::less_equal(prev_t, &t) {
597                    chains.push(chain);
598                    chain = VecDeque::new();
599                }
600                chain.push_back((d, t, r));
601            }
602            chains.push(chain);
603        }
604        Self { chains }
605    }
606}
607
608#[cfg(test)]
609mod test {
610    use std::sync::atomic::AtomicUsize;
611    use std::sync::mpsc::{Receiver, TryRecvError};
612
613    use differential_dataflow::consolidation;
614    use differential_dataflow::input::{Input, InputSession};
615    use serde::{Deserialize, Serialize};
616    use timely::dataflow::operators::capture::{Event, Extract};
617    use timely::dataflow::operators::vec::UnorderedInput;
618    use timely::dataflow::operators::vec::unordered_input::UnorderedHandle;
619    use timely::dataflow::operators::{ActivateCapability, Capture};
620    use timely::progress::PathSummary;
621    use timely::progress::timestamp::Refines;
622    use timely::worker::Worker;
623
624    use crate::capture::PusherCapture;
625    use crate::order::Partitioned;
626
627    use super::*;
628
629    type Diff = Overflowing<i64>;
630    type FromTime = Partitioned<u64, u64>;
631    type IntoTime = u64;
632    type BindingHandle<FromTime> = InputSession<IntoTime, FromTime, Diff>;
633    type DataHandle<D, FromTime> = (
634        UnorderedHandle<FromTime, (D, FromTime, Diff)>,
635        ActivateCapability<FromTime>,
636    );
637    type ReclockedStream<D> = Receiver<Event<IntoTime, Vec<(D, IntoTime, Diff)>>>;
638
639    /// A helper function that sets up a dataflow program to test the reclocking operator. Each
640    /// test provides a test logic closure which accepts four arguments:
641    ///
642    /// * A reference to the worker that allows the test to step the computation
643    /// * A [`BindingHandle`] that allows the test to manipulate the remap bindings
644    /// * A [`DataHandle`] that allows the test to submit the data to be reclocked
645    /// * A [`ReclockedStream`] that allows observing the result of the reclocking process
646    ///
647    /// Note that the `DataHandle` contains a capability that should be dropped or downgraded before
648    /// calling [`step`] to process data at the time.
649    fn harness<FromTime, D, F, R>(as_of: Antichain<IntoTime>, test_logic: F) -> R
650    where
651        FromTime: Timestamp + Refines<()>,
652        D: ExchangeData,
653        F: FnOnce(
654                &mut Worker,
655                BindingHandle<FromTime>,
656                DataHandle<D, FromTime>,
657                ReclockedStream<D>,
658            ) -> R
659            + Send
660            + Sync
661            + 'static,
662        R: Send + 'static,
663    {
664        timely::execute_directly(move |worker| {
665            let (bindings, data, data_cap, reclocked) = worker.dataflow::<(), _, _>(|scope| {
666                let (bindings, data_pusher, reclocked) =
667                    scope.scoped::<IntoTime, _, _>("IntoScope", move |scope| {
668                        let (binding_handle, binding_collection) = scope.new_collection();
669                        let (data_pusher, reclocked_collection) =
670                            reclock(binding_collection, as_of);
671                        let reclocked_capture = reclocked_collection.inner.capture();
672                        (binding_handle, data_pusher, reclocked_capture)
673                    });
674
675                let (data, data_cap) = scope.scoped::<FromTime, _, _>("FromScope", move |scope| {
676                    let ((handle, cap), data) = scope.new_unordered_input::<(D, FromTime, Diff)>();
677                    data.capture_into(PusherCapture(data_pusher));
678                    (handle, cap)
679                });
680
681                (bindings, data, data_cap, reclocked)
682            });
683
684            test_logic(worker, bindings, (data, data_cap), reclocked)
685        })
686    }
687
688    /// Steps the worker four times which is the required number of times for both data and
689    /// frontier updates to propagate across the two scopes and into the probing channels.
690    fn step(worker: &mut Worker) {
691        for _ in 0..4 {
692            worker.step();
693        }
694    }
695
696    #[mz_ore::test]
697    fn basic_reclocking() {
698        let as_of = Antichain::from_elem(IntoTime::minimum());
699        harness::<FromTime, _, _, _>(
700            as_of,
701            |worker, bindings, (mut data, data_cap), reclocked| {
702                // Reclock everything at the minimum IntoTime
703                bindings.close();
704                data.activate()
705                    .session(&data_cap)
706                    .give(('a', Partitioned::minimum(), Diff::ONE));
707                drop(data_cap);
708                step(worker);
709                let extracted = reclocked.extract();
710                let expected = vec![(0, vec![('a', 0, Diff::ONE)])];
711                assert_eq!(extracted, expected);
712            },
713        )
714    }
715
716    /// Generates a `Partitioned<u64, u64>` Antichain where all the provided
717    /// partitions are at the specified offset and the gaps in between are filled with range
718    /// timestamps at offset zero.
719    fn partitioned_frontier<I>(items: I) -> Antichain<Partitioned<u64, u64>>
720    where
721        I: IntoIterator<Item = (u64, u64)>,
722    {
723        let mut frontier = Antichain::new();
724        let mut prev = 0;
725        for (pid, offset) in items {
726            if prev < pid {
727                frontier.insert(Partitioned::new_range(prev, pid - 1, 0));
728            }
729            frontier.insert(Partitioned::new_singleton(pid, offset));
730            prev = pid + 1
731        }
732        frontier.insert(Partitioned::new_range(prev, u64::MAX, 0));
733        frontier
734    }
735
736    #[mz_ore::test]
737    fn test_basic_usage() {
738        let as_of = Antichain::from_elem(IntoTime::minimum());
739        harness(
740            as_of,
741            |worker, mut bindings, (mut data, data_cap), reclocked| {
742                // Reclock offsets 1 and 3 to timestamp 1000
743                bindings.update_at(Partitioned::minimum(), 0, Diff::ONE);
744                bindings.update_at(Partitioned::minimum(), 1000, Diff::MINUS_ONE);
745                for time in partitioned_frontier([(0, 4)]) {
746                    bindings.update_at(time, 1000, Diff::ONE);
747                }
748                bindings.advance_to(1001);
749                bindings.flush();
750                data.activate().session(&data_cap).give_iterator(
751                    vec![
752                        (1, Partitioned::new_singleton(0, 1), Diff::ONE),
753                        (1, Partitioned::new_singleton(0, 1), Diff::ONE),
754                        (3, Partitioned::new_singleton(0, 3), Diff::ONE),
755                    ]
756                    .into_iter(),
757                );
758
759                step(worker);
760                assert_eq!(
761                    reclocked.try_recv(),
762                    Ok(Event::Messages(
763                        0u64,
764                        vec![
765                            (1, 1000, Diff::ONE),
766                            (1, 1000, Diff::ONE),
767                            (3, 1000, Diff::ONE)
768                        ]
769                    ))
770                );
771                assert_eq!(
772                    reclocked.try_recv(),
773                    Ok(Event::Progress(vec![(0, -1), (1000, 1)]))
774                );
775
776                // Reclock more messages for offsets 3 to the same timestamp
777                data.activate().session(&data_cap).give_iterator(
778                    vec![
779                        (3, Partitioned::new_singleton(0, 3), Diff::ONE),
780                        (3, Partitioned::new_singleton(0, 3), Diff::ONE),
781                    ]
782                    .into_iter(),
783                );
784                step(worker);
785                assert_eq!(
786                    reclocked.try_recv(),
787                    Ok(Event::Messages(
788                        1000u64,
789                        vec![(3, 1000, Diff::ONE), (3, 1000, Diff::ONE)]
790                    ))
791                );
792
793                // Drop the capability which should advance the reclocked frontier to 1001.
794                drop(data_cap);
795                step(worker);
796                assert_eq!(reclocked.try_recv(), Ok(Event::Progress(vec![(1000, -1)])));
797            },
798        );
799    }
800
801    // Test that once the input stream has reached the empty frontier with no pending data we
802    // eagerly downgrade the output frontier to the empty frontier regardless of the remap stream.
803    #[mz_ore::test]
804    fn test_finalized_source() {
805        let as_of = Antichain::from_elem(IntoTime::minimum());
806        harness(
807            as_of,
808            |worker, mut bindings, (mut data, data_cap), reclocked| {
809                // The starts by producing data in the source stream and advancing the frontier to
810                // the empty frontier. The reclock operator buffers the data and holds the
811                // appropriate capability to reclock the pending data in the future.
812                data.activate().session(&data_cap).give_iterator(
813                    vec![
814                        (1, Partitioned::new_singleton(0, 1), Diff::ONE),
815                        (2, Partitioned::new_singleton(0, 2), Diff::ONE),
816                    ]
817                    .into_iter(),
818                );
819                drop(data_cap);
820                step(worker);
821
822                // Reclock offset 0 and 1 to timestamp 1000
823                bindings.update_at(Partitioned::minimum(), 0, Diff::ONE);
824                bindings.update_at(Partitioned::minimum(), 1000, Diff::MINUS_ONE);
825                for time in partitioned_frontier([(0, 2)]) {
826                    bindings.update_at(time, 1000, Diff::ONE);
827                }
828                bindings.advance_to(1001);
829                bindings.flush();
830                step(worker);
831
832                // At this point the reclock operator is able to release one of the updates and
833                // downgrade the output capability to timestamp 1001. It can't yet downgrade to the
834                // empty frontier since the other update is still pending.
835                assert_eq!(
836                    reclocked.try_recv(),
837                    Ok(Event::Messages(0u64, vec![(1, 1000, Diff::ONE),]))
838                );
839                assert_eq!(
840                    reclocked.try_recv(),
841                    Ok(Event::Progress(vec![(0, -1), (1001, 1)]))
842                );
843
844                // Reclock offset 2 to timestamp 2000
845                for time in partitioned_frontier([(0, 2)]) {
846                    bindings.update_at(time, 2000, Diff::MINUS_ONE);
847                }
848                for time in partitioned_frontier([(0, 3)]) {
849                    bindings.update_at(time, 2000, Diff::ONE);
850                }
851                bindings.advance_to(2001);
852                bindings.flush();
853                step(worker);
854
855                // Now the reclock operator reclocks the last data item and since the source
856                // frontier is empty there will be no more work to be done, so the output frontier
857                // downgrades to the empty frontier even though the frontier of the remap bindings
858                // is still at [2001].
859                assert_eq!(
860                    reclocked.try_recv(),
861                    Ok(Event::Messages(1001u64, vec![(2, 2000, Diff::ONE),]))
862                );
863                assert_eq!(reclocked.try_recv(), Ok(Event::Progress(vec![(1001, -1)])));
864            },
865        );
866    }
867
868    #[mz_ore::test]
869    fn test_reclock_frontier() {
870        let as_of = Antichain::from_elem(IntoTime::minimum());
871        harness::<_, (), _, _>(
872            as_of,
873            |worker, mut bindings, (_data, data_cap), reclocked| {
874                // Initialize the bindings such that the minimum IntoTime contains the minimum FromTime
875                // frontier.
876                bindings.update_at(Partitioned::minimum(), 0, Diff::ONE);
877                bindings.advance_to(1);
878                bindings.flush();
879                step(worker);
880                assert_eq!(
881                    reclocked.try_recv(),
882                    Ok(Event::Progress(vec![(0, -1), (1, 1)]))
883                );
884
885                // Mint a couple of bindings for multiple partitions
886                bindings.update_at(Partitioned::minimum(), 1000, Diff::MINUS_ONE);
887                for time in partitioned_frontier([(1, 10)]) {
888                    bindings.update_at(time.clone(), 1000, Diff::ONE);
889                    bindings.update_at(time, 2000, Diff::MINUS_ONE);
890                }
891                for time in partitioned_frontier([(1, 10), (2, 10)]) {
892                    bindings.update_at(time, 2000, Diff::ONE);
893                }
894                bindings.advance_to(2001);
895                bindings.flush();
896
897                // The initial frontier should now map to the minimum between the two partitions
898                step(worker);
899                step(worker);
900                assert_eq!(
901                    reclocked.try_recv(),
902                    Ok(Event::Progress(vec![(1, -1), (1000, 1)]))
903                );
904
905                // Downgrade data frontier such that only one of the partitions is advanced
906                let mut part1_cap = data_cap.delayed(&Partitioned::new_singleton(1, 9));
907                let mut part2_cap = data_cap.delayed(&Partitioned::new_singleton(2, 0));
908                let _rest_cap = data_cap.delayed(&Partitioned::new_range(3, u64::MAX, 0));
909                drop(data_cap);
910                step(worker);
911                assert_eq!(reclocked.try_recv(), Err(TryRecvError::Empty));
912
913                // Downgrade the data frontier past the first binding
914                part1_cap.downgrade(&Partitioned::new_singleton(1, 10));
915                step(worker);
916                assert_eq!(
917                    reclocked.try_recv(),
918                    Ok(Event::Progress(vec![(1000, -1), (2000, 1)]))
919                );
920
921                // Downgrade the data frontier past the second binding
922                part2_cap.downgrade(&Partitioned::new_singleton(2, 10));
923                step(worker);
924                assert_eq!(
925                    reclocked.try_recv(),
926                    Ok(Event::Progress(vec![(2000, -1), (2001, 1)]))
927                );
928
929                // Advance the binding frontier and confirm that we get to the next timestamp
930                bindings.advance_to(3001);
931                bindings.flush();
932                step(worker);
933                assert_eq!(
934                    reclocked.try_recv(),
935                    Ok(Event::Progress(vec![(2001, -1), (3001, 1)]))
936                );
937            },
938        );
939    }
940
941    #[mz_ore::test]
942    fn test_reclock() {
943        let as_of = Antichain::from_elem(IntoTime::minimum());
944        harness(
945            as_of,
946            |worker, mut bindings, (mut data, data_cap), reclocked| {
947                // Initialize the bindings such that the minimum IntoTime contains the minimum FromTime
948                // frontier.
949                bindings.update_at(Partitioned::minimum(), 0, Diff::ONE);
950
951                // Setup more precise capabilities for the rest of the test
952                let mut part0_cap = data_cap.delayed(&Partitioned::new_singleton(0, 0));
953                let rest_cap = data_cap.delayed(&Partitioned::new_range(1, u64::MAX, 0));
954                drop(data_cap);
955
956                // Reclock offsets 1 and 2 to timestamp 1000
957                data.activate().session(&part0_cap).give_iterator(
958                    vec![
959                        (1, Partitioned::new_singleton(0, 1), Diff::ONE),
960                        (2, Partitioned::new_singleton(0, 2), Diff::ONE),
961                    ]
962                    .into_iter(),
963                );
964
965                part0_cap.downgrade(&Partitioned::new_singleton(0, 3));
966                bindings.update_at(Partitioned::minimum(), 1000, Diff::MINUS_ONE);
967                bindings.update_at(part0_cap.time().clone(), 1000, Diff::ONE);
968                bindings.update_at(rest_cap.time().clone(), 1000, Diff::ONE);
969                bindings.advance_to(1001);
970                bindings.flush();
971                step(worker);
972                assert_eq!(
973                    reclocked.try_recv(),
974                    Ok(Event::Messages(
975                        0,
976                        vec![(1, 1000, Diff::ONE), (2, 1000, Diff::ONE)]
977                    ))
978                );
979                assert_eq!(
980                    reclocked.try_recv(),
981                    Ok(Event::Progress(vec![(0, -1), (1000, 1)]))
982                );
983                assert_eq!(
984                    reclocked.try_recv(),
985                    Ok(Event::Progress(vec![(1000, -1), (1001, 1)]))
986                );
987
988                // Reclock offsets 3 and 4 to timestamp 2000
989                data.activate().session(&part0_cap).give_iterator(
990                    vec![
991                        (3, Partitioned::new_singleton(0, 3), Diff::ONE),
992                        (3, Partitioned::new_singleton(0, 3), Diff::ONE),
993                        (4, Partitioned::new_singleton(0, 4), Diff::ONE),
994                    ]
995                    .into_iter(),
996                );
997                bindings.update_at(part0_cap.time().clone(), 2000, Diff::MINUS_ONE);
998                part0_cap.downgrade(&Partitioned::new_singleton(0, 5));
999                bindings.update_at(part0_cap.time().clone(), 2000, Diff::ONE);
1000                bindings.advance_to(2001);
1001                bindings.flush();
1002                step(worker);
1003                assert_eq!(
1004                    reclocked.try_recv(),
1005                    Ok(Event::Messages(
1006                        1001,
1007                        vec![
1008                            (3, 2000, Diff::ONE),
1009                            (3, 2000, Diff::ONE),
1010                            (4, 2000, Diff::ONE)
1011                        ]
1012                    ))
1013                );
1014                assert_eq!(
1015                    reclocked.try_recv(),
1016                    Ok(Event::Progress(vec![(1001, -1), (2000, 1)]))
1017                );
1018                assert_eq!(
1019                    reclocked.try_recv(),
1020                    Ok(Event::Progress(vec![(2000, -1), (2001, 1)]))
1021                );
1022            },
1023        );
1024    }
1025
1026    #[mz_ore::test]
1027    fn test_reclock_gh16318() {
1028        let as_of = Antichain::from_elem(IntoTime::minimum());
1029        harness(
1030            as_of,
1031            |worker, mut bindings, (mut data, data_cap), reclocked| {
1032                // Initialize the bindings such that the minimum IntoTime contains the minimum FromTime
1033                // frontier.
1034                bindings.update_at(Partitioned::minimum(), 0, Diff::ONE);
1035                // First mint bindings for 0 at timestamp 1000
1036                bindings.update_at(Partitioned::minimum(), 1000, Diff::MINUS_ONE);
1037                for time in partitioned_frontier([(0, 50)]) {
1038                    bindings.update_at(time, 1000, Diff::ONE);
1039                }
1040                // Then only for 1 at timestamp 2000
1041                for time in partitioned_frontier([(0, 50)]) {
1042                    bindings.update_at(time, 2000, Diff::MINUS_ONE);
1043                }
1044                for time in partitioned_frontier([(0, 50), (1, 50)]) {
1045                    bindings.update_at(time, 2000, Diff::ONE);
1046                }
1047                // Then again only for 0 at timestamp 3000
1048                for time in partitioned_frontier([(0, 50), (1, 50)]) {
1049                    bindings.update_at(time, 3000, Diff::MINUS_ONE);
1050                }
1051                for time in partitioned_frontier([(0, 100), (1, 50)]) {
1052                    bindings.update_at(time, 3000, Diff::ONE);
1053                }
1054                bindings.advance_to(3001);
1055                bindings.flush();
1056
1057                // Reclockng (0, 50) must ignore the updates on the FromTime frontier that happened at
1058                // timestamp 2000 since those are completely unrelated
1059                data.activate().session(&data_cap).give((
1060                    50,
1061                    Partitioned::new_singleton(0, 50),
1062                    Diff::ONE,
1063                ));
1064                drop(data_cap);
1065                step(worker);
1066                assert_eq!(
1067                    reclocked.try_recv(),
1068                    Ok(Event::Messages(0, vec![(50, 3000, Diff::ONE),]))
1069                );
1070                assert_eq!(
1071                    reclocked.try_recv(),
1072                    Ok(Event::Progress(vec![(0, -1), (1000, 1)]))
1073                );
1074                assert_eq!(reclocked.try_recv(), Ok(Event::Progress(vec![(1000, -1)])));
1075            },
1076        );
1077    }
1078
1079    /// Test that compact(reclock(remap, source)) == reclock(compact(remap), source)
1080    #[mz_ore::test]
1081    fn test_compaction() {
1082        let mut remap = vec![];
1083        remap.push((Partitioned::minimum(), 0, Diff::ONE));
1084        // Reclock offsets 1 and 2 to timestamp 1000
1085        remap.push((Partitioned::minimum(), 1000, Diff::MINUS_ONE));
1086        for time in partitioned_frontier([(0, 3)]) {
1087            remap.push((time, 1000, Diff::ONE));
1088        }
1089        // Reclock offsets 3 and 4 to timestamp 2000
1090        for time in partitioned_frontier([(0, 3)]) {
1091            remap.push((time, 2000, Diff::MINUS_ONE));
1092        }
1093        for time in partitioned_frontier([(0, 5)]) {
1094            remap.push((time, 2000, Diff::ONE));
1095        }
1096
1097        let source_updates = vec![
1098            (1, Partitioned::new_singleton(0, 1), Diff::ONE),
1099            (2, Partitioned::new_singleton(0, 2), Diff::ONE),
1100            (3, Partitioned::new_singleton(0, 3), Diff::ONE),
1101            (4, Partitioned::new_singleton(0, 4), Diff::ONE),
1102        ];
1103
1104        let since = Antichain::from_elem(1500);
1105
1106        // Compute reclock(remap, source)
1107        let as_of = Antichain::from_elem(IntoTime::minimum());
1108        let remap1 = remap.clone();
1109        let source_updates1 = source_updates.clone();
1110        let reclock_remap = harness(
1111            as_of,
1112            move |worker, mut bindings, (mut data, data_cap), reclocked| {
1113                for (from_ts, into_ts, diff) in remap1 {
1114                    bindings.update_at(from_ts, into_ts, diff);
1115                }
1116                bindings.close();
1117                data.activate()
1118                    .session(&data_cap)
1119                    .give_iterator(source_updates1.iter().cloned());
1120                drop(data_cap);
1121                step(worker);
1122                reclocked.extract()
1123            },
1124        );
1125        // Compute compact(reclock(remap, source))
1126        let mut compact_reclock_remap = reclock_remap;
1127        for (t, updates) in compact_reclock_remap.iter_mut() {
1128            t.advance_by(since.borrow());
1129            for (_, t, _) in updates.iter_mut() {
1130                t.advance_by(since.borrow());
1131            }
1132        }
1133
1134        // Compute compact(remap)
1135        let mut compact_remap = remap;
1136        for (_, t, _) in compact_remap.iter_mut() {
1137            t.advance_by(since.borrow());
1138        }
1139        consolidation::consolidate_updates(&mut compact_remap);
1140        // Compute reclock(compact(remap), source)
1141        let reclock_compact_remap = harness(
1142            since,
1143            move |worker, mut bindings, (mut data, data_cap), reclocked| {
1144                for (from_ts, into_ts, diff) in compact_remap {
1145                    bindings.update_at(from_ts, into_ts, diff);
1146                }
1147                bindings.close();
1148                data.activate()
1149                    .session(&data_cap)
1150                    .give_iterator(source_updates.iter().cloned());
1151                drop(data_cap);
1152                step(worker);
1153                reclocked.extract()
1154            },
1155        );
1156
1157        let expected = vec![(
1158            1500,
1159            vec![
1160                (1, 1500, Diff::ONE),
1161                (2, 1500, Diff::ONE),
1162                (3, 2000, Diff::ONE),
1163                (4, 2000, Diff::ONE),
1164            ],
1165        )];
1166        assert_eq!(expected, reclock_compact_remap);
1167        assert_eq!(expected, compact_reclock_remap);
1168    }
1169
1170    #[mz_ore::test]
1171    fn test_chainbatch_merge() {
1172        let a = ChainBatch::from_iter([('a', 0, 1)]);
1173        let b = ChainBatch::from_iter([('a', 0, -1), ('a', 1, 1)]);
1174        assert_eq!(a.merge_with(b), ChainBatch::from_iter([('a', 1, 1)]));
1175    }
1176
1177    #[mz_ore::test]
1178    #[cfg_attr(miri, ignore)] // too slow
1179    fn test_binding_consolidation() {
1180        use std::sync::atomic::Ordering;
1181
1182        #[derive(Debug, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize, Deserialize)]
1183        struct Time(u64);
1184
1185        // A counter of the number of active Time instances
1186        static INSTANCES: AtomicUsize = AtomicUsize::new(0);
1187
1188        impl Time {
1189            fn new(time: u64) -> Self {
1190                INSTANCES.fetch_add(1, Ordering::Relaxed);
1191                Self(time)
1192            }
1193        }
1194
1195        impl Clone for Time {
1196            fn clone(&self) -> Self {
1197                INSTANCES.fetch_add(1, Ordering::Relaxed);
1198                Self(self.0)
1199            }
1200        }
1201
1202        impl Drop for Time {
1203            fn drop(&mut self) {
1204                INSTANCES.fetch_sub(1, Ordering::Relaxed);
1205            }
1206        }
1207
1208        impl Timestamp for Time {
1209            type Summary = ();
1210
1211            fn minimum() -> Self {
1212                Time::new(0)
1213            }
1214        }
1215
1216        impl PathSummary<Time> for () {
1217            fn results_in(&self, src: &Time) -> Option<Time> {
1218                Some(src.clone())
1219            }
1220
1221            fn followed_by(&self, _other: &()) -> Option<Self> {
1222                Some(())
1223            }
1224        }
1225
1226        impl Refines<()> for Time {
1227            fn to_inner(_: ()) -> Self {
1228                Self::minimum()
1229            }
1230            fn to_outer(self) -> () {}
1231            fn summarize(_path: ()) {}
1232        }
1233
1234        impl PartialOrder for Time {
1235            fn less_equal(&self, other: &Self) -> bool {
1236                self.0.less_equal(&other.0)
1237            }
1238        }
1239
1240        let as_of = 1000;
1241
1242        // Test that supplying a single big batch of unconsolidated bindings gets
1243        // consolidated after a single worker step.
1244        harness::<Time, u64, _, _>(
1245            Antichain::from_elem(as_of),
1246            move |worker, mut bindings, _, _| {
1247                step(worker);
1248                let instances_before = INSTANCES.load(Ordering::Relaxed);
1249                for ts in 0..as_of {
1250                    if ts > 0 {
1251                        bindings.update_at(Time::new(ts - 1), ts, Diff::MINUS_ONE);
1252                    }
1253                    bindings.update_at(Time::new(ts), ts, Diff::ONE);
1254                }
1255                bindings.advance_to(as_of);
1256                bindings.flush();
1257                step(worker);
1258                let instances_after = INSTANCES.load(Ordering::Relaxed);
1259                // The extra instances live in a ChangeBatch which considers compaction when more
1260                // than 32 elements are inside.
1261                assert!(instances_after - instances_before < 32);
1262            },
1263        );
1264
1265        // Test that a slow feed of uncompacted bindings over multiple steps never leads to an
1266        // excessive number of bindings held in memory.
1267        harness::<Time, u64, _, _>(
1268            Antichain::from_elem(as_of),
1269            move |worker, mut bindings, _, _| {
1270                step(worker);
1271                let instances_before = INSTANCES.load(Ordering::Relaxed);
1272                for ts in 0..as_of {
1273                    if ts > 0 {
1274                        bindings.update_at(Time::new(ts - 1), ts, Diff::MINUS_ONE);
1275                    }
1276                    bindings.update_at(Time::new(ts), ts, Diff::ONE);
1277                    bindings.advance_to(ts + 1);
1278                    bindings.flush();
1279                    step(worker);
1280                    let instances_now = INSTANCES.load(Ordering::Relaxed);
1281                    // The extra instances live in a ChangeBatch which considers compaction when
1282                    // more than 32 elements are inside.
1283                    assert!(instances_now - instances_before < 32);
1284                }
1285            },
1286        );
1287    }
1288
1289    #[cfg(feature = "count-allocations")]
1290    #[mz_ore::test]
1291    #[cfg_attr(miri, ignore)] // too slow
1292    fn test_shrinking() {
1293        let as_of = 1000_u64;
1294
1295        // This workflow accumulates updates in remap_trace, advances the source frontier,
1296        // and validates that memory was reclaimed.  To avoid errant test failures due to
1297        // optimizations, this only validates that memory is reclaimed, not how much.
1298        harness::<FromTime, u64, _, _>(
1299            Antichain::from_elem(0),
1300            move |worker, mut bindings, (_data, mut data_cap), _| {
1301                let info1 = allocation_counter::measure(|| {
1302                    step(worker);
1303                    for ts in 0..as_of {
1304                        if ts > 0 {
1305                            bindings.update_at(
1306                                Partitioned::new_singleton(0, ts - 1),
1307                                ts,
1308                                Diff::MINUS_ONE,
1309                            );
1310                        }
1311                        bindings.update_at(Partitioned::new_singleton(0, ts), ts, Diff::ONE);
1312                        bindings.advance_to(ts + 1);
1313                        bindings.flush();
1314                        step(worker);
1315                    }
1316                });
1317                println!("info = {info1:?}");
1318
1319                let info2 = allocation_counter::measure(|| {
1320                    data_cap.downgrade(&Partitioned::new_singleton(0, as_of));
1321                    step(worker);
1322                });
1323                println!("info = {info2:?}");
1324                assert!(info2.bytes_current < 0);
1325            },
1326        );
1327    }
1328}