Skip to main content

mz_storage/
upsert_continual_feedback.rs

1// Copyright Materialize, Inc. and contributors. All rights reserved.
2//
3// Use of this software is governed by the Business Source License
4// included in the LICENSE file.
5//
6// As of the Change Date specified in that file, in accordance with
7// the Business Source License, use of this software will be governed
8// by the Apache License, Version 2.0.
9
10//! Implementation of feedback UPSERT operator and associated helpers. See
11//! [`upsert_inner`] for a description of how the operator works and why.
12
13use std::cmp::Reverse;
14use std::fmt::Debug;
15use std::sync::Arc;
16
17use differential_dataflow::hashable::Hashable;
18use differential_dataflow::{AsCollection, VecCollection};
19use indexmap::map::Entry;
20use itertools::Itertools;
21use mz_repr::{Diff, GlobalId, Row};
22use mz_storage_types::errors::{DataflowError, EnvelopeError};
23use mz_timely_util::builder_async::{
24    Event as AsyncEvent, OperatorBuilder as AsyncOperatorBuilder, PressOnDropButton,
25};
26use std::convert::Infallible;
27use timely::container::CapacityContainerBuilder;
28use timely::dataflow::StreamVec;
29use timely::dataflow::channels::pact::Exchange;
30use timely::dataflow::operators::{Capability, CapabilitySet};
31use timely::order::{PartialOrder, TotalOrder};
32use timely::progress::timestamp::Refines;
33use timely::progress::{Antichain, Timestamp};
34
35use crate::healthcheck::HealthStatusUpdate;
36use crate::metrics::upsert::UpsertMetrics;
37use crate::upsert::UpsertConfig;
38use crate::upsert::UpsertErrorEmitter;
39use crate::upsert::UpsertKey;
40use crate::upsert::UpsertValue;
41use crate::upsert::types::UpsertValueAndSize;
42use crate::upsert::types::{self as upsert_types, ValueMetadata};
43use crate::upsert::types::{StateValue, UpsertState, UpsertStateBackend};
44
45/// An operator that transforms an input stream of upserts (updates to key-value
46/// pairs), which represents an imaginary key-value state, into a differential
47/// collection. It keeps an internal map-like state which keeps the latest value
48/// for each key, such that it can emit the retractions and additions implied by
49/// a new update for a given key.
50///
51/// This operator is intended to be used in an ingestion pipeline that reads
52/// from an external source, and the output of this operator is eventually
53/// written to persist.
54///
55/// The operator has two inputs: a) the source input, of upserts, and b) a
56/// persist input that feeds back the upsert state to the operator. Below, there
57/// is a section for each input that describes how and why we process updates
58/// from each input.
59///
60/// An important property of this operator is that it does _not_ update the
61/// map-like state that it keeps for translating the stream of upserts into a
62/// differential collection when it processes source input. It _only_ updates
63/// the map-like state based on updates from the persist (feedback) input. We do
64/// this because the operator is expected to be used in cases where there are
65/// multiple concurrent instances of the same ingestion pipeline, and the
66/// different instances might see different input because of concurrency and
67/// non-determinism. All instances of the upsert operator must produce output
68/// that is consistent with the current state of the output (that all instances
69/// produce "collaboratively"). This global state is what the operator
70/// continually learns about via updates from the persist input.
71///
72/// ## Processing the Source Input
73///
74/// Updates on the source input are stashed/staged until they can be processed.
75/// Whether or not an update can be processed depends both on the upper frontier
76/// of the source input and on the upper frontier of the persist input:
77///
78///  - Input updates are only processed once their timestamp is "done", that is
79///  the input upper is no longer `less_equal` their timestamp.
80///
81///  - Input updates are only processed once they are at the persist upper, that
82///  is we have emitted and written down updates for all previous times and we
83///  have updated our map-like state to the latest global state of the output of
84///  the ingestion pipeline. We know this is the case when the persist upper is
85///  no longer `less_than` their timestamp.
86///
87/// As an optimization, we allow processing input updates when they are right at
88/// the input frontier. This is called _partial emission_ because we are
89/// emitting updates that might be retracted when processing more updates from
90/// the same timestamp. In order to be able to process these updates we keep
91/// _provisional values_ in our upsert state. These will be overwritten when we
92/// get the final upsert values on the persist input.
93///
94/// ## Processing the Persist Input
95///
96/// We continually ingest updates from the persist input into our state using
97/// `UpsertState::consolidate_chunk`. We might be ingesting updates from the
98/// initial snapshot (when starting the operator) that are not consolidated or
99/// we might be ingesting updates from a partial emission (see above). In either
100/// case, our input might not be consolidated and `consolidate_chunk` is able to
101/// handle that.
102pub fn upsert_inner<'scope, T, FromTime, F, Fut, US>(
103    input: VecCollection<'scope, T, (UpsertKey, Option<UpsertValue>, FromTime), Diff>,
104    key_indices: Vec<usize>,
105    resume_upper: Antichain<T>,
106    persist_ok: VecCollection<'scope, T, Row, Diff>,
107    persist_err: VecCollection<'scope, T, DataflowError, Diff>,
108    mut persist_token: Option<Vec<PressOnDropButton>>,
109    upsert_metrics: UpsertMetrics,
110    source_config: crate::source::SourceExportCreationConfig,
111    state_fn: F,
112    upsert_config: UpsertConfig,
113    prevent_snapshot_buffering: bool,
114    snapshot_buffering_max: Option<usize>,
115) -> (
116    VecCollection<'scope, T, Result<Row, DataflowError>, Diff>,
117    StreamVec<'scope, T, (Option<GlobalId>, HealthStatusUpdate)>,
118    StreamVec<'scope, T, Infallible>,
119    PressOnDropButton,
120)
121where
122    T: Timestamp + Refines<mz_repr::Timestamp> + TotalOrder + Sync,
123    F: FnOnce() -> Fut + 'static,
124    Fut: std::future::Future<Output = US>,
125    US: UpsertStateBackend<T, FromTime>,
126    FromTime: Debug + timely::ExchangeData + Clone + Ord + Sync,
127{
128    let mut builder = AsyncOperatorBuilder::new("Upsert".to_string(), input.scope());
129
130    let persist_input = crate::upsert::key_persist_feedback(persist_ok, persist_err, key_indices);
131    let (output_handle, output) = builder.new_output::<CapacityContainerBuilder<_>>();
132
133    // An output that just reports progress of the snapshot consolidation process upstream to the
134    // persist source to ensure that backpressure is applied
135    let (_snapshot_handle, snapshot_stream) =
136        builder.new_output::<CapacityContainerBuilder<Vec<Infallible>>>();
137
138    let (mut health_output, health_stream) = builder.new_output();
139    let mut input = builder.new_input_for(
140        input.inner,
141        Exchange::new(move |((key, _, _), _, _)| UpsertKey::hashed(key)),
142        &output_handle,
143    );
144
145    let mut persist_input = builder.new_disconnected_input(
146        persist_input.inner,
147        Exchange::new(|((key, _), _, _)| UpsertKey::hashed(key)),
148    );
149
150    let upsert_shared_metrics = Arc::clone(&upsert_metrics.shared);
151
152    let shutdown_button = builder.build(move |caps| async move {
153        let [output_cap, snapshot_cap, health_cap]: [_; 3] = caps.try_into().unwrap();
154        drop(output_cap);
155        let mut snapshot_cap = CapabilitySet::from_elem(snapshot_cap);
156
157        let mut state = UpsertState::<_, T, FromTime>::new(
158            state_fn().await,
159            upsert_shared_metrics,
160            &upsert_metrics,
161            source_config.source_statistics.clone(),
162            upsert_config.shrink_upsert_unused_buffers_by_ratio,
163        );
164
165        // True while we're still reading the initial "snapshot" (a whole bunch
166        // of updates, all at the same initial timestamp) from our persist
167        // input or while we're reading the initial snapshot from the upstream
168        // source.
169        let mut hydrating = true;
170
171        // A re-usable buffer of changes, per key. This is an `IndexMap`
172        // because it has to be `drain`-able and have a consistent iteration
173        // order.
174        let mut commands_state: indexmap::IndexMap<
175            _,
176            upsert_types::UpsertValueAndSize<T, FromTime>,
177        > = indexmap::IndexMap::new();
178        let mut multi_get_scratch = Vec::new();
179
180        // For stashing source input while it's not eligible for processing.
181        let mut stash = vec![];
182        // A capability suitable for emitting any updates based on stash. No capability is held
183        // when the stash is empty.
184        let mut stash_cap: Option<Capability<T>> = None;
185        let mut input_upper = Antichain::from_elem(Timestamp::minimum());
186        let mut partial_drain_time = None;
187
188        // For our persist/feedback input, both of these.
189        let mut persist_stash = vec![];
190        let mut persist_upper = Antichain::from_elem(Timestamp::minimum());
191
192        // We keep track of the largest timestamp seen on the persist input so
193        // that we can block processing source input while that timestamp is
194        // beyond the persist frontier. While ingesting updates of a timestamp,
195        // our upsert state is in a consolidating state, and trying to read it
196        // at that time would yield a panic.
197        //
198        // NOTE(aljoscha): You would think that it cannot happen that we even
199        // attempt to process source updates while the state is in a
200        // consolidating state, because we always wait until the persist
201        // frontier "catches up" with the timestamp of the source input. If
202        // there is only this here UPSERT operator and no concurrent instances,
203        // this is true. But with concurrent instances it can happen that an
204        // operator that is faster than us makes it so updates get written to
205        // persist. And we would then be ingesting them.
206        let mut largest_seen_persist_ts: Option<T> = None;
207
208        // A buffer for our output.
209        let mut output_updates = vec![];
210
211        let mut error_emitter = (&mut health_output, &health_cap);
212
213        loop {
214            tokio::select! {
215                _ = persist_input.ready() => {
216                    // Read away as much input as we can.
217                    while let Some(persist_event) = persist_input.next_sync() {
218                        match persist_event {
219                            AsyncEvent::Data(time, data) => {
220                                tracing::trace!(
221                                    worker_id = %source_config.worker_id,
222                                    source_id = %source_config.id,
223                                    time=?time,
224                                    updates=%data.len(),
225                                    "received persist data");
226
227                                persist_stash.extend(data.into_iter().map(
228                                    |((key, value), ts, diff)| {
229                                        largest_seen_persist_ts =
230                                            std::cmp::max(
231                                                largest_seen_persist_ts
232                                                    .clone(),
233                                                Some(ts.clone()),
234                                            );
235                                        (key, value, ts, diff)
236                                    },
237                                ));
238                            }
239                            AsyncEvent::Progress(upper) => {
240                                tracing::trace!(
241                                    worker_id = %source_config.worker_id,
242                                    source_id = %source_config.id,
243                                    ?upper,
244                                    "received persist progress");
245                                persist_upper = upper;
246                            }
247                        }
248                    }
249
250                    let last_rehydration_chunk =
251                        hydrating && PartialOrder::less_equal(&resume_upper, &persist_upper);
252
253                    tracing::debug!(
254                        worker_id = %source_config.worker_id,
255                        source_id = %source_config.id,
256                        persist_stash = %persist_stash.len(),
257                        %hydrating,
258                        %last_rehydration_chunk,
259                        ?resume_upper,
260                        ?persist_upper,
261                        "ingesting persist snapshot chunk");
262
263                    // Log any (key, ts) pairs in this batch that have a suspicious
264                    // net diff, to help diagnose how diff_sum corruption enters the
265                    // system.
266                    //
267                    // Consolidating by key alone is too noisy during hydration,
268                    // because a single batch can legitimately contain multiple
269                    // timestamps for the same key. The suspicious shape for this
270                    // bug is multiple net updates for the same key at one logical
271                    // timestamp.
272                    {
273                        let mut key_ts_diffs: Vec<(
274                            (UpsertKey, T),
275                            mz_repr::Diff
276                        )> = persist_stash
277                            .iter()
278                            .map(|(key, _val, ts, diff)| ((*key, ts.clone()), *diff))
279                            .collect();
280                        differential_dataflow::consolidation::consolidate(&mut key_ts_diffs);
281                        for ((key, ts), net_diff) in &key_ts_diffs {
282                            if net_diff.into_inner() > 1 || net_diff.into_inner() < -1 {
283                                tracing::warn!(
284                                    worker_id = %source_config.worker_id,
285                                    source_id = %source_config.id,
286                                    ?key,
287                                    ?ts,
288                                    net_diff = net_diff.into_inner(),
289                                    %hydrating,
290                                    ?persist_upper,
291                                    "persist feedback batch has (key, ts) with suspicious net diff \
292                                    (expected -1, 0, or 1 after per-(key, ts) consolidation)"
293                                );
294                            }
295                        }
296                    }
297
298                    let persist_stash_iter = persist_stash
299                        .drain(..)
300                        .map(|(key, val, _ts, diff)| (key, val, diff));
301
302                    match state
303                        .consolidate_chunk(
304                            persist_stash_iter,
305                            last_rehydration_chunk,
306                        )
307                        .await
308                    {
309                        Ok(_) => {}
310                        Err(e) => {
311                            // Make sure our persist source can shut down.
312                            persist_token.take();
313                            snapshot_cap.downgrade(&[]);
314                            UpsertErrorEmitter::<T>::emit(
315                                &mut error_emitter,
316                                "Failed to rehydrate state".to_string(),
317                                e,
318                            )
319                            .await;
320                        }
321                    }
322
323                    tracing::debug!(
324                        worker_id = %source_config.worker_id,
325                        source_id = %source_config.id,
326                        ?resume_upper,
327                        ?persist_upper,
328                        "downgrading snapshot cap",
329                    );
330
331                    // Only downgrade this _after_ ingesting the data, because
332                    // that can actually take quite some time, and we don't want
333                    // to announce that we're done ingesting the initial
334                    // snapshot too early.
335                    //
336                    // When we finish ingesting our initial persist snapshot,
337                    // during "re-hydration", we downgrade this to the empty
338                    // frontier, so we need to be lenient to this failing from
339                    // then on.
340                    let _ = snapshot_cap.try_downgrade(persist_upper.iter());
341
342
343
344                    if last_rehydration_chunk {
345                        hydrating = false;
346
347                        tracing::info!(
348                            worker_id = %source_config.worker_id,
349                            source_id = %source_config.id,
350                            "upsert source finished rehydration",
351                        );
352
353                        snapshot_cap.downgrade(&[]);
354                    }
355
356                }
357                _ = input.ready() => {
358                    let mut events_processed = 0;
359                    while let Some(event) = input.next_sync() {
360                        match event {
361                            AsyncEvent::Data(cap, mut data) => {
362                                tracing::trace!(
363                                    worker_id = %source_config.worker_id,
364                                    source_id = %source_config.id,
365                                    time=?cap.time(),
366                                    updates=%data.len(),
367                                    "received data");
368
369                                let event_time = cap.time().clone();
370
371                                stage_input(
372                                    &mut stash,
373                                    &mut data,
374                                    &input_upper,
375                                    &resume_upper,
376                                );
377                                if !stash.is_empty() {
378                                    // Update the stashed capability to the minimum
379                                    stash_cap = match stash_cap {
380                                        Some(stash_cap) => {
381                                            if cap.time() < stash_cap.time() {
382                                                Some(cap)
383                                            } else {
384                                                Some(stash_cap)
385                                            }
386                                        }
387                                        None => Some(cap)
388                                    };
389                                }
390
391                                if prevent_snapshot_buffering
392                                    && input_upper.as_option()
393                                        == Some(&event_time)
394                                {
395                                    tracing::debug!(
396                                        worker_id = %source_config.worker_id,
397                                        source_id = %source_config.id,
398                                        ?event_time,
399                                        ?resume_upper,
400                                        ?input_upper,
401                                        "allowing partial drain");
402                                    partial_drain_time = Some(event_time.clone());
403                                } else {
404                                    tracing::debug!(
405                                        worker_id = %source_config.worker_id,
406                                        source_id = %source_config.id,
407                                        %prevent_snapshot_buffering,
408                                        ?event_time,
409                                        ?resume_upper,
410                                        ?input_upper,
411                                        "not allowing partial drain");
412                                }
413                            }
414                            AsyncEvent::Progress(upper) => {
415                                tracing::trace!(
416                                    worker_id = %source_config.worker_id,
417                                    source_id = %source_config.id,
418                                    ?upper,
419                                    "received progress");
420
421                                // Ignore progress updates before the `resume_upper`, which is our initial
422                                // capability post-snapshotting.
423                                if PartialOrder::less_than(&upper, &resume_upper) {
424                                    tracing::trace!(
425                                        worker_id = %source_config.worker_id,
426                                        source_id = %source_config.id,
427                                        ?upper,
428                                        ?resume_upper,
429                                        "ignoring progress updates before resume_upper");
430                                    continue;
431                                }
432
433                                // Disable partial drain, because this progress
434                                // update has moved the frontier. We might allow
435                                // it again once we receive data right at the
436                                // frontier again.
437                                partial_drain_time = None;
438                                input_upper = upper;
439                            }
440                        }
441
442                        events_processed += 1;
443                        if let Some(max) = snapshot_buffering_max {
444                            if events_processed >= max {
445                                break;
446                            }
447                        }
448                    }
449                }
450            };
451
452            // While we have partially ingested updates of a timestamp our state
453            // is in an inconsistent/consolidating state and accessing it would
454            // panic.
455            if let Some(largest_seen_persist_ts) = largest_seen_persist_ts.as_ref() {
456                let largest_seen_outer_persist_ts = largest_seen_persist_ts.clone().to_outer();
457                let outer_persist_upper = persist_upper.iter().map(|ts| ts.clone().to_outer());
458                let outer_persist_upper = Antichain::from_iter(outer_persist_upper);
459                if outer_persist_upper.less_equal(&largest_seen_outer_persist_ts) {
460                    continue;
461                }
462            }
463
464            // We try and drain from our stash every time we go through the
465            // loop. More of our stash can become eligible for draining both
466            // when the source-input frontier advances or when the persist
467            // frontier advances.
468            if !stash.is_empty() {
469                let cap = stash_cap
470                    .as_mut()
471                    .expect("missing capability for non-empty stash");
472
473                tracing::trace!(
474                    worker_id = %source_config.worker_id,
475                    source_id = %source_config.id,
476                    ?cap,
477                    ?stash,
478                    "stashed updates");
479
480                let mut min_remaining_time = drain_staged_input::<_, _, _, _>(
481                    &mut stash,
482                    &mut commands_state,
483                    &mut output_updates,
484                    &mut multi_get_scratch,
485                    DrainStyle::ToUpper {
486                        input_upper: &input_upper,
487                        persist_upper: &persist_upper,
488                    },
489                    &mut error_emitter,
490                    &mut state,
491                    &source_config,
492                )
493                .await;
494
495                tracing::trace!(
496                    worker_id = %source_config.worker_id,
497                    source_id = %source_config.id,
498                    output_updates = %output_updates.len(),
499                    "output updates for complete timestamp");
500
501                for (update, ts, diff) in output_updates.drain(..) {
502                    output_handle.give(cap, (update, ts, diff));
503                }
504
505                if !stash.is_empty() {
506                    let min_remaining_time = min_remaining_time
507                        .take()
508                        .expect("we still have updates left");
509                    cap.downgrade(&min_remaining_time);
510                } else {
511                    stash_cap = None;
512                }
513            }
514
515            if input_upper.is_empty() {
516                tracing::debug!(
517                    worker_id = %source_config.worker_id,
518                    source_id = %source_config.id,
519                    "input exhausted, shutting down");
520                break;
521            };
522
523            // If there were staged events that occurred at the capability time, drain
524            // them. This is safe because out-of-order updates to the same key that are
525            // drained in separate calls to `drain_staged_input` are correctly ordered by
526            // their `FromTime` in `drain_staged_input`.
527            //
528            // Note also that this may result in more updates in the output collection than
529            // the minimum. However, because the frontier only advances on `Progress` updates,
530            // the collection always accumulates correctly for all keys.
531            if let Some(partial_drain_time) = &partial_drain_time {
532                if !stash.is_empty() {
533                    let cap = stash_cap
534                        .as_mut()
535                        .expect("missing capability for non-empty stash");
536
537                    tracing::trace!(
538                        worker_id = %source_config.worker_id,
539                        source_id = %source_config.id,
540                        ?cap,
541                        ?stash,
542                        "stashed updates");
543
544                    let mut min_remaining_time = drain_staged_input::<_, _, _, _>(
545                        &mut stash,
546                        &mut commands_state,
547                        &mut output_updates,
548                        &mut multi_get_scratch,
549                        DrainStyle::AtTime {
550                            time: partial_drain_time.clone(),
551                            persist_upper: &persist_upper,
552                        },
553                        &mut error_emitter,
554                        &mut state,
555                        &source_config,
556                    )
557                    .await;
558
559                    tracing::trace!(
560                        worker_id = %source_config.worker_id,
561                        source_id = %source_config.id,
562                        output_updates = %output_updates.len(),
563                        "output updates for partial timestamp");
564
565                    for (update, ts, diff) in output_updates.drain(..) {
566                        output_handle.give(cap, (update, ts, diff));
567                    }
568
569                    if !stash.is_empty() {
570                        let min_remaining_time = min_remaining_time
571                            .take()
572                            .expect("we still have updates left");
573                        cap.downgrade(&min_remaining_time);
574                    } else {
575                        stash_cap = None;
576                    }
577                }
578            }
579        }
580    });
581
582    (
583        output
584            .as_collection()
585            .map(|result: UpsertValue| match result {
586                Ok(ok) => Ok(ok),
587                Err(err) => Err(DataflowError::from(EnvelopeError::Upsert(*err))),
588            }),
589        health_stream,
590        snapshot_stream,
591        shutdown_button.press_on_drop(),
592    )
593}
594
595/// Helper method for [`upsert_inner`] used to stage `data` updates
596/// from the input/source timely edge.
597#[allow(clippy::disallowed_types)]
598fn stage_input<T, FromTime>(
599    stash: &mut Vec<(T, UpsertKey, Reverse<FromTime>, Option<UpsertValue>)>,
600    data: &mut Vec<((UpsertKey, Option<UpsertValue>, FromTime), T, Diff)>,
601    input_upper: &Antichain<T>,
602    resume_upper: &Antichain<T>,
603) where
604    T: PartialOrder + timely::progress::Timestamp,
605    FromTime: Ord,
606{
607    if PartialOrder::less_equal(input_upper, resume_upper) {
608        data.retain(|(_, ts, _)| resume_upper.less_equal(ts));
609    }
610
611    stash.extend(data.drain(..).map(|((key, value, order), time, diff)| {
612        assert!(diff.is_positive(), "invalid upsert input");
613        (time, key, Reverse(order), value)
614    }));
615}
616
617/// The style of drain we are performing on the stash. `AtTime`-drains cannot
618/// assume that all values have been seen, and must leave tombstones behind for deleted values.
619#[derive(Debug)]
620enum DrainStyle<'a, T> {
621    ToUpper {
622        input_upper: &'a Antichain<T>,
623        persist_upper: &'a Antichain<T>,
624    },
625    // For partial draining when taking the source snapshot.
626    AtTime {
627        time: T,
628        persist_upper: &'a Antichain<T>,
629    },
630}
631
632/// Helper method for [`upsert_inner`] used to stage `data` updates
633/// from the input timely edge.
634///
635/// Returns the minimum observed time across the updates that remain in the
636/// stash or `None` if none are left.
637///
638/// ## Correctness
639///
640/// It is safe to call this function multiple times with the same `persist_upper` provided that the
641/// drain style is `AtTime`, which updates the state such that past actions are remembered and can
642/// be undone in subsequent calls.
643///
644/// It is *not* safe to call this function more than once with the same `persist_upper` and a
645/// `ToUpper` drain style. Doing so causes all calls except the first one to base their work on
646/// stale state, since in this drain style no modifications to the state are made.
647async fn drain_staged_input<S, T, FromTime, E>(
648    stash: &mut Vec<(T, UpsertKey, Reverse<FromTime>, Option<UpsertValue>)>,
649    commands_state: &mut indexmap::IndexMap<UpsertKey, UpsertValueAndSize<T, FromTime>>,
650    output_updates: &mut Vec<(UpsertValue, T, Diff)>,
651    multi_get_scratch: &mut Vec<UpsertKey>,
652    drain_style: DrainStyle<'_, T>,
653    error_emitter: &mut E,
654    state: &mut UpsertState<'_, S, T, FromTime>,
655    source_config: &crate::source::SourceExportCreationConfig,
656) -> Option<T>
657where
658    S: UpsertStateBackend<T, FromTime>,
659    T: Timestamp + TotalOrder + timely::ExchangeData + Clone + Debug + Ord + Sync,
660    FromTime: timely::ExchangeData + Clone + Ord + Sync,
661    E: UpsertErrorEmitter<T>,
662{
663    let mut min_remaining_time = Antichain::new();
664
665    let mut eligible_updates = stash
666        .extract_if(.., |(ts, _, _, _)| {
667            let eligible = match &drain_style {
668                DrainStyle::ToUpper {
669                    input_upper,
670                    persist_upper,
671                } => {
672                    // We make sure that a) we only process updates when we know their
673                    // timestamp is complete, that is there will be no more updates for
674                    // that timestamp, and b) that "previous" times in the persist
675                    // input are complete. The latter makes sure that we emit updates
676                    // for the next timestamp that are consistent with the global state
677                    // in the output persist shard, which also serves as a persistent
678                    // copy of our in-memory/on-disk upsert state.
679                    !input_upper.less_equal(ts) && !persist_upper.less_than(ts)
680                }
681                DrainStyle::AtTime {
682                    time,
683                    persist_upper,
684                } => {
685                    // Even when emitting partial updates, we still need to wait
686                    // until "previous" times in the persist input are complete.
687                    *ts <= *time && !persist_upper.less_than(ts)
688                }
689            };
690
691            if !eligible {
692                min_remaining_time.insert(ts.clone());
693            }
694
695            eligible
696        })
697        .filter(|(ts, _, _, _)| {
698            let persist_upper = match &drain_style {
699                DrainStyle::ToUpper {
700                    input_upper: _,
701                    persist_upper,
702                } => persist_upper,
703                DrainStyle::AtTime {
704                    time: _,
705                    persist_upper,
706                } => persist_upper,
707            };
708
709            // Any update that is "in the past" of the persist upper is not
710            // relevant anymore. We _can_ emit changes for it, but the
711            // downstream persist_sink would filter these updates out because
712            // the shard upper is already further ahead.
713            //
714            // Plus, our upsert state is up-to-date to the persist_upper, so we
715            // wouldn't be able to emit correct retractions for incoming
716            // commands whose `ts` is in the past of that.
717            let relevant = persist_upper.less_equal(ts);
718            relevant
719        })
720        .collect_vec();
721
722    tracing::debug!(
723        worker_id = %source_config.worker_id,
724        source_id = %source_config.id,
725        ?drain_style,
726        remaining = %stash.len(),
727        eligible = eligible_updates.len(),
728        "draining stash");
729
730    // Sort the eligible updates by (key, time, Reverse(from_time)) so that
731    // deduping by (key, time) gives the latest change for that key.
732    eligible_updates.sort_unstable_by(|a, b| {
733        let (ts1, key1, from_ts1, val1) = a;
734        let (ts2, key2, from_ts2, val2) = b;
735        Ord::cmp(&(ts1, key1, from_ts1, val1), &(ts2, key2, from_ts2, val2))
736    });
737
738    // Read the previous values _per key_ out of `state`, recording it
739    // along with the value with the _latest timestamp for that key_.
740    commands_state.clear();
741    for (_, key, _, _) in eligible_updates.iter() {
742        commands_state.entry(*key).or_default();
743    }
744
745    // These iterators iterate in the same order because `commands_state`
746    // is an `IndexMap`.
747    multi_get_scratch.clear();
748    multi_get_scratch.extend(commands_state.iter().map(|(k, _)| *k));
749    match state
750        .multi_get(multi_get_scratch.drain(..), commands_state.values_mut())
751        .await
752    {
753        Ok(_) => {}
754        Err(e) => {
755            error_emitter
756                .emit("Failed to fetch records from state".to_string(), e)
757                .await;
758        }
759    }
760
761    // From the prefix that can be emitted we can deduplicate based on (ts, key) in
762    // order to only process the command with the maximum order within the (ts,
763    // key) group. This is achieved by wrapping order in `Reverse(FromTime)` above.;
764    let mut commands = eligible_updates.into_iter().dedup_by(|a, b| {
765        let ((a_ts, a_key, _, _), (b_ts, b_key, _, _)) = (a, b);
766        a_ts == b_ts && a_key == b_key
767    });
768
769    let bincode_opts = upsert_types::upsert_bincode_opts();
770    // Upsert the values into `commands_state`, by recording the latest
771    // value (or deletion). These will be synced at the end to the `state`.
772    //
773    // Note that we are effectively doing "mini-upsert" here, using
774    // `command_state`. This "mini-upsert" is seeded with data from `state`, using
775    // a single `multi_get` above, and the final state is written out into
776    // `state` using a single `multi_put`. This simplifies `UpsertStateBackend`
777    // implementations, and reduces the number of reads and write we need to do.
778    //
779    // This "mini-upsert" technique is actually useful in `UpsertState`'s
780    // `consolidate_snapshot_read_write_inner` implementation, minimizing gets and puts on
781    // the `UpsertStateBackend` implementations. In some sense, its "upsert all the way down".
782    while let Some((ts, key, from_time, value)) = commands.next() {
783        let mut command_state = if let Entry::Occupied(command_state) = commands_state.entry(key) {
784            command_state
785        } else {
786            panic!("key missing from commands_state");
787        };
788
789        let existing_state_cell = &mut command_state.get_mut().value;
790
791        if let Some(cs) = existing_state_cell.as_mut() {
792            cs.ensure_decoded(bincode_opts, source_config.id, Some(&key));
793        }
794
795        // Skip this command if its order key is below the one in the upsert state.
796        // Note that the existing order key may be `None` if the existing value
797        // is from snapshotting, which always sorts below new values/deletes.
798        let existing_order = existing_state_cell
799            .as_ref()
800            .and_then(|cs| cs.provisional_order(&ts));
801        if existing_order >= Some(&from_time.0) {
802            // Skip this update. If no later updates adjust this key, then we just
803            // end up writing the same value back to state. If there
804            // is nothing in the state, `existing_order` is `None`, and this
805            // does not occur.
806            continue;
807        }
808
809        match value {
810            Some(value) => {
811                if let Some(old_value) = existing_state_cell.as_ref() {
812                    if let Some(old_value) = old_value.provisional_value_ref(&ts) {
813                        output_updates.push((old_value.clone(), ts.clone(), Diff::MINUS_ONE));
814                    }
815                }
816
817                match &drain_style {
818                    DrainStyle::AtTime { .. } => {
819                        let existing_value = existing_state_cell.take();
820
821                        let new_value = match existing_value {
822                            Some(existing_value) => existing_value.clone().into_provisional_value(
823                                value.clone(),
824                                ts.clone(),
825                                from_time.0.clone(),
826                            ),
827                            None => StateValue::new_provisional_value(
828                                value.clone(),
829                                ts.clone(),
830                                from_time.0.clone(),
831                            ),
832                        };
833
834                        existing_state_cell.replace(new_value);
835                    }
836                    DrainStyle::ToUpper { .. } => {
837                        // Not writing down provisional values, or anything.
838                    }
839                };
840
841                output_updates.push((value, ts, Diff::ONE));
842            }
843            None => {
844                if let Some(old_value) = existing_state_cell.as_ref() {
845                    if let Some(old_value) = old_value.provisional_value_ref(&ts) {
846                        output_updates.push((old_value.clone(), ts.clone(), Diff::MINUS_ONE));
847                    }
848                }
849
850                match &drain_style {
851                    DrainStyle::AtTime { .. } => {
852                        let existing_value = existing_state_cell.take();
853
854                        let new_value = match existing_value {
855                            Some(existing_value) => existing_value
856                                .into_provisional_tombstone(ts.clone(), from_time.0.clone()),
857                            None => StateValue::new_provisional_tombstone(
858                                ts.clone(),
859                                from_time.0.clone(),
860                            ),
861                        };
862
863                        existing_state_cell.replace(new_value);
864                    }
865                    DrainStyle::ToUpper { .. } => {
866                        // Not writing down provisional values, or anything.
867                    }
868                }
869            }
870        }
871    }
872
873    match &drain_style {
874        DrainStyle::AtTime { .. } => {
875            match state
876                .multi_put(
877                    // We don't want to update per-record stats, like size of
878                    // records indexed or count of records indexed.
879                    //
880                    // We only add provisional values and these will be
881                    // overwritten once we receive updates for state from the
882                    // persist input. And the merge functionality cannot know
883                    // what was in state before merging, so it cannot correctly
884                    // retract/update stats added here.
885                    //
886                    // Mostly, the merge functionality can't update those stats
887                    // because merging happens in a function that we pass to
888                    // rocksdb which doesn't have access to any external
889                    // context. And in general, with rocksdb we do blind writes
890                    // rather than inspect what was there before when
891                    // updating/inserting.
892                    false,
893                    commands_state.drain(..).map(|(k, cv)| {
894                        (
895                            k,
896                            upsert_types::PutValue {
897                                value: cv.value.map(|cv| cv.into_decoded()),
898                                previous_value_metadata: cv.metadata.map(|v| ValueMetadata {
899                                    size: v.size.try_into().expect("less than i64 size"),
900                                    is_tombstone: v.is_tombstone,
901                                }),
902                            },
903                        )
904                    }),
905                )
906                .await
907            {
908                Ok(_) => {}
909                Err(e) => {
910                    error_emitter
911                        .emit("Failed to update records in state".to_string(), e)
912                        .await;
913                }
914            }
915        }
916        style @ DrainStyle::ToUpper { .. } => {
917            tracing::trace!(
918                worker_id = %source_config.worker_id,
919                source_id = %source_config.id,
920                "not doing state update for drain style {:?}", style);
921        }
922    }
923
924    min_remaining_time.into_option()
925}
926
927#[cfg(test)]
928mod test {
929    //! No test drives errors through the persist feedback, so every harness below closes
930    //! the error input by dropping its handle.
931
932    use std::sync::mpsc;
933
934    use mz_ore::metrics::MetricsRegistry;
935    use mz_persist_types::ShardId;
936    use mz_repr::{Datum, Timestamp as MzTimestamp};
937    use mz_rocksdb::{RocksDBConfig, ValueIterator};
938    use mz_storage_operators::persist_source::Subtime;
939    use mz_storage_types::sources::SourceEnvelope;
940    use mz_storage_types::sources::envelope::{KeyEnvelope, UpsertEnvelope, UpsertStyle};
941    use rocksdb::Env;
942    use timely::dataflow::operators::capture::Extract;
943    use timely::dataflow::operators::{Capture, Input, Probe};
944    use timely::progress::Timestamp;
945
946    use crate::metrics::StorageMetrics;
947    use crate::metrics::upsert::UpsertMetricDefs;
948    use crate::source::SourceExportCreationConfig;
949    use crate::statistics::{SourceStatistics, SourceStatisticsMetricDefs};
950    use crate::upsert::memory::InMemoryHashMap;
951    use crate::upsert::types::{BincodeOpts, consolidating_merge_function, upsert_bincode_opts};
952
953    use super::*;
954
955    #[mz_ore::test]
956    #[cfg_attr(miri, ignore)]
957    fn gh_9160_repro() {
958        // Helper to wrap timestamps in the appropriate types
959        let new_ts = |ts| (MzTimestamp::new(ts), Subtime::minimum());
960
961        let output_handle = timely::execute_directly(move |worker| {
962            let (mut input_handle, mut persist_handle, output_handle) = worker
963                .dataflow::<MzTimestamp, _, _>(|scope| {
964                    // Enter a subscope since the upsert operator expects to work a backpressure
965                    // enabled scope.
966                    scope.scoped::<(MzTimestamp, Subtime), _, _>("upsert", |scope| {
967                        let (input_handle, input) = scope.new_input();
968                        let (persist_handle, persist_input) = scope.new_input();
969                        let (_persist_err_handle, persist_err_input) = scope.new_input();
970                        let upsert_config = UpsertConfig {
971                            shrink_upsert_unused_buffers_by_ratio: 0,
972                        };
973                        let source_id = GlobalId::User(0);
974                        let metrics_registry = MetricsRegistry::new();
975                        let upsert_metrics_defs =
976                            UpsertMetricDefs::register_with(&metrics_registry);
977                        let upsert_metrics =
978                            UpsertMetrics::new(&upsert_metrics_defs, source_id, 0, None);
979
980                        let metrics_registry = MetricsRegistry::new();
981                        let storage_metrics = StorageMetrics::register_with(&metrics_registry);
982
983                        let metrics_registry = MetricsRegistry::new();
984                        let source_statistics_defs =
985                            SourceStatisticsMetricDefs::register_with(&metrics_registry);
986                        let envelope = SourceEnvelope::Upsert(UpsertEnvelope {
987                            source_arity: 2,
988                            style: UpsertStyle::Default(KeyEnvelope::Flattened),
989                            key_indices: vec![0],
990                        });
991                        let source_statistics = SourceStatistics::new(
992                            source_id,
993                            0,
994                            &source_statistics_defs,
995                            source_id,
996                            &ShardId::new(),
997                            envelope,
998                            Antichain::from_elem(Timestamp::minimum()),
999                        );
1000
1001                        let source_config = SourceExportCreationConfig {
1002                            id: GlobalId::User(0),
1003                            worker_id: 0,
1004                            metrics: storage_metrics,
1005                            source_statistics,
1006                        };
1007
1008                        let (output, _, _, button) = upsert_inner(
1009                            input.as_collection(),
1010                            vec![0],
1011                            Antichain::from_elem(Timestamp::minimum()),
1012                            persist_input.as_collection(),
1013                            persist_err_input.as_collection(),
1014                            None,
1015                            upsert_metrics,
1016                            source_config,
1017                            || async { InMemoryHashMap::default() },
1018                            upsert_config,
1019                            true,
1020                            None,
1021                        );
1022                        std::mem::forget(button);
1023
1024                        (input_handle, persist_handle, output.inner.capture())
1025                    })
1026                });
1027
1028            // We work with a hypothetical schema of (key int, value int).
1029
1030            // The input will contain records for two keys, 0 and 1.
1031            let key0 = UpsertKey::from_key(Ok(&Row::pack_slice(&[Datum::Int64(0)])));
1032            let key1 = UpsertKey::from_key(Ok(&Row::pack_slice(&[Datum::Int64(1)])));
1033
1034            // We will assume that the kafka topic contains the following messages with their
1035            // associated reclocked timestamp:
1036            //  1. {offset=1, key=0, value=0}    @ mz_time = 0
1037            //  2. {offset=2, key=1, value=NULL} @ mz_time = 2  // <- deletion of unrelated key. Causes the operator
1038            //                                                  //    to maintain the associated cap to time 2
1039            //  3. {offset=3, key=0, value=1}    @ mz_time = 3
1040            //  4. {offset=4, key=0, value=2}    @ mz_time = 3  // <- messages 2 and 3 are reclocked to time 3
1041            let value1 = Row::pack_slice(&[Datum::Int64(0), Datum::Int64(0)]);
1042            let value3 = Row::pack_slice(&[Datum::Int64(0), Datum::Int64(1)]);
1043            let value4 = Row::pack_slice(&[Datum::Int64(0), Datum::Int64(2)]);
1044            let msg1 = (key0, Some(Ok(value1.clone())), 1);
1045            let msg2 = (key1, None, 2);
1046            let msg3 = (key0, Some(Ok(value3)), 3);
1047            let msg4 = (key0, Some(Ok(value4)), 4);
1048
1049            // The first message will initialize the upsert state such that key 0 has value 0 and
1050            // produce an output update to that effect.
1051            input_handle.send((msg1, new_ts(0), Diff::ONE));
1052            input_handle.advance_to(new_ts(2));
1053            worker.step();
1054
1055            // We assume this worker succesfully CAAs the update to the shard so we send it back
1056            // through the persist_input
1057            persist_handle.send((value1, new_ts(0), Diff::ONE));
1058            persist_handle.advance_to(new_ts(1));
1059            worker.step();
1060
1061            // Then, messages 2 and 3 are sent as one batch with capability = 2
1062            input_handle.send_batch(&mut vec![
1063                (msg2, new_ts(2), Diff::ONE),
1064                (msg3, new_ts(3), Diff::ONE),
1065            ]);
1066            // Advance our capability to 3
1067            input_handle.advance_to(new_ts(3));
1068            // Message 4 is sent with capability 3
1069            input_handle.send_batch(&mut vec![(msg4, new_ts(3), Diff::ONE)]);
1070            // Advance our capability to 4
1071            input_handle.advance_to(new_ts(4));
1072            // We now step the worker so that the pending data is received. This causes the
1073            // operator to store internally the following map from capabilities to updates:
1074            // cap=2 => [ msg2, msg3 ]
1075            // cap=3 => [ msg4 ]
1076            worker.step();
1077
1078            // We now assume that another replica raced us and processed msg1 at time 2, which in
1079            // this test is a no-op so the persist frontier advances to time 3 without new data.
1080            persist_handle.advance_to(new_ts(3));
1081            // We now step this worker again, which will notice that the persist upper is {3} and
1082            // wlil attempt to process msg3 and msg4 *separately*, causing it to produce a double
1083            // retraction.
1084            worker.step();
1085
1086            output_handle
1087        });
1088
1089        let mut actual_output = output_handle
1090            .extract()
1091            .into_iter()
1092            .flat_map(|(_cap, container)| container)
1093            .collect();
1094        differential_dataflow::consolidation::consolidate_updates(&mut actual_output);
1095
1096        // The expected consolidated output contains only updates for key 0 which has the value 0
1097        // at timestamp 0 and the value 2 at timestamp 3
1098        let value1 = Row::pack_slice(&[Datum::Int64(0), Datum::Int64(0)]);
1099        let value4 = Row::pack_slice(&[Datum::Int64(0), Datum::Int64(2)]);
1100        let expected_output: Vec<(Result<Row, DataflowError>, _, _)> = vec![
1101            (Ok(value1.clone()), new_ts(0), Diff::ONE),
1102            (Ok(value1), new_ts(3), Diff::MINUS_ONE),
1103            (Ok(value4), new_ts(3), Diff::ONE),
1104        ];
1105        assert_eq!(actual_output, expected_output);
1106    }
1107
1108    #[mz_ore::test]
1109    #[cfg_attr(miri, ignore)]
1110    fn gh_9540_repro() {
1111        // Helper to wrap timestamps in the appropriate types
1112        let mz_ts = |ts| (MzTimestamp::new(ts), Subtime::minimum());
1113        let (tx, rx) = mpsc::channel::<std::thread::JoinHandle<()>>();
1114
1115        let rocksdb_dir = tempfile::tempdir().unwrap();
1116        let output_handle = timely::execute_directly(move |worker| {
1117            let tx = tx.clone();
1118            let (mut input_handle, mut persist_handle, output_probe, output_handle) =
1119                worker.dataflow::<MzTimestamp, _, _>(|scope| {
1120                    // Enter a subscope since the upsert operator expects to work a backpressure
1121                    // enabled scope.
1122                    scope.scoped::<(MzTimestamp, Subtime), _, _>("upsert", |scope| {
1123                        let (input_handle, input) = scope.new_input();
1124                        let (persist_handle, persist_input) = scope.new_input();
1125                        let (_persist_err_handle, persist_err_input) = scope.new_input();
1126                        let upsert_config = UpsertConfig {
1127                            shrink_upsert_unused_buffers_by_ratio: 0,
1128                        };
1129                        let source_id = GlobalId::User(0);
1130                        let metrics_registry = MetricsRegistry::new();
1131                        let upsert_metrics_defs =
1132                            UpsertMetricDefs::register_with(&metrics_registry);
1133                        let upsert_metrics =
1134                            UpsertMetrics::new(&upsert_metrics_defs, source_id, 0, None);
1135                        let rocksdb_shared_metrics = Arc::clone(&upsert_metrics.rocksdb_shared);
1136                        let rocksdb_instance_metrics =
1137                            Arc::clone(&upsert_metrics.rocksdb_instance_metrics);
1138
1139                        let metrics_registry = MetricsRegistry::new();
1140                        let storage_metrics = StorageMetrics::register_with(&metrics_registry);
1141
1142                        let metrics_registry = MetricsRegistry::new();
1143                        let source_statistics_defs =
1144                            SourceStatisticsMetricDefs::register_with(&metrics_registry);
1145                        let envelope = SourceEnvelope::Upsert(UpsertEnvelope {
1146                            source_arity: 2,
1147                            style: UpsertStyle::Default(KeyEnvelope::Flattened),
1148                            key_indices: vec![0],
1149                        });
1150                        let source_statistics = SourceStatistics::new(
1151                            source_id,
1152                            0,
1153                            &source_statistics_defs,
1154                            source_id,
1155                            &ShardId::new(),
1156                            envelope,
1157                            Antichain::from_elem(Timestamp::minimum()),
1158                        );
1159
1160                        let source_config = SourceExportCreationConfig {
1161                            id: GlobalId::User(0),
1162                            worker_id: 0,
1163                            metrics: storage_metrics,
1164                            source_statistics,
1165                        };
1166
1167                        // A closure that will initialize and return a configured RocksDB instance
1168                        let rocksdb_init_fn = move || async move {
1169                            let merge_operator = Some((
1170                                "upsert_state_snapshot_merge_v1".to_string(),
1171                                |a: &[u8],
1172                                 b: ValueIterator<
1173                                    BincodeOpts,
1174                                    StateValue<(MzTimestamp, Subtime), u64>,
1175                                >| {
1176                                    consolidating_merge_function::<(MzTimestamp, Subtime), u64>(
1177                                        a.into(),
1178                                        b,
1179                                    )
1180                                },
1181                            ));
1182                            let rocksdb_cleanup_tries = 5;
1183                            let tuning = RocksDBConfig::new(Default::default(), None);
1184                            let mut rocksdb_inst = mz_rocksdb::RocksDBInstance::new(
1185                                rocksdb_dir.path(),
1186                                mz_rocksdb::InstanceOptions::new(
1187                                    Env::mem_env().unwrap(),
1188                                    rocksdb_cleanup_tries,
1189                                    merge_operator,
1190                                    // For now, just use the same config as the one used for
1191                                    // merging snapshots.
1192                                    upsert_bincode_opts(),
1193                                ),
1194                                tuning,
1195                                rocksdb_shared_metrics,
1196                                rocksdb_instance_metrics,
1197                            )
1198                            .unwrap();
1199
1200                            let handle = rocksdb_inst.take_core_loop_handle().expect("join handle");
1201                            tx.send(handle).expect("sent joinhandle");
1202                            crate::upsert::rocksdb::RocksDB::new(rocksdb_inst)
1203                        };
1204
1205                        let (output, _, _, button) = upsert_inner(
1206                            input.as_collection(),
1207                            vec![0],
1208                            Antichain::from_elem(Timestamp::minimum()),
1209                            persist_input.as_collection(),
1210                            persist_err_input.as_collection(),
1211                            None,
1212                            upsert_metrics,
1213                            source_config,
1214                            rocksdb_init_fn,
1215                            upsert_config,
1216                            true,
1217                            None,
1218                        );
1219                        std::mem::forget(button);
1220
1221                        let (probe, stream) = output.inner.probe();
1222                        (input_handle, persist_handle, probe, stream.capture())
1223                    })
1224                });
1225
1226            // We work with a hypothetical schema of (key int, value int).
1227
1228            // The input will contain records for two keys, 0 and 1.
1229            let key0 = UpsertKey::from_key(Ok(&Row::pack_slice(&[Datum::Int64(0)])));
1230
1231            // We will assume that the kafka topic contains the following messages with their
1232            // associated reclocked timestamp:
1233            //  1. {offset=1, key=0, value=0}    @ mz_time = 0
1234            //  2. {offset=2, key=0, value=NULL} @ mz_time = 1
1235            //  3. {offset=3, key=0, value=0}    @ mz_time = 2
1236            //  4. {offset=4, key=0, value=NULL} @ mz_time = 2  // <- messages 3 and 4 are *BOTH* reclocked to time 2
1237            let value1 = Row::pack_slice(&[Datum::Int64(0), Datum::Int64(0)]);
1238            let msg1 = ((key0, Some(Ok(value1.clone())), 1), mz_ts(0), Diff::ONE);
1239            let msg2 = ((key0, None, 2), mz_ts(1), Diff::ONE);
1240            let msg3 = ((key0, Some(Ok(value1.clone())), 3), mz_ts(2), Diff::ONE);
1241            let msg4 = ((key0, None, 4), mz_ts(2), Diff::ONE);
1242
1243            // The first message will initialize the upsert state such that key 0 has value 0 and
1244            // produce an output update to that effect.
1245            input_handle.send(msg1);
1246            input_handle.advance_to(mz_ts(1));
1247            while output_probe.less_than(&mz_ts(1)) {
1248                worker.step_or_park(None);
1249            }
1250            // Feedback the produced output..
1251            persist_handle.send((value1.clone(), mz_ts(0), Diff::ONE));
1252            persist_handle.advance_to(mz_ts(1));
1253            // ..and send the next upsert command that deletes the key.
1254            input_handle.send(msg2);
1255            input_handle.advance_to(mz_ts(2));
1256            while output_probe.less_than(&mz_ts(2)) {
1257                worker.step_or_park(None);
1258            }
1259
1260            // Feedback the produced output..
1261            persist_handle.send((value1, mz_ts(1), Diff::MINUS_ONE));
1262            persist_handle.advance_to(mz_ts(2));
1263            // ..and send the next *out of order* upsert command that deletes the key. Here msg4
1264            // happens at offset 4 and the operator should rememeber that.
1265            input_handle.send(msg4);
1266            input_handle.flush();
1267            // Run the worker for enough steps to process these events. We can't guide the
1268            // execution with the probe here since the frontier does not advance, only provisional
1269            // updates are produced.
1270            for _ in 0..5 {
1271                worker.step();
1272            }
1273
1274            // Send the missing message that will now confuse the operator because it has lost
1275            // track that for key 0 it has already seen a command for offset 4, and therefore msg3
1276            // should be skipped.
1277            input_handle.send(msg3);
1278            input_handle.flush();
1279            input_handle.advance_to(mz_ts(3));
1280
1281            output_handle
1282        });
1283
1284        let mut actual_output = output_handle
1285            .extract()
1286            .into_iter()
1287            .flat_map(|(_cap, container)| container)
1288            .collect();
1289        differential_dataflow::consolidation::consolidate_updates(&mut actual_output);
1290
1291        // The expected consolidated output contains only updates for key 0 which has the value 0
1292        // at timestamp 0 and the value 2 at timestamp 3
1293        let value1 = Row::pack_slice(&[Datum::Int64(0), Datum::Int64(0)]);
1294        let expected_output: Vec<(Result<Row, DataflowError>, _, _)> = vec![
1295            (Ok(value1.clone()), mz_ts(0), Diff::ONE),
1296            (Ok(value1), mz_ts(1), Diff::MINUS_ONE),
1297        ];
1298        assert_eq!(actual_output, expected_output);
1299
1300        while let Ok(handle) = rx.recv() {
1301            handle.join().expect("threads completed successfully");
1302        }
1303    }
1304
1305    /// v1 counterpart of v2's `lagging_replacement_below_upper_strands_data`,
1306    /// under the IDENTICAL setup: the external writer has advanced the feedback
1307    /// `persist_upper` to T = 10 with no operator output, and the lagging
1308    /// replacement then produces source data at ts BELOW that upper (5 and 7).
1309    ///
1310    /// v1's drain classifies any `ts` in the past of `persist_upper` as not
1311    /// relevant (`relevant = persist_upper.less_equal(ts)`) and DROPS it — the
1312    /// downstream persist_sink would filter it anyway since the shard upper is
1313    /// already further ahead. So v1 strands nothing and its output frontier
1314    /// advances cleanly past persist_upper (here to the input upper, 11),
1315    /// unlike v2, which without its drop-fix pins below.
1316    #[mz_ore::test]
1317    #[cfg_attr(miri, ignore)]
1318    fn lagging_replacement_below_upper_is_dropped() {
1319        let mz_ts = |ts| (MzTimestamp::new(ts), Subtime::minimum());
1320        let (tx, rx) = mpsc::channel::<std::thread::JoinHandle<()>>();
1321
1322        let rocksdb_dir = tempfile::tempdir().unwrap();
1323        let (frontier, capture) = timely::execute_directly(move |worker| {
1324            let tx = tx.clone();
1325            let (mut input, mut persist, probe, capture) =
1326                worker.dataflow::<MzTimestamp, _, _>(|scope| {
1327                    scope.scoped::<(MzTimestamp, Subtime), _, _>("upsert", |scope| {
1328                        let (input_handle, input) = scope.new_input();
1329                        let (persist_handle, persist_input) = scope.new_input();
1330                        let (_persist_err_handle, persist_err_input) = scope.new_input();
1331                        let upsert_config = UpsertConfig {
1332                            shrink_upsert_unused_buffers_by_ratio: 0,
1333                        };
1334                        let source_id = GlobalId::User(0);
1335                        let metrics_registry = MetricsRegistry::new();
1336                        let upsert_metrics_defs =
1337                            UpsertMetricDefs::register_with(&metrics_registry);
1338                        let upsert_metrics =
1339                            UpsertMetrics::new(&upsert_metrics_defs, source_id, 0, None);
1340                        let rocksdb_shared_metrics = Arc::clone(&upsert_metrics.rocksdb_shared);
1341                        let rocksdb_instance_metrics =
1342                            Arc::clone(&upsert_metrics.rocksdb_instance_metrics);
1343
1344                        let metrics_registry = MetricsRegistry::new();
1345                        let storage_metrics = StorageMetrics::register_with(&metrics_registry);
1346
1347                        let metrics_registry = MetricsRegistry::new();
1348                        let source_statistics_defs =
1349                            SourceStatisticsMetricDefs::register_with(&metrics_registry);
1350                        let envelope = SourceEnvelope::Upsert(UpsertEnvelope {
1351                            source_arity: 2,
1352                            style: UpsertStyle::Default(KeyEnvelope::Flattened),
1353                            key_indices: vec![0],
1354                        });
1355                        let source_statistics = SourceStatistics::new(
1356                            source_id,
1357                            0,
1358                            &source_statistics_defs,
1359                            source_id,
1360                            &ShardId::new(),
1361                            envelope,
1362                            Antichain::from_elem(Timestamp::minimum()),
1363                        );
1364
1365                        let source_config = SourceExportCreationConfig {
1366                            id: GlobalId::User(0),
1367                            worker_id: 0,
1368                            metrics: storage_metrics,
1369                            source_statistics,
1370                        };
1371
1372                        let rocksdb_init_fn = move || async move {
1373                            let merge_operator = Some((
1374                                "upsert_state_snapshot_merge_v1".to_string(),
1375                                |a: &[u8],
1376                                 b: ValueIterator<
1377                                    BincodeOpts,
1378                                    StateValue<(MzTimestamp, Subtime), u64>,
1379                                >| {
1380                                    consolidating_merge_function::<(MzTimestamp, Subtime), u64>(
1381                                        a.into(),
1382                                        b,
1383                                    )
1384                                },
1385                            ));
1386                            let rocksdb_cleanup_tries = 5;
1387                            let tuning = RocksDBConfig::new(Default::default(), None);
1388                            let mut rocksdb_inst = mz_rocksdb::RocksDBInstance::new(
1389                                rocksdb_dir.path(),
1390                                mz_rocksdb::InstanceOptions::new(
1391                                    Env::mem_env().unwrap(),
1392                                    rocksdb_cleanup_tries,
1393                                    merge_operator,
1394                                    upsert_bincode_opts(),
1395                                ),
1396                                tuning,
1397                                rocksdb_shared_metrics,
1398                                rocksdb_instance_metrics,
1399                            )
1400                            .unwrap();
1401
1402                            let handle = rocksdb_inst.take_core_loop_handle().expect("join handle");
1403                            tx.send(handle).expect("sent joinhandle");
1404                            crate::upsert::rocksdb::RocksDB::new(rocksdb_inst)
1405                        };
1406
1407                        let (output, _, _, button) = upsert_inner(
1408                            input.as_collection(),
1409                            vec![0],
1410                            Antichain::from_elem(Timestamp::minimum()),
1411                            persist_input.as_collection(),
1412                            persist_err_input.as_collection(),
1413                            None,
1414                            upsert_metrics,
1415                            source_config,
1416                            rocksdb_init_fn,
1417                            upsert_config,
1418                            true,
1419                            None,
1420                        );
1421                        std::mem::forget(button);
1422
1423                        let (probe, stream) = output.inner.probe();
1424                        (input_handle, persist_handle, probe, stream.capture())
1425                    })
1426                });
1427
1428            let key0 = UpsertKey::from_key(Ok(&Row::pack_slice(&[Datum::Int64(0)])));
1429            let key1 = UpsertKey::from_key(Ok(&Row::pack_slice(&[Datum::Int64(1)])));
1430            let value0 = Row::pack_slice(&[Datum::Int64(0), Datum::Int64(1)]);
1431            let value1 = Row::pack_slice(&[Datum::Int64(1), Datum::Int64(2)]);
1432
1433            // External writer has advanced the feedback persist_upper to T = 10
1434            // WITHOUT the operator emitting anything itself.
1435            persist.advance_to(mz_ts(10));
1436            for _ in 0..20 {
1437                worker.step();
1438            }
1439
1440            // Lagging replacement produces source data at ts BELOW the current
1441            // persist_upper (5 and 7 while persist_upper = 10).
1442            input.send(((key0, Some(Ok(value0)), 1), mz_ts(5), Diff::ONE));
1443            input.send(((key1, Some(Ok(value1)), 2), mz_ts(7), Diff::ONE));
1444            input.advance_to(mz_ts(11));
1445            for _ in 0..20 {
1446                worker.step();
1447            }
1448
1449            (probe.with_frontier(|f| f.to_vec()), capture)
1450        });
1451
1452        let mut emitted: Vec<_> = capture
1453            .extract()
1454            .into_iter()
1455            .flat_map(|(_cap, c)| c)
1456            .collect();
1457        differential_dataflow::consolidation::consolidate_updates(&mut emitted);
1458
1459        // v1 drops the below-upper data and lets its output frontier advance
1460        // freely (here to the input upper, 11): nothing is stranded, and the
1461        // frontier is never pinned below the shard upper (10).
1462        assert!(emitted.is_empty(), "v1 emitted {emitted:?}");
1463        assert_eq!(
1464            frontier,
1465            vec![mz_ts(11)],
1466            "v1 output frontier should advance to the input upper, not pin below \
1467             persist_upper as v2 does"
1468        );
1469        assert!(
1470            frontier[0] >= mz_ts(10),
1471            "v1 output frontier {frontier:?} should reach at least persist_upper (10)"
1472        );
1473
1474        while let Ok(handle) = rx.recv() {
1475            handle.join().expect("threads completed successfully");
1476        }
1477    }
1478}