Skip to main content

mz_storage/
upsert.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
10use std::cell::RefCell;
11use std::cmp::Reverse;
12use std::convert::AsRef;
13use std::fmt::Debug;
14use std::hash::{Hash, Hasher};
15use std::path::PathBuf;
16use std::sync::Arc;
17
18use differential_dataflow::hashable::Hashable;
19use differential_dataflow::{AsCollection, VecCollection};
20use futures::StreamExt;
21use futures::future::FutureExt;
22use indexmap::map::Entry;
23use itertools::Itertools;
24use mz_ore::error::ErrorExt;
25use mz_repr::{Datum, DatumVec, Diff, GlobalId, Row};
26use mz_rocksdb::ValueIterator;
27use mz_sql_server_util::cdc::Lsn;
28use mz_storage_types::configuration::StorageConfiguration;
29use mz_storage_types::dyncfgs;
30use mz_storage_types::errors::{DataflowError, EnvelopeError, UpsertError};
31use mz_storage_types::sources::MzOffset;
32use mz_storage_types::sources::envelope::UpsertEnvelope;
33use mz_storage_types::sources::kafka::{KafkaTimestamp, RangeBound};
34use mz_storage_types::sources::mysql::GtidPartition;
35use mz_timely_util::builder_async::{
36    AsyncOutputHandle, Event as AsyncEvent, OperatorBuilder as AsyncOperatorBuilder,
37    PressOnDropButton,
38};
39use serde::{Deserialize, Serialize};
40use sha2::{Digest, Sha256};
41use timely::dataflow::channels::pact::Exchange;
42use timely::dataflow::operators::{Capability, InputCapability, Operator};
43use timely::dataflow::{Scope, StreamVec};
44use timely::order::{PartialOrder, TotalOrder};
45use timely::progress::timestamp::Refines;
46use timely::progress::{Antichain, Timestamp};
47
48use crate::healthcheck::HealthStatusUpdate;
49use crate::metrics::upsert::{UpsertBackpressureMetrics, UpsertMetrics};
50use crate::storage_state::StorageInstanceContext;
51use crate::{upsert_continual_feedback, upsert_continual_feedback_v2};
52use types::{
53    BincodeOpts, StateValue, UpsertState, UpsertStateBackend, consolidating_merge_function,
54    upsert_bincode_opts,
55};
56
57#[cfg(any(test, feature = "fuzzing"))]
58pub mod memory;
59pub(crate) mod rocksdb;
60// TODO(aljoscha): Move next to upsert module, rename to upsert_types.
61pub(crate) mod types;
62
63pub type UpsertValue = Result<Row, Box<UpsertError>>;
64
65#[derive(
66    Copy,
67    Clone,
68    Hash,
69    PartialEq,
70    Eq,
71    PartialOrd,
72    Ord,
73    Serialize,
74    Deserialize,
75    bytemuck::AnyBitPattern,
76    bytemuck::NoUninit
77)]
78#[repr(transparent)]
79pub struct UpsertKey([u8; 32]);
80
81impl columnation::Columnation for UpsertKey {
82    type InnerRegion = columnation::CopyRegion<UpsertKey>;
83}
84
85/// Columnar (the `columnar` crate, distinct from `columnation`) support for
86/// `UpsertKey`, so the upsert-v2 source stash can use a paged columnar merge
87/// batcher keyed natively by `UpsertKey` (no `Row` packing). `UpsertKey` is a
88/// POD `[u8; 32]` newtype, so the container is a fixed-stride byte column.
89///
90/// This is hand-rolled rather than `#[derive(Columnar)]`d for one load-bearing
91/// reason: the reference type must be `&UpsertKey`. `&UpsertKey` is `Copy + Ord`
92/// (the lexicographic `[u8; 32]` order the persist-feedback trace is keyed on),
93/// which both satisfies the merge batcher's `Ref: Copy + Ord` requirement and —
94/// crucially — matches the read item of the feedback arrangement's
95/// `ColumnationStack<UpsertKey>` key container, so the paged `ValRow` builder
96/// can reconcile the `Column` input against the spine (its `BuilderInput` bound
97/// is `ReadItem: PartialEq<Ref<UpsertKey>>`). A derived impl would yield a
98/// generated `UpsertKeyReference` (and route `[u8; 32]` through the generic,
99/// non-fixed-stride array container), breaking that reconciliation.
100mod columnar_upsert_key {
101    use super::UpsertKey;
102    use columnar::Columnar;
103    use mz_ore::cast::CastFrom;
104    use std::ops::Range;
105
106    /// A newtype wrapper for a vector of `UpsertKey` values.
107    #[derive(Clone, Copy, Default, Debug)]
108    pub struct UpsertKeys<T>(T);
109    impl<D, T: columnar::Push<D>> columnar::Push<D> for UpsertKeys<T> {
110        #[inline(always)]
111        fn push(&mut self, item: D) {
112            self.0.push(item)
113        }
114    }
115    impl<T: columnar::Clear> columnar::Clear for UpsertKeys<T> {
116        #[inline(always)]
117        fn clear(&mut self) {
118            self.0.clear()
119        }
120    }
121    impl<T: columnar::Len> columnar::Len for UpsertKeys<T> {
122        #[inline(always)]
123        fn len(&self) -> usize {
124            self.0.len()
125        }
126    }
127    impl<'a> columnar::Index for UpsertKeys<&'a [UpsertKey]> {
128        type Ref = &'a UpsertKey;
129
130        #[inline(always)]
131        fn get(&self, index: usize) -> Self::Ref {
132            &self.0[index]
133        }
134    }
135
136    impl Columnar for UpsertKey {
137        #[inline(always)]
138        fn into_owned<'a>(other: columnar::Ref<'a, Self>) -> Self {
139            *other
140        }
141        type Container = UpsertKeys<Vec<UpsertKey>>;
142        #[inline(always)]
143        fn reborrow<'b, 'a: 'b>(thing: columnar::Ref<'a, Self>) -> columnar::Ref<'b, Self>
144        where
145            Self: 'a,
146        {
147            thing
148        }
149    }
150
151    impl columnar::Borrow for UpsertKeys<Vec<UpsertKey>> {
152        type Ref<'a> = &'a UpsertKey;
153        type Borrowed<'a>
154            = UpsertKeys<&'a [UpsertKey]>
155        where
156            Self: 'a;
157        #[inline(always)]
158        fn borrow<'a>(&'a self) -> Self::Borrowed<'a> {
159            UpsertKeys(self.0.as_slice())
160        }
161        #[inline(always)]
162        fn reborrow<'b, 'a: 'b>(item: Self::Borrowed<'a>) -> Self::Borrowed<'b>
163        where
164            Self: 'a,
165        {
166            UpsertKeys(item.0)
167        }
168        #[inline(always)]
169        fn reborrow_ref<'b, 'a: 'b>(item: Self::Ref<'a>) -> Self::Ref<'b>
170        where
171            Self: 'a,
172        {
173            item
174        }
175    }
176
177    impl columnar::Container for UpsertKeys<Vec<UpsertKey>> {
178        #[inline(always)]
179        fn extend_from_self(&mut self, other: Self::Borrowed<'_>, range: Range<usize>) {
180            self.0.extend_from_self(other.0, range)
181        }
182        #[inline(always)]
183        fn reserve_for<'a, I>(&mut self, selves: I)
184        where
185            Self: 'a,
186            I: Iterator<Item = Self::Borrowed<'a>> + Clone,
187        {
188            self.0.reserve_for(selves.map(|s| s.0));
189        }
190    }
191
192    impl<'a> columnar::AsBytes<'a> for UpsertKeys<&'a [UpsertKey]> {
193        const SLICE_COUNT: usize = 1;
194        #[inline(always)]
195        fn get_byte_slice(&self, index: usize) -> (u64, &'a [u8]) {
196            mz_ore::soft_assert_no_log!(index < Self::SLICE_COUNT);
197            (
198                u64::cast_from(align_of::<UpsertKey>()),
199                bytemuck::cast_slice(self.0),
200            )
201        }
202        #[inline(always)]
203        fn as_bytes(&self) -> impl Iterator<Item = (u64, &'a [u8])> {
204            std::iter::once((
205                u64::cast_from(align_of::<UpsertKey>()),
206                bytemuck::cast_slice(self.0),
207            ))
208        }
209    }
210    impl<'a> columnar::FromBytes<'a> for UpsertKeys<&'a [UpsertKey]> {
211        const SLICE_COUNT: usize = 1;
212        #[inline(always)]
213        fn from_bytes(bytes: &mut impl Iterator<Item = &'a [u8]>) -> Self {
214            UpsertKeys(bytemuck::cast_slice(
215                bytes.next().expect("Iterator exhausted prematurely"),
216            ))
217        }
218    }
219}
220
221/// Projects a source's native `FromTime` to a columnar, totally-ordered key
222/// used by the upsert source stash to keep the latest update per `(key, time)`.
223///
224/// The upsert stash is a paged columnar merge batcher; its diff carries this
225/// projection rather than the raw `FromTime`, so the only columnar type the
226/// stash needs is `Order` — never the (possibly structurally complex) source
227/// timestamp. This is what keeps the columnar requirement off the generic
228/// source-render path: that path only ever needs `FromTime: UpsertSourceTime`.
229/// Only the relative order matters; the value is never read back.
230///
231/// The upsert envelope is rendered for Kafka and the KEY VALUE load generator,
232/// so those source times (`KafkaTimestamp`, `MzOffset`) project to a real order
233/// key. The remaining source times implement the trait only for coherence on
234/// the generic render path — their sources never render upsert — so their
235/// projection is a panicking guard rather than a real key.
236pub trait UpsertSourceTime
237where
238    for<'a> columnar::Ref<'a, Self::Order>: Ord,
239{
240    /// Columnar order key. Must order consistently with the source time it is
241    /// projected from.
242    type Order: columnar::Columnar + Clone + Default + Ord + Send + Sync + 'static;
243    /// Project the source time onto its order key.
244    fn upsert_order(&self) -> Self::Order;
245}
246
247impl UpsertSourceTime for KafkaTimestamp {
248    /// Per-record Kafka source times are exact singletons (a single partition
249    /// at a single offset; see the source reader), and `KafkaTimestamp`'s
250    /// derived `Ord` is lexicographic on `(partition, offset)`, so this flat
251    /// projection is order-preserving. `RangeBound`'s infinities map to the
252    /// `i64` extrema to remain order-consistent for any non-singleton bound.
253    type Order = (i64, u64);
254    fn upsert_order(&self) -> (i64, u64) {
255        let partition = match self.interval().lower {
256            RangeBound::NegInfinity => i64::MIN,
257            RangeBound::Elem(p, _) => i64::from(p),
258            RangeBound::PosInfinity => i64::MAX,
259        };
260        (partition, self.timestamp().offset)
261    }
262}
263
264/// Load-generator (and Postgres) source time. The KEY VALUE load generator is
265/// the one non-Kafka source that renders the upsert envelope (see
266/// `apply_source_envelope_encoding` in the planner), so this projects to the
267/// record offset: offsets increase with each update, so "max order wins" is
268/// exactly "latest update wins" for dedup.
269impl UpsertSourceTime for MzOffset {
270    type Order = u64;
271    fn upsert_order(&self) -> u64 {
272        self.offset
273    }
274}
275
276/// Source times whose sources never render the upsert envelope (MySQL and SQL
277/// Server CDC). `Order = ()` keeps the generic render path free of any columnar
278/// requirement on the source time. The projection panics rather than returning:
279/// `()` would collapse every `from_time` to equal, silently breaking "latest
280/// offset wins" dedup, so if such a source ever reaches the upsert path we want
281/// a loud failure, not arbitrary per-key output.
282macro_rules! upsert_source_time_unit {
283    ($($ty:ty),+ $(,)?) => {$(
284        impl UpsertSourceTime for $ty {
285            type Order = ();
286            fn upsert_order(&self) {
287                unreachable!(
288                    "upsert source stash is not rendered for this source, but \
289                     {} reached the projection",
290                    std::any::type_name::<Self>(),
291                )
292            }
293        }
294    )+};
295}
296upsert_source_time_unit!(GtidPartition, Lsn);
297
298/// Storage's leg of the process-wide chunk spill gate, used by the chunked
299/// upsert-v2 stash flavor.
300///
301/// In that flavor the source stash and feedback arrangement spill through the
302/// process buffer pool ([`mz_timely_util::columnar::chunk`]): committed chunk
303/// bodies land in the pool once compute's config handler has installed and
304/// budgeted it (storage and compute run in the same `clusterd` process).
305///
306/// The gate is process-wide with one leg per subsystem, and chunks spill
307/// while either leg is set. Storage sets its leg from
308/// `enable_upsert_paged_spill`, so that flag alone cannot veto spilling
309/// enabled by compute's leg. The gate is consulted at every chunk commit, so
310/// flips apply to running dataflows.
311pub mod upsert_stash_spill {
312    /// Enable or disable spilling of upsert chunk bodies to the buffer pool.
313    pub fn set_enabled(enabled: bool) {
314        mz_timely_util::columnar::chunk::set_storage_spill_enabled(enabled);
315    }
316}
317
318/// Pager for the paged upsert-v2 stash flavor.
319///
320/// This draws from the same process-wide [`TieredPolicy`] budget pool as the
321/// compute column-paged batcher — there is one budget and one underlying
322/// `mz_ore::pager` — but whether the stash *uses* it is gated by storage's own
323/// `enable_upsert_paged_spill` flag, independently of compute's
324/// `enable_column_paged_batcher_spill`. The shared pool's budget / backend /
325/// codec are configured by compute's `apply_tiered_config` (storage and compute
326/// run in the same `clusterd` process).
327///
328/// Flipping the flag takes effect on dataflows created after the change: the
329/// paged flavor captures the pager once at operator construction.
330///
331/// [`TieredPolicy`]: mz_timely_util::column_pager::policy::TieredPolicy
332pub mod upsert_stash_pager {
333    use std::sync::{LazyLock, RwLock};
334
335    use mz_timely_util::column_pager::{ColumnPager, shared_pager};
336
337    /// Active pager handed to upsert source-stash batchers. Defaults to
338    /// disabled (every chunk resident) until [`set_enabled`] turns it on.
339    static PAGER: LazyLock<RwLock<ColumnPager>> =
340        LazyLock::new(|| RwLock::new(ColumnPager::disabled()));
341
342    /// Enable or disable the stash's use of the shared column pager. When
343    /// enabled, the stash spills through the shared budget pool; when disabled
344    /// it keeps every chunk resident.
345    pub fn set_enabled(enabled: bool) {
346        *PAGER.write().expect("upsert stash pager poisoned") = shared_pager(enabled);
347    }
348
349    /// The current upsert-stash pager. Cheap: clones the inner `Arc`.
350    pub fn pager() -> ColumnPager {
351        PAGER.read().expect("upsert stash pager poisoned").clone()
352    }
353}
354
355impl Debug for UpsertKey {
356    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
357        write!(f, "0x")?;
358        for byte in self.0 {
359            write!(f, "{:02x}", byte)?;
360        }
361        Ok(())
362    }
363}
364
365impl AsRef<[u8]> for UpsertKey {
366    #[inline(always)]
367    // Note we do 1 `multi_get` and 1 `multi_put` while processing a _batch of updates_. Within the
368    // batch, we effectively consolidate each key, before persisting that consolidated value.
369    // Easy!!
370    fn as_ref(&self) -> &[u8] {
371        &self.0
372    }
373}
374
375impl From<&[u8]> for UpsertKey {
376    fn from(bytes: &[u8]) -> Self {
377        UpsertKey(bytes.try_into().expect("invalid key length"))
378    }
379}
380
381/// The hash function used to map upsert keys. It is important that this hash is a cryptographic
382/// hash so that there is no risk of collisions. Collisions on SHA256 have a probability of 2^128
383/// which is many orders of magnitude smaller than many other events that we don't even think about
384/// (e.g bit flips). In short, we can safely assume that sha256(a) == sha256(b) iff a == b.
385type KeyHash = Sha256;
386
387impl UpsertKey {
388    pub fn from_key(key: Result<&Row, &UpsertError>) -> Self {
389        Self::from_iter(key.map(|r| r.iter()))
390    }
391
392    pub fn from_value(value: Result<&Row, &UpsertError>, key_indices: &[usize]) -> Self {
393        thread_local! {
394            /// A thread-local datum cache used to calculate hashes
395            static VALUE_DATUMS: RefCell<DatumVec> = RefCell::new(DatumVec::new());
396        }
397        VALUE_DATUMS.with(|value_datums| {
398            let mut value_datums = value_datums.borrow_mut();
399            let value = value.map(|v| value_datums.borrow_with(v));
400            let key = match value {
401                Ok(ref datums) => Ok(key_indices.iter().map(|&idx| datums[idx])),
402                Err(err) => Err(err),
403            };
404            Self::from_iter(key)
405        })
406    }
407
408    pub fn from_iter<'a, 'b>(
409        key: Result<impl Iterator<Item = Datum<'a>> + 'b, &UpsertError>,
410    ) -> Self {
411        thread_local! {
412            /// A thread-local datum cache used to calculate hashes
413            static KEY_DATUMS: RefCell<DatumVec> = RefCell::new(DatumVec::new());
414        }
415        KEY_DATUMS.with(|key_datums| {
416            let mut key_datums = key_datums.borrow_mut();
417            // Borrowing the DatumVec gives us a temporary buffer to store datums in that will be
418            // automatically cleared on Drop. See the DatumVec docs for more details.
419            let mut key_datums = key_datums.borrow();
420            let key: Result<&[Datum], Datum> = match key {
421                Ok(key) => {
422                    for datum in key {
423                        key_datums.push(datum);
424                    }
425                    Ok(&*key_datums)
426                }
427                Err(UpsertError::Value(err)) => {
428                    key_datums.extend(err.for_key.iter());
429                    Ok(&*key_datums)
430                }
431                Err(UpsertError::KeyDecode(err)) => Err(Datum::Bytes(&err.raw)),
432                Err(UpsertError::NullKey(_)) => Err(Datum::Null),
433            };
434            let mut hasher = DigestHasher(KeyHash::new());
435            key.hash(&mut hasher);
436            Self(hasher.0.finalize().into())
437        })
438    }
439}
440
441struct DigestHasher<H: Digest>(H);
442
443impl<H: Digest> Hasher for DigestHasher<H> {
444    fn write(&mut self, bytes: &[u8]) {
445        self.0.update(bytes);
446    }
447
448    fn finish(&self) -> u64 {
449        panic!("digest wrapper used to produce a hash");
450    }
451}
452
453use std::convert::Infallible;
454use timely::container::CapacityContainerBuilder;
455use timely::dataflow::channels::pact::Pipeline;
456
457use self::types::ValueMetadata;
458
459/// This leaf operator drops `token` after the input reaches the `resume_upper`.
460/// This is useful to take coordinated actions across all workers, after the `upsert`
461/// operator has rehydrated.
462pub fn rehydration_finished<'scope, T: Timestamp>(
463    scope: Scope<'scope, T>,
464    source_config: &crate::source::RawSourceCreationConfig,
465    // A token that we can drop to signal we are finished rehydrating.
466    token: impl std::any::Any + 'static,
467    resume_upper: Antichain<T>,
468    input: StreamVec<'scope, T, Infallible>,
469) {
470    let worker_id = source_config.worker_id;
471    let id = source_config.id;
472    let mut builder = AsyncOperatorBuilder::new(format!("rehydration_finished({id}"), scope);
473    let mut input = builder.new_disconnected_input(input, Pipeline);
474
475    builder.build(move |_capabilities| async move {
476        let mut input_upper = Antichain::from_elem(Timestamp::minimum());
477        // Ensure this operator finishes if the resume upper is `[0]`
478        while !PartialOrder::less_equal(&resume_upper, &input_upper) {
479            let Some(event) = input.next().await else {
480                break;
481            };
482            if let AsyncEvent::Progress(upper) = event {
483                input_upper = upper;
484            }
485        }
486        tracing::info!(
487            %worker_id,
488            source_id = %id,
489            "upsert source has downgraded past the resume upper ({resume_upper:?}) across all workers",
490        );
491        drop(token);
492    });
493}
494
495/// Keys this operator's previous output, read back from persist, for retraction.
496///
497/// Only [`UpsertError`] survives the error side. It is the one error this operator can have
498/// written, and therefore the one it can retract; anything else entered the shard from
499/// elsewhere and is not ours to take back.
500pub(crate) fn key_persist_feedback<'scope, T>(
501    ok: VecCollection<'scope, T, Row, Diff>,
502    err: VecCollection<'scope, T, DataflowError, Diff>,
503    key_indices: Vec<usize>,
504) -> VecCollection<'scope, T, (UpsertKey, UpsertValue), Diff>
505where
506    T: Timestamp,
507{
508    let keyed_ok = {
509        let key_indices = key_indices.clone();
510        ok.map(move |row| {
511            let key = UpsertKey::from_value(Ok(&row), &key_indices);
512            (key, Ok(row))
513        })
514    };
515    let keyed_err = err.flat_map(move |err| {
516        let err = match err {
517            DataflowError::EnvelopeError(err) => match *err {
518                EnvelopeError::Upsert(err) => Box::new(err),
519                EnvelopeError::Flat(_) => return None,
520            },
521            _ => return None,
522        };
523        let key = UpsertKey::from_value(Err(&err), &key_indices);
524        Some((key, Err(err)))
525    });
526    keyed_ok.concat(keyed_err)
527}
528
529/// Resumes an upsert computation at `resume_upper` given as inputs a collection of upsert commands
530/// and the collection of the previous output of this operator.
531/// Returns a tuple of
532/// - A collection of the computed upsert operator and,
533/// - A health update stream to propagate errors
534pub(crate) fn upsert<'scope, T, FromTime>(
535    input: VecCollection<'scope, T, (UpsertKey, Option<UpsertValue>, FromTime), Diff>,
536    upsert_envelope: UpsertEnvelope,
537    resume_upper: Antichain<T>,
538    previous_ok: VecCollection<'scope, T, Row, Diff>,
539    previous_err: VecCollection<'scope, T, DataflowError, Diff>,
540    previous_token: Option<Vec<PressOnDropButton>>,
541    source_config: crate::source::SourceExportCreationConfig,
542    instance_context: &StorageInstanceContext,
543    storage_configuration: &StorageConfiguration,
544    dataflow_paramters: &crate::internal_control::DataflowParameters,
545    backpressure_metrics: Option<UpsertBackpressureMetrics>,
546) -> (
547    VecCollection<'scope, T, Result<Row, DataflowError>, Diff>,
548    StreamVec<'scope, T, (Option<GlobalId>, HealthStatusUpdate)>,
549    StreamVec<'scope, T, Infallible>,
550    PressOnDropButton,
551)
552where
553    T: Timestamp + TotalOrder + Sync,
554    T: Refines<mz_repr::Timestamp> + TotalOrder + Sync,
555    FromTime: Timestamp + Clone + Sync,
556{
557    let upsert_metrics = source_config.metrics.get_upsert_metrics(
558        source_config.id,
559        source_config.worker_id,
560        backpressure_metrics,
561    );
562
563    let rocksdb_cleanup_tries =
564        dyncfgs::STORAGE_ROCKSDB_CLEANUP_TRIES.get(storage_configuration.config_set());
565
566    // Whether or not to partially drain the input buffer
567    // to prevent buffering of the _upstream_ snapshot.
568    let prevent_snapshot_buffering =
569        dyncfgs::STORAGE_UPSERT_PREVENT_SNAPSHOT_BUFFERING.get(storage_configuration.config_set());
570    // If the above is true, the number of timely batches to process at once.
571    let snapshot_buffering_max = dyncfgs::STORAGE_UPSERT_MAX_SNAPSHOT_BATCH_BUFFERING
572        .get(storage_configuration.config_set());
573
574    // Whether we should provide the upsert state merge operator to the RocksDB instance
575    // (for faster performance during snapshot hydration).
576    let rocksdb_use_native_merge_operator =
577        dyncfgs::STORAGE_ROCKSDB_USE_MERGE_OPERATOR.get(storage_configuration.config_set());
578
579    let upsert_config = UpsertConfig {
580        shrink_upsert_unused_buffers_by_ratio: storage_configuration
581            .parameters
582            .shrink_upsert_unused_buffers_by_ratio,
583    };
584
585    let thin_input = upsert_thinning(input);
586
587    let tuning = dataflow_paramters.upsert_rocksdb_tuning_config.clone();
588
589    // When running RocksDB in memory, the file system is emulated. However, we still need to
590    // pick a path that exists because RocksDB will attempt to create the working directory
591    // (see https://github.com/rust-rocksdb/rust-rocksdb/issues/1015) and write a lock file,
592    // so we need to ensure the directory is unique per worker.
593    let rocksdb_dir = instance_context
594        .scratch_directory
595        .clone()
596        .unwrap_or_else(|| PathBuf::from("/tmp"))
597        .join("storage")
598        .join("upsert")
599        .join(source_config.id.to_string())
600        .join(source_config.worker_id.to_string());
601
602    tracing::info!(
603        worker_id = %source_config.worker_id,
604        source_id = %source_config.id,
605        ?rocksdb_dir,
606        ?tuning,
607        ?rocksdb_use_native_merge_operator,
608        "rendering upsert source"
609    );
610
611    let rocksdb_shared_metrics = Arc::clone(&upsert_metrics.rocksdb_shared);
612    let rocksdb_instance_metrics = Arc::clone(&upsert_metrics.rocksdb_instance_metrics);
613
614    let env = instance_context
615        .rocksdb_env()
616        .expect("failed to create rocksdb env");
617
618    // A closure that will initialize and return a configured RocksDB instance
619    let rocksdb_init_fn = move || async move {
620        let merge_operator = if rocksdb_use_native_merge_operator {
621            Some((
622                "upsert_state_snapshot_merge_v1".to_string(),
623                |a: &[u8], b: ValueIterator<BincodeOpts, StateValue<T, FromTime>>| {
624                    consolidating_merge_function::<T, FromTime>(a.into(), b)
625                },
626            ))
627        } else {
628            None
629        };
630        rocksdb::RocksDB::new(
631            mz_rocksdb::RocksDBInstance::new(
632                &rocksdb_dir,
633                mz_rocksdb::InstanceOptions::new(
634                    env,
635                    rocksdb_cleanup_tries,
636                    merge_operator,
637                    // For now, just use the same config as the one used for
638                    // merging snapshots.
639                    upsert_bincode_opts(),
640                ),
641                tuning,
642                rocksdb_shared_metrics,
643                rocksdb_instance_metrics,
644            )
645            .unwrap(),
646        )
647    };
648
649    upsert_operator(
650        thin_input,
651        upsert_envelope.key_indices,
652        resume_upper,
653        previous_ok,
654        previous_err,
655        previous_token,
656        upsert_metrics,
657        source_config,
658        rocksdb_init_fn,
659        upsert_config,
660        storage_configuration,
661        prevent_snapshot_buffering,
662        snapshot_buffering_max,
663    )
664}
665
666/// An experimental upsert implementation loosely described in this doc:
667/// [Upsert V2 Much Simpler Boogaloo](https://www.notion.so/materialize/Upsert-V2-Much-Simpler-Boogaloo-31913f48d37b807fa88bdeafc27c02d9?source=copy_link)
668///
669/// Instead of using rocksdb as a state backend, this implementation uses a differential dataflow collection to hold the key state,
670/// and performs consolidation of updates with matching keys and MZ timestamps, using max FromTime to choose winners,
671/// resulting in only one record per key per time.
672pub(crate) fn upsert_v2<'scope, T, FromTime>(
673    input: VecCollection<'scope, T, (UpsertKey, Option<UpsertValue>, FromTime), Diff>,
674    upsert_envelope: UpsertEnvelope,
675    resume_upper: Antichain<T>,
676    previous_ok: VecCollection<'scope, T, Row, Diff>,
677    previous_err: VecCollection<'scope, T, DataflowError, Diff>,
678    previous_token: Option<Vec<PressOnDropButton>>,
679    source_config: crate::source::SourceExportCreationConfig,
680    backpressure_metrics: Option<UpsertBackpressureMetrics>,
681    stash_flavor: upsert_continual_feedback_v2::UpsertStashFlavor,
682) -> (
683    VecCollection<'scope, T, Result<Row, DataflowError>, Diff>,
684    StreamVec<'scope, T, (Option<GlobalId>, HealthStatusUpdate)>,
685    StreamVec<'scope, T, Infallible>,
686    PressOnDropButton,
687)
688where
689    T: Timestamp + TotalOrder + Sync,
690    T: Refines<mz_repr::Timestamp> + differential_dataflow::lattice::Lattice,
691    T: columnation::Columnation,
692    T: columnar::Columnar + Default,
693    for<'a> columnar::Ref<'a, T>: Copy + Ord,
694    FromTime: Timestamp + Clone + Sync,
695    FromTime: UpsertSourceTime,
696{
697    let upsert_metrics = source_config.metrics.get_upsert_metrics(
698        source_config.id,
699        source_config.worker_id,
700        backpressure_metrics,
701    );
702
703    let thin_input = upsert_thinning(input);
704
705    tracing::info!(
706        worker_id = %source_config.worker_id,
707        source_id = %source_config.id,
708        ?stash_flavor,
709        "rendering upsert source (btreemap backend)"
710    );
711
712    upsert_continual_feedback_v2::upsert_inner(
713        stash_flavor,
714        thin_input,
715        upsert_envelope.key_indices,
716        resume_upper,
717        previous_ok,
718        previous_err,
719        previous_token,
720        upsert_metrics,
721        source_config,
722    )
723}
724
725// A shim so we can dispatch based on the dyncfg that tells us which upsert
726// operator to use.
727fn upsert_operator<'scope, T, FromTime, F, Fut, US>(
728    input: VecCollection<'scope, T, (UpsertKey, Option<UpsertValue>, FromTime), Diff>,
729    key_indices: Vec<usize>,
730    resume_upper: Antichain<T>,
731    persist_ok: VecCollection<'scope, T, Row, Diff>,
732    persist_err: VecCollection<'scope, T, DataflowError, Diff>,
733    persist_token: Option<Vec<PressOnDropButton>>,
734    upsert_metrics: UpsertMetrics,
735    source_config: crate::source::SourceExportCreationConfig,
736    state: F,
737    upsert_config: UpsertConfig,
738    _storage_configuration: &StorageConfiguration,
739    prevent_snapshot_buffering: bool,
740    snapshot_buffering_max: Option<usize>,
741) -> (
742    VecCollection<'scope, T, Result<Row, DataflowError>, Diff>,
743    StreamVec<'scope, T, (Option<GlobalId>, HealthStatusUpdate)>,
744    StreamVec<'scope, T, Infallible>,
745    PressOnDropButton,
746)
747where
748    T: Timestamp + TotalOrder + Sync,
749    T: Refines<mz_repr::Timestamp> + TotalOrder + Sync,
750    F: FnOnce() -> Fut + 'static,
751    Fut: std::future::Future<Output = US>,
752    US: UpsertStateBackend<T, FromTime>,
753    FromTime: Debug + timely::ExchangeData + Clone + Ord + Sync,
754{
755    // Hard-coded to true because classic UPSERT cannot be used safely with
756    // concurrent ingestions, which we need for both 0dt upgrades and
757    // multi-replica ingestions.
758    let use_continual_feedback_upsert = true;
759
760    tracing::info!(id = %source_config.id, %use_continual_feedback_upsert, "upsert operator implementation");
761
762    if use_continual_feedback_upsert {
763        upsert_continual_feedback::upsert_inner(
764            input,
765            key_indices,
766            resume_upper,
767            persist_ok,
768            persist_err,
769            persist_token,
770            upsert_metrics,
771            source_config,
772            state,
773            upsert_config,
774            prevent_snapshot_buffering,
775            snapshot_buffering_max,
776        )
777    } else {
778        upsert_classic(
779            input,
780            key_indices,
781            resume_upper,
782            persist_ok,
783            persist_err,
784            persist_token,
785            upsert_metrics,
786            source_config,
787            state,
788            upsert_config,
789            prevent_snapshot_buffering,
790            snapshot_buffering_max,
791        )
792    }
793}
794
795/// Renders an operator that discards updates that are known to not affect the outcome of upsert in
796/// a streaming fashion. For each distinct (key, time) in the input it emits the value with the
797/// highest from_time. Its purpose is to thin out data as much as possible before exchanging them
798/// across workers.
799fn upsert_thinning<'scope, T, K, V, FromTime>(
800    input: VecCollection<'scope, T, (K, V, FromTime), Diff>,
801) -> VecCollection<'scope, T, (K, V, FromTime), Diff>
802where
803    T: Timestamp + TotalOrder,
804    K: timely::ExchangeData + Clone + Eq + Ord,
805    V: timely::ExchangeData + Clone,
806    FromTime: Timestamp,
807{
808    input
809        .inner
810        .unary(Pipeline, "UpsertThinning", |_, _| {
811            // A capability suitable to emit all updates in `updates`, if any.
812            let mut capability: Option<InputCapability<T>> = None;
813            // A batch of received updates
814            let mut updates = Vec::new();
815            move |input, output| {
816                input.for_each(|cap, data| {
817                    assert!(
818                        data.iter().all(|(_, _, diff)| diff.is_positive()),
819                        "invalid upsert input"
820                    );
821                    updates.append(data);
822                    match capability.as_mut() {
823                        Some(capability) => {
824                            if cap.time() <= capability.time() {
825                                *capability = cap;
826                            }
827                        }
828                        None => capability = Some(cap),
829                    }
830                });
831                if let Some(capability) = capability.take() {
832                    // Sort by (key, time, Reverse(from_time)) so that deduping by (key, time) gives
833                    // the latest change for that key.
834                    updates.sort_unstable_by(|a, b| {
835                        let ((key1, _, from_time1), time1, _) = a;
836                        let ((key2, _, from_time2), time2, _) = b;
837                        Ord::cmp(
838                            &(key1, time1, Reverse(from_time1)),
839                            &(key2, time2, Reverse(from_time2)),
840                        )
841                    });
842                    let mut session = output.session(&capability);
843                    session.give_iterator(updates.drain(..).dedup_by(|a, b| {
844                        let ((key1, _, _), time1, _) = a;
845                        let ((key2, _, _), time2, _) = b;
846                        (key1, time1) == (key2, time2)
847                    }))
848                }
849            }
850        })
851        .as_collection()
852}
853
854/// Helper method for `upsert_classic` used to stage `data` updates
855/// from the input/source timely edge.
856fn stage_input<T, FromTime>(
857    stash: &mut Vec<(T, UpsertKey, Reverse<FromTime>, Option<UpsertValue>)>,
858    data: &mut Vec<((UpsertKey, Option<UpsertValue>, FromTime), T, Diff)>,
859    input_upper: &Antichain<T>,
860    resume_upper: &Antichain<T>,
861    storage_shrink_upsert_unused_buffers_by_ratio: usize,
862) where
863    T: PartialOrder,
864    FromTime: Ord,
865{
866    if PartialOrder::less_equal(input_upper, resume_upper) {
867        data.retain(|(_, ts, _)| resume_upper.less_equal(ts));
868    }
869
870    stash.extend(data.drain(..).map(|((key, value, order), time, diff)| {
871        assert!(diff.is_positive(), "invalid upsert input");
872        (time, key, Reverse(order), value)
873    }));
874
875    if storage_shrink_upsert_unused_buffers_by_ratio > 0 {
876        let reduced_capacity = stash.capacity() / storage_shrink_upsert_unused_buffers_by_ratio;
877        if reduced_capacity > stash.len() {
878            stash.shrink_to(reduced_capacity);
879        }
880    }
881}
882
883/// The style of drain we are performing on the stash. `AtTime`-drains cannot
884/// assume that all values have been seen, and must leave tombstones behind for deleted values.
885#[derive(Debug)]
886enum DrainStyle<'a, T> {
887    ToUpper(&'a Antichain<T>),
888    AtTime(T),
889}
890
891/// Helper method for `upsert_inner` used to stage `data` updates
892/// from the input timely edge.
893async fn drain_staged_input<S, T, FromTime, E>(
894    stash: &mut Vec<(T, UpsertKey, Reverse<FromTime>, Option<UpsertValue>)>,
895    commands_state: &mut indexmap::IndexMap<UpsertKey, types::UpsertValueAndSize<T, FromTime>>,
896    output_updates: &mut Vec<(UpsertValue, T, Diff)>,
897    multi_get_scratch: &mut Vec<UpsertKey>,
898    drain_style: DrainStyle<'_, T>,
899    error_emitter: &mut E,
900    state: &mut UpsertState<'_, S, T, FromTime>,
901    source_config: &crate::source::SourceExportCreationConfig,
902) where
903    S: UpsertStateBackend<T, FromTime>,
904    T: PartialOrder + Ord + Clone + Send + Sync + Serialize + Debug + 'static,
905    FromTime: timely::ExchangeData + Clone + Ord + Sync,
906    E: UpsertErrorEmitter<T>,
907{
908    stash.sort_unstable();
909
910    // Find the prefix that we can emit
911    let idx = stash.partition_point(|(ts, _, _, _)| match &drain_style {
912        DrainStyle::ToUpper(upper) => !upper.less_equal(ts),
913        DrainStyle::AtTime(time) => ts <= time,
914    });
915
916    tracing::trace!(?drain_style, updates = idx, "draining stash in upsert");
917
918    // Read the previous values _per key_ out of `state`, recording it
919    // along with the value with the _latest timestamp for that key_.
920    commands_state.clear();
921    for (_, key, _, _) in stash.iter().take(idx) {
922        commands_state.entry(*key).or_default();
923    }
924
925    // These iterators iterate in the same order because `commands_state`
926    // is an `IndexMap`.
927    multi_get_scratch.clear();
928    multi_get_scratch.extend(commands_state.iter().map(|(k, _)| *k));
929    match state
930        .multi_get(multi_get_scratch.drain(..), commands_state.values_mut())
931        .await
932    {
933        Ok(_) => {}
934        Err(e) => {
935            error_emitter
936                .emit("Failed to fetch records from state".to_string(), e)
937                .await;
938        }
939    }
940
941    // From the prefix that can be emitted we can deduplicate based on (ts, key) in
942    // order to only process the command with the maximum order within the (ts,
943    // key) group. This is achieved by wrapping order in `Reverse(FromTime)` above.;
944    let mut commands = stash.drain(..idx).dedup_by(|a, b| {
945        let ((a_ts, a_key, _, _), (b_ts, b_key, _, _)) = (a, b);
946        a_ts == b_ts && a_key == b_key
947    });
948
949    let bincode_opts = types::upsert_bincode_opts();
950    // Upsert the values into `commands_state`, by recording the latest
951    // value (or deletion). These will be synced at the end to the `state`.
952    //
953    // Note that we are effectively doing "mini-upsert" here, using
954    // `command_state`. This "mini-upsert" is seeded with data from `state`, using
955    // a single `multi_get` above, and the final state is written out into
956    // `state` using a single `multi_put`. This simplifies `UpsertStateBackend`
957    // implementations, and reduces the number of reads and write we need to do.
958    //
959    // This "mini-upsert" technique is actually useful in `UpsertState`'s
960    // `consolidate_snapshot_read_write_inner` implementation, minimizing gets and puts on
961    // the `UpsertStateBackend` implementations. In some sense, its "upsert all the way down".
962    while let Some((ts, key, from_time, value)) = commands.next() {
963        let mut command_state = if let Entry::Occupied(command_state) = commands_state.entry(key) {
964            command_state
965        } else {
966            panic!("key missing from commands_state");
967        };
968
969        let existing_value = &mut command_state.get_mut().value;
970
971        if let Some(cs) = existing_value.as_mut() {
972            cs.ensure_decoded(bincode_opts, source_config.id, Some(&key));
973        }
974
975        // Skip this command if its order key is below the one in the upsert state.
976        // Note that the existing order key may be `None` if the existing value
977        // is from snapshotting, which always sorts below new values/deletes.
978        let existing_order = existing_value
979            .as_ref()
980            .and_then(|cs| cs.provisional_order(&ts));
981        if existing_order >= Some(&from_time.0) {
982            // Skip this update. If no later updates adjust this key, then we just
983            // end up writing the same value back to state. If there
984            // is nothing in the state, `existing_order` is `None`, and this
985            // does not occur.
986            continue;
987        }
988
989        match value {
990            Some(value) => {
991                if let Some(old_value) =
992                    existing_value.replace(StateValue::finalized_value(value.clone()))
993                {
994                    if let Some(old_value) = old_value.into_decoded().finalized {
995                        output_updates.push((old_value, ts.clone(), Diff::MINUS_ONE));
996                    }
997                }
998                output_updates.push((value, ts, Diff::ONE));
999            }
1000            None => {
1001                if let Some(old_value) = existing_value.take() {
1002                    if let Some(old_value) = old_value.into_decoded().finalized {
1003                        output_updates.push((old_value, ts, Diff::MINUS_ONE));
1004                    }
1005                }
1006
1007                // Record a tombstone for deletes.
1008                *existing_value = Some(StateValue::tombstone());
1009            }
1010        }
1011    }
1012
1013    match state
1014        .multi_put(
1015            true, // Do update per-update stats.
1016            commands_state.drain(..).map(|(k, cv)| {
1017                (
1018                    k,
1019                    types::PutValue {
1020                        value: cv.value.map(|cv| cv.into_decoded()),
1021                        previous_value_metadata: cv.metadata.map(|v| ValueMetadata {
1022                            size: v.size.try_into().expect("less than i64 size"),
1023                            is_tombstone: v.is_tombstone,
1024                        }),
1025                    },
1026                )
1027            }),
1028        )
1029        .await
1030    {
1031        Ok(_) => {}
1032        Err(e) => {
1033            error_emitter
1034                .emit("Failed to update records in state".to_string(), e)
1035                .await;
1036        }
1037    }
1038}
1039
1040/// A no-op-ish error emitter for the fuzzing hook. With the in-memory backend
1041/// and the well-formed inputs the fuzzer builds, `multi_get`/`multi_put` never
1042/// error, so reaching this is itself a finding.
1043#[cfg(feature = "fuzzing")]
1044struct PanicErrorEmitter;
1045
1046#[cfg(feature = "fuzzing")]
1047#[async_trait::async_trait(?Send)]
1048impl<T> UpsertErrorEmitter<T> for PanicErrorEmitter {
1049    async fn emit(&mut self, context: String, e: anyhow::Error) {
1050        panic!("unexpected upsert state error during fuzzing: {context}: {e}");
1051    }
1052}
1053
1054/// Fuzzing hook: run a single `drain_staged_input` over `commands` (each a
1055/// `(timestamp, key, order, value)`, where `value == None` is a delete) against
1056/// a fresh empty in-memory state, draining everything strictly below
1057/// `drain_to`. Returns the emitted output updates and the final finalized value
1058/// of each key in `all_keys`. Exposed only for fuzzing. Not a stable public
1059/// API.
1060#[cfg(feature = "fuzzing")]
1061pub async fn fuzz_drain_staged_input(
1062    parts: &types::FuzzUpsertParts,
1063    source_config: &crate::source::SourceExportCreationConfig,
1064    commands: Vec<(u64, UpsertKey, u64, Option<UpsertValue>)>,
1065    drain_to: u64,
1066    all_keys: &[UpsertKey],
1067) -> (Vec<(UpsertValue, u64, Diff)>, Vec<Option<UpsertValue>>) {
1068    let mut state = parts.state();
1069    let mut stash: Vec<(u64, UpsertKey, Reverse<u64>, Option<UpsertValue>)> = commands
1070        .into_iter()
1071        .map(|(ts, key, order, value)| (ts, key, Reverse(order), value))
1072        .collect();
1073    let mut commands_state = indexmap::IndexMap::new();
1074    let mut output = Vec::new();
1075    let mut multi_get_scratch = Vec::new();
1076    let mut emitter = PanicErrorEmitter;
1077
1078    drain_staged_input(
1079        &mut stash,
1080        &mut commands_state,
1081        &mut output,
1082        &mut multi_get_scratch,
1083        DrainStyle::ToUpper(&Antichain::from_elem(drain_to)),
1084        &mut emitter,
1085        &mut state,
1086        source_config,
1087    )
1088    .await;
1089
1090    let bincode_opts = types::upsert_bincode_opts();
1091    let mut results = vec![types::UpsertValueAndSize::default(); all_keys.len()];
1092    state
1093        .multi_get(all_keys.iter().copied(), results.iter_mut())
1094        .await
1095        .expect("multi_get in fuzz hook should not error");
1096    let final_state = results
1097        .into_iter()
1098        .map(|r| match r.value {
1099            None => None,
1100            Some(mut sv) => {
1101                sv.ensure_decoded(bincode_opts, GlobalId::User(0), None);
1102                sv.into_decoded().finalized
1103            }
1104        })
1105        .collect();
1106
1107    (output, final_state)
1108}
1109
1110// Created a struct to hold the configs for upserts.
1111// So that new configs don't require a new method parameter.
1112pub(crate) struct UpsertConfig {
1113    pub shrink_upsert_unused_buffers_by_ratio: usize,
1114}
1115
1116fn upsert_classic<'scope, T, FromTime, F, Fut, US>(
1117    input: VecCollection<'scope, T, (UpsertKey, Option<UpsertValue>, FromTime), Diff>,
1118    key_indices: Vec<usize>,
1119    resume_upper: Antichain<T>,
1120    previous_ok: VecCollection<'scope, T, Row, Diff>,
1121    previous_err: VecCollection<'scope, T, DataflowError, Diff>,
1122    previous_token: Option<Vec<PressOnDropButton>>,
1123    upsert_metrics: UpsertMetrics,
1124    source_config: crate::source::SourceExportCreationConfig,
1125    state: F,
1126    upsert_config: UpsertConfig,
1127    prevent_snapshot_buffering: bool,
1128    snapshot_buffering_max: Option<usize>,
1129) -> (
1130    VecCollection<'scope, T, Result<Row, DataflowError>, Diff>,
1131    StreamVec<'scope, T, (Option<GlobalId>, HealthStatusUpdate)>,
1132    StreamVec<'scope, T, Infallible>,
1133    PressOnDropButton,
1134)
1135where
1136    T: Timestamp + TotalOrder + Sync,
1137    F: FnOnce() -> Fut + 'static,
1138    Fut: std::future::Future<Output = US>,
1139    US: UpsertStateBackend<T, FromTime>,
1140    FromTime: timely::ExchangeData + Clone + Ord + Sync,
1141{
1142    let mut builder = AsyncOperatorBuilder::new("Upsert".to_string(), input.scope());
1143
1144    let previous = key_persist_feedback(previous_ok, previous_err, key_indices);
1145    let (output_handle, output) = builder.new_output();
1146
1147    // An output that just reports progress of the snapshot consolidation process upstream to the
1148    // persist source to ensure that backpressure is applied
1149    let (_snapshot_handle, snapshot_stream) =
1150        builder.new_output::<CapacityContainerBuilder<Vec<Infallible>>>();
1151
1152    let (mut health_output, health_stream) = builder.new_output();
1153    let mut input = builder.new_input_for(
1154        input.inner,
1155        Exchange::new(move |((key, _, _), _, _)| UpsertKey::hashed(key)),
1156        &output_handle,
1157    );
1158
1159    let mut previous = builder.new_input_for(
1160        previous.inner,
1161        Exchange::new(|((key, _), _, _)| UpsertKey::hashed(key)),
1162        &output_handle,
1163    );
1164
1165    let upsert_shared_metrics = Arc::clone(&upsert_metrics.shared);
1166    let shutdown_button = builder.build(move |caps| async move {
1167        let [mut output_cap, mut snapshot_cap, health_cap]: [_; 3] = caps.try_into().unwrap();
1168
1169        let mut state = UpsertState::<_, _, FromTime>::new(
1170            state().await,
1171            upsert_shared_metrics,
1172            &upsert_metrics,
1173            source_config.source_statistics.clone(),
1174            upsert_config.shrink_upsert_unused_buffers_by_ratio,
1175        );
1176        let mut events = vec![];
1177        let mut snapshot_upper = Antichain::from_elem(Timestamp::minimum());
1178
1179        let mut stash = vec![];
1180
1181        let mut error_emitter = (&mut health_output, &health_cap);
1182
1183        tracing::info!(
1184            ?resume_upper,
1185            ?snapshot_upper,
1186            "timely-{} upsert source {} starting rehydration",
1187            source_config.worker_id,
1188            source_config.id
1189        );
1190        // Read and consolidate the snapshot from the 'previous' input until it
1191        // reaches the `resume_upper`.
1192        while !PartialOrder::less_equal(&resume_upper, &snapshot_upper) {
1193            previous.ready().await;
1194            while let Some(event) = previous.next_sync() {
1195                match event {
1196                    AsyncEvent::Data(_cap, data) => {
1197                        events.extend(data.into_iter().filter_map(|((key, value), ts, diff)| {
1198                            if !resume_upper.less_equal(&ts) {
1199                                Some((key, value, diff))
1200                            } else {
1201                                None
1202                            }
1203                        }))
1204                    }
1205                    AsyncEvent::Progress(upper) => {
1206                        snapshot_upper = upper;
1207                    }
1208                };
1209            }
1210
1211            match state
1212                .consolidate_chunk(
1213                    events.drain(..),
1214                    PartialOrder::less_equal(&resume_upper, &snapshot_upper),
1215                )
1216                .await
1217            {
1218                Ok(_) => {
1219                    if let Some(ts) = snapshot_upper.clone().into_option() {
1220                        // As we shutdown, we could ostensibly get data from later than the
1221                        // `resume_upper`, which we ignore above. We don't want our output capability to make
1222                        // it further than the `resume_upper`.
1223                        if !resume_upper.less_equal(&ts) {
1224                            snapshot_cap.downgrade(&ts);
1225                            output_cap.downgrade(&ts);
1226                        }
1227                    }
1228                }
1229                Err(e) => {
1230                    UpsertErrorEmitter::<T>::emit(
1231                        &mut error_emitter,
1232                        "Failed to rehydrate state".to_string(),
1233                        e,
1234                    )
1235                    .await;
1236                }
1237            }
1238        }
1239
1240        drop(events);
1241        drop(previous_token);
1242        drop(snapshot_cap);
1243
1244        // Exchaust the previous input. It is expected to immediately reach the empty
1245        // antichain since we have dropped its token.
1246        //
1247        // Note that we do not need to also process the `input` during this, as the dropped token
1248        // will shutdown the `backpressure` operator
1249        while let Some(_event) = previous.next().await {}
1250
1251        // After snapshotting, our output frontier is exactly the `resume_upper`
1252        if let Some(ts) = resume_upper.as_option() {
1253            output_cap.downgrade(ts);
1254        }
1255
1256        tracing::info!(
1257            "timely-{} upsert source {} finished rehydration",
1258            source_config.worker_id,
1259            source_config.id
1260        );
1261
1262        // A re-usable buffer of changes, per key. This is an `IndexMap` because it has to be `drain`-able
1263        // and have a consistent iteration order.
1264        let mut commands_state: indexmap::IndexMap<_, types::UpsertValueAndSize<T, FromTime>> =
1265            indexmap::IndexMap::new();
1266        let mut multi_get_scratch = Vec::new();
1267
1268        // Now can can resume consuming the collection
1269        let mut output_updates = vec![];
1270        let mut input_upper = Antichain::from_elem(Timestamp::minimum());
1271
1272        while let Some(event) = input.next().await {
1273            // Buffer as many events as possible. This should be bounded, as new data can't be
1274            // produced in this worker until we yield to timely.
1275            let events = [event]
1276                .into_iter()
1277                .chain(std::iter::from_fn(|| input.next().now_or_never().flatten()))
1278                .enumerate();
1279
1280            let mut partial_drain_time = None;
1281            for (i, event) in events {
1282                match event {
1283                    AsyncEvent::Data(cap, mut data) => {
1284                        tracing::trace!(
1285                            time=?cap.time(),
1286                            updates=%data.len(),
1287                            "received data in upsert"
1288                        );
1289                        stage_input(
1290                            &mut stash,
1291                            &mut data,
1292                            &input_upper,
1293                            &resume_upper,
1294                            upsert_config.shrink_upsert_unused_buffers_by_ratio,
1295                        );
1296
1297                        let event_time = cap.time();
1298                        // If the data is at _exactly_ the output frontier, we can preemptively drain it into the state.
1299                        // Data within this set events strictly beyond this time are staged as
1300                        // normal.
1301                        //
1302                        // This is a load-bearing optimization, as it is required to avoid buffering
1303                        // the entire source snapshot in the `stash`.
1304                        if prevent_snapshot_buffering && output_cap.time() == event_time {
1305                            partial_drain_time = Some(event_time.clone());
1306                        }
1307                    }
1308                    AsyncEvent::Progress(upper) => {
1309                        tracing::trace!(?upper, "received progress in upsert");
1310                        // Ignore progress updates before the `resume_upper`, which is our initial
1311                        // capability post-snapshotting.
1312                        if PartialOrder::less_than(&upper, &resume_upper) {
1313                            continue;
1314                        }
1315
1316                        // Disable the partial drain as this progress event covers
1317                        // the `output_cap` time.
1318                        partial_drain_time = None;
1319                        drain_staged_input::<_, _, _, _>(
1320                            &mut stash,
1321                            &mut commands_state,
1322                            &mut output_updates,
1323                            &mut multi_get_scratch,
1324                            DrainStyle::ToUpper(&upper),
1325                            &mut error_emitter,
1326                            &mut state,
1327                            &source_config,
1328                        )
1329                        .await;
1330
1331                        output_handle.give_container(&output_cap, &mut output_updates);
1332
1333                        if let Some(ts) = upper.as_option() {
1334                            output_cap.downgrade(ts);
1335                        }
1336                        input_upper = upper;
1337                    }
1338                }
1339                let events_processed = i + 1;
1340                if let Some(max) = snapshot_buffering_max {
1341                    if events_processed >= max {
1342                        break;
1343                    }
1344                }
1345            }
1346
1347            // If there were staged events that occurred at the capability time, drain
1348            // them. This is safe because out-of-order updates to the same key that are
1349            // drained in separate calls to `drain_staged_input` are correctly ordered by
1350            // their `FromTime` in `drain_staged_input`.
1351            //
1352            // Note also that this may result in more updates in the output collection than
1353            // the minimum. However, because the frontier only advances on `Progress` updates,
1354            // the collection always accumulates correctly for all keys.
1355            if let Some(partial_drain_time) = partial_drain_time {
1356                drain_staged_input::<_, _, _, _>(
1357                    &mut stash,
1358                    &mut commands_state,
1359                    &mut output_updates,
1360                    &mut multi_get_scratch,
1361                    DrainStyle::AtTime(partial_drain_time),
1362                    &mut error_emitter,
1363                    &mut state,
1364                    &source_config,
1365                )
1366                .await;
1367
1368                output_handle.give_container(&output_cap, &mut output_updates);
1369            }
1370        }
1371    });
1372
1373    (
1374        output.as_collection().map(|result| match result {
1375            Ok(ok) => Ok(ok),
1376            Err(err) => Err(DataflowError::from(EnvelopeError::Upsert(*err))),
1377        }),
1378        health_stream,
1379        snapshot_stream,
1380        shutdown_button.press_on_drop(),
1381    )
1382}
1383
1384#[async_trait::async_trait(?Send)]
1385pub(crate) trait UpsertErrorEmitter<T> {
1386    async fn emit(&mut self, context: String, e: anyhow::Error);
1387}
1388
1389#[async_trait::async_trait(?Send)]
1390impl<T: Timestamp> UpsertErrorEmitter<T>
1391    for (
1392        &mut AsyncOutputHandle<
1393            T,
1394            CapacityContainerBuilder<Vec<(Option<GlobalId>, HealthStatusUpdate)>>,
1395        >,
1396        &Capability<T>,
1397    )
1398{
1399    async fn emit(&mut self, context: String, e: anyhow::Error) {
1400        process_upsert_state_error::<T>(context, e, self.0, self.1).await
1401    }
1402}
1403
1404/// Emit the given error, and stall till the dataflow is restarted.
1405async fn process_upsert_state_error<T: Timestamp>(
1406    context: String,
1407    e: anyhow::Error,
1408    health_output: &AsyncOutputHandle<
1409        T,
1410        CapacityContainerBuilder<Vec<(Option<GlobalId>, HealthStatusUpdate)>>,
1411    >,
1412    health_cap: &Capability<T>,
1413) {
1414    let update = HealthStatusUpdate::halting(e.context(context).to_string_with_causes(), None);
1415    health_output.give(health_cap, (None, update));
1416    std::future::pending::<()>().await;
1417    unreachable!("pending future never returns");
1418}