Skip to main content

mz_storage/
upsert_continual_feedback_v2.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 the feedback UPSERT operator.
11//!
12//! # Architecture
13//!
14//! The operator converts a stream of upsert commands `(key, Option<value>)` into
15//! a differential collection of `(key, value)` pairs, using a feedback loop
16//! through persist to maintain the "previous value" state needed for computing
17//! retractions.
18//!
19//! ## Dataflow topology
20//!
21//! ```text
22//!   Source input ──► ┌──────────┐ ──► Output ──► Persist
23//!                    │  Upsert  │
24//!   Persist read ──► └──────────┘
25//!       ▲                                           │
26//!       └───────────── feedback ────────────────────┘
27//! ```
28//!
29//! ## Operator loop (each iteration)
30//!
31//! 1. **Ingest source data.** Read upsert commands from the source input,
32//!    wrap each in an `UpsertDiff` (carrying a columnar order key projected
33//!    from `FromTime` via [`UpsertSourceTime`] for dedup), and push into the
34//!    source-stash batcher. The batcher consolidates entries for the same
35//!    `(key, time)` via the `UpsertDiff` Semigroup, keeping the update with
36//!    the highest order key (latest source offset), through amortized
37//!    geometric merging as data is pushed in, and cold stash state leaves RSS
38//!    through the flavor's spill path (see below). This bounds resident
39//!    memory even during large source snapshots.
40//!
41//! 2. **Read persist frontier.** Check the probe on the persist arrangement
42//!    to learn which times have been committed. When the persist frontier
43//!    reaches the resume upper, rehydration is complete.
44//!
45//! 3. **Seal & drain.** Call `batcher.seal(input_upper)` to extract all
46//!    source-finalized entries as sorted, consolidated chunks. Each entry is
47//!    classified:
48//!    - **Eligible** (at the persist frontier): the persist trace has the
49//!      correct "before" state for this time. Look up the old value in the
50//!      feedback arrangement, emit a retraction if present, and emit the new
51//!      value.
52//!    - **Ineligible** (between persist and input frontiers): persist hasn't
53//!      caught up yet. Push back into the batcher for the next iteration.
54//!    - **Already persisted** (below the persist frontier): some writer has
55//!      already advanced the shard past this time, so it is dropped. See the
56//!      drain functions for why re-stashing it would strand the data and pin
57//!      the output frontier below the shard upper.
58//!
59//! 4. **Capability management.** Downgrade the output capability to the
60//!    minimum time of any remaining buffered data (in the batcher or pushed
61//!    back as ineligible). Drop the capability entirely when the batcher is
62//!    empty.
63//!
64//! ## Stash flavors
65//!
66//! [`UpsertStashFlavor`], resolved from `enable_upsert_chunked_stash` at
67//! operator construction, selects between two instantiations of the same
68//! loop:
69//!
70//! * **Chunked**: the stash is differential's chunk merge batcher over
71//!   `ColumnChunk`s and the feedback arrangement is a spine of chunk batches.
72//!   Committed chunk bodies spill to the process buffer pool, the drain
73//!   loads sealed chunks one at a time, and prior state comes back through
74//!   bulk probes of the trace's batches.
75//! * **Paged**: the stash is the paged columnar merge batcher and the
76//!   feedback arrangement is a `ValRowSpine`. Cold chains page out of RSS
77//!   through the storage-owned column pager, and prior state comes back
78//!   through a trace cursor.
79//!
80//! Both flavors' spill paths are gated by `enable_upsert_paged_spill`.
81//!
82//! ## Eligibility condition (total order)
83//!
84//! For a total-order timestamp with `input_upper = {i}` and
85//! `persist_upper = {p}`, an entry at time `ts` is eligible when
86//! `ts == p < i` — the source has finalized it and persist is exactly at
87//! that time, so the feedback trace holds the correct prior state. An entry
88//! with `p < ts` is ineligible (persist hasn't caught up), and one with
89//! `ts < p` is already persisted and dropped.
90
91use std::fmt::Debug;
92
93use differential_dataflow::difference::{IsZero, Semigroup};
94use differential_dataflow::hashable::Hashable;
95use differential_dataflow::lattice::Lattice;
96use differential_dataflow::logging::Logger;
97use differential_dataflow::operators::arrange::agent::TraceAgent;
98use differential_dataflow::operators::arrange::arrangement::{Arranged, arrange_core};
99use differential_dataflow::trace::chunk::{ChunkBatch, ChunkBatcher, ChunkBuilder};
100use differential_dataflow::trace::{BatchReader, Batcher, Cursor, Description, TraceReader};
101use differential_dataflow::{AsCollection, VecCollection};
102use mz_dyncfg::ConfigSet;
103use mz_repr::{Datum, Diff, GlobalId, Row};
104// Only the fuzzing-gated `datum_seq_to_upsert_value` takes a `DatumSeq`.
105#[cfg(feature = "fuzzing")]
106use mz_row_spine::DatumSeq;
107use mz_row_spine::{FundedValRowSpine, ValRowColPagedBuilder};
108use mz_storage_types::dyncfgs::ENABLE_UPSERT_CHUNKED_STASH;
109use mz_storage_types::errors::{DataflowError, EnvelopeError, UpsertError};
110use mz_timely_util::builder_async::{
111    AsyncOutputHandle, Event as AsyncEvent, OperatorBuilder as AsyncOperatorBuilder,
112    PressOnDropButton,
113};
114use mz_timely_util::columnar::batcher::ColumnChunker;
115use mz_timely_util::columnar::body::ColumnBody;
116use mz_timely_util::columnar::builder::ColumnBuilder;
117use mz_timely_util::columnar::chunk::{ChunkChunker, ColumnChunk};
118use mz_timely_util::columnar::merge_batcher::{ColumnMergeBatcher, PagedChunker};
119use mz_timely_util::columnar::unload::UnloadBatch;
120use mz_timely_util::columnar::{Col2ValPagedBatcher, Column};
121use mz_timely_util::containers::stack::FueledBuilder;
122use std::convert::Infallible;
123use timely::container::{CapacityContainerBuilder, PushInto};
124use timely::dataflow::channels::pact::{Exchange, Pipeline};
125use timely::dataflow::operators::generic::Operator;
126use timely::dataflow::operators::{Capability, CapabilitySet, Exchange as _};
127use timely::dataflow::{Stream, StreamVec};
128use timely::order::{PartialOrder, TotalOrder};
129use timely::progress::frontier::AntichainRef;
130use timely::progress::timestamp::Refines;
131use timely::progress::{Antichain, Timestamp};
132
133use crate::healthcheck::HealthStatusUpdate;
134use crate::metrics::upsert::UpsertMetrics;
135use crate::statistics::SourceStatistics;
136use crate::upsert::UpsertKey;
137use crate::upsert::UpsertSourceTime;
138use crate::upsert::UpsertValue;
139
140/// Which stash and feedback-arrangement representation the upsert-v2
141/// operator instantiates. The two flavors run the same operator loop; they
142/// differ in the batcher, the feedback trace, and how the drain reads prior
143/// state. See the module docs for the comparison.
144#[derive(Clone, Copy, Debug)]
145pub enum UpsertStashFlavor {
146    /// Paged columnar merge batcher stash, `ValRowSpine` feedback
147    /// arrangement, cursor-based drain. Spills through the storage-owned
148    /// column pager.
149    Paged,
150    /// Chunk merge batcher stash, chunk-spine feedback arrangement,
151    /// bulk-probe drain. Spills through the process buffer pool.
152    Chunked,
153}
154
155impl UpsertStashFlavor {
156    /// Resolve the flavor from the replica's config set. Call this once per
157    /// source at operator construction time, so a dataflow keeps one flavor
158    /// for its whole life even if the flag flips underneath it.
159    pub fn from_config(config: &ConfigSet) -> Self {
160        if ENABLE_UPSERT_CHUNKED_STASH.get(config) {
161            Self::Chunked
162        } else {
163            Self::Paged
164        }
165    }
166}
167
168/// The paged flavor's persist-feedback batcher, wrapping
169/// [`Col2ValPagedBatcher`] only to capture the storage upsert-stash pager at
170/// construction.
171///
172/// `arrange_core` builds its batcher via [`Batcher::new`], which has no pager
173/// hook, so a plain `Col2ValPagedBatcher` falls back to the process-global
174/// (compute) pager, meaning the feedback arrangement's spill would be gated
175/// by compute's `enable_column_paged_batcher_spill` rather than storage's
176/// `enable_upsert_paged_spill`. Injecting `upsert_stash_pager::pager()` in
177/// `new` puts the feedback arrangement under the same flag as the source
178/// stash. Every other method delegates to the inner batcher unchanged.
179struct UpsertFeedbackBatcher<T: columnar::Columnar>(Col2ValPagedBatcher<UpsertKey, Row, T, Diff>);
180
181impl<T> Batcher for UpsertFeedbackBatcher<T>
182where
183    T: Timestamp + columnar::Columnar + Default + PartialOrder,
184    for<'a> columnar::Ref<'a, T>: Copy + Ord,
185{
186    type Output = Column<((UpsertKey, Row), T, Diff)>;
187    type Time = T;
188
189    fn new(logger: Option<Logger>, operator_id: usize) -> Self {
190        let mut batcher =
191            <Col2ValPagedBatcher<UpsertKey, Row, T, Diff> as Batcher>::new(logger, operator_id);
192        batcher.set_pager(crate::upsert::upsert_stash_pager::pager());
193        Self(batcher)
194    }
195
196    fn seal(&mut self, upper: Antichain<T>) -> (Vec<Self::Output>, Description<T>) {
197        self.0.seal(upper)
198    }
199
200    fn frontier(&mut self) -> AntichainRef<'_, T> {
201        self.0.frontier()
202    }
203}
204
205impl<T> PushInto<Column<((UpsertKey, Row), T, Diff)>> for UpsertFeedbackBatcher<T>
206where
207    T: Timestamp + columnar::Columnar + Default + PartialOrder,
208    for<'a> columnar::Ref<'a, T>: Copy + Ord,
209{
210    fn push_into(&mut self, chunk: Column<((UpsertKey, Row), T, Diff)>) {
211        self.0.push_into(chunk)
212    }
213}
214
215/// One persist-feedback update: a key, its current value row, the time, and
216/// an additive count.
217type FeedbackUpdate<T> = ((UpsertKey, Row), T, Diff);
218
219/// One feedback-arrangement chunk: a sorted, consolidated run of updates,
220/// resident or spilled to the buffer pool. Both batcher chains and sealed
221/// spine batches are sequences of these, so the arrangement's state pages out
222/// of RSS under the pool's budget, and the drain reads it back through the
223/// bulk [`UnloadChunk`](mz_timely_util::columnar::unload::UnloadChunk)
224/// surface: copy-out probes, no cursor borrows.
225type FeedbackChunk<T> = ColumnChunk<(UpsertKey, Row), T, Diff>;
226
227/// The feedback arrangement's trace: a funded spine of `Rc`-shared chunk
228/// batches. See [`mz_timely_util::funded_spine`].
229type FeedbackSpine<T> =
230    mz_timely_util::funded_spine::Spine<std::rc::Rc<ChunkBatch<FeedbackChunk<T>>>>;
231
232// The source stash carries the upsert payload in a custom diff type so the
233// merge batcher consolidates by (key, time), keeping the update with the
234// highest `FromTime` (latest source offset) per group. The diff is `Columnar`
235// so the paged merge batcher can store it in a `Column` and page it out of RSS.
236//
237// The value is a tag-encoded `Row` (see `upsert_value_to_row`) rather than an
238// `UpsertValue`: folding both the `Ok` and `Err` arms into one `Row` lets the
239// value share a single columnar byte container, and `Row` already implements
240// `Columnar`. `None` is a deletion tombstone.
241
242// Derive ordering on the generated `UpsertDiffReference` too: the paged merge
243// batcher requires `Ref: Ord` to sort the `(key, time, diff)` columns it
244// consolidates. The derived order (by `from_time`, then `value`) is fine —
245// "max FromTime wins" can tie only between equal `from_time`s, and a source
246// never emits two distinct values for the same `(key, time, from_time)`, so
247// the consolidated result doesn't depend on the fold order of equal
248// `(key, time)` runs.
249#[derive(Clone, Debug, Default, columnar::Columnar)]
250#[columnar(derive(PartialEq, Eq, PartialOrd, Ord))]
251struct UpsertDiff<O> {
252    from_time: O,
253    value: Option<Row>,
254}
255
256impl<O> IsZero for UpsertDiff<O> {
257    fn is_zero(&self) -> bool {
258        false
259    }
260}
261
262impl<O: Ord + Clone> Semigroup for UpsertDiff<O> {
263    fn plus_equals(&mut self, rhs: &Self) {
264        if rhs.from_time > self.from_time {
265            *self = rhs.clone();
266        }
267    }
268}
269
270// Accumulate a borrowed columnar reference: the paged merge batcher consolidates
271// `Column`-resident diffs through this path on every fold of an equal
272// `(key, time)` run. Materialize only the order key to decide the "max FromTime
273// wins" comparison — copying the value `Row` out of the column solely when `rhs`
274// wins. Losing folds (the common case for a repeatedly-updated key) then pay no
275// `Row` copy at all.
276impl<'a, O> Semigroup<columnar::Ref<'a, UpsertDiff<O>>> for UpsertDiff<O>
277where
278    O: columnar::Columnar + Ord + Clone,
279{
280    fn plus_equals(&mut self, rhs: &columnar::Ref<'a, UpsertDiff<O>>) {
281        let rhs_from_time = <O as columnar::Columnar>::into_owned(rhs.from_time);
282        if rhs_from_time > self.from_time {
283            self.from_time = rhs_from_time;
284            self.value = <Option<Row> as columnar::Columnar>::into_owned(rhs.value);
285        }
286    }
287}
288
289/// One source-stash update: a key, its dataflow time, and the payload diff.
290/// `O` is the columnar order key projected from the source `FromTime` (see
291/// [`UpsertSourceTime`]).
292type UpsertUpdate<T, O> = (UpsertKey, T, UpsertDiff<O>);
293
294/// One stash chunk: a sorted, consolidated run of updates, resident or
295/// spilled to the buffer pool.
296type UpsertChunk<T, O> = ColumnChunk<UpsertKey, T, UpsertDiff<O>>;
297
298/// The chunked flavor's stash: differential's chunk merge batcher over
299/// `ColumnChunk`s. Data is pushed in unsorted. The batcher maintains
300/// geometrically-sized sorted chains and consolidates via the UpsertDiff
301/// Semigroup automatically. Committed chunks spill their bodies to the
302/// process buffer pool (see `mz_timely_util::columnar::chunk`), so the
303/// not-yet-eligible backlog (the snapshot / persist-lag window) pages out of
304/// RSS instead of growing it.
305type UpsertChunkBatcher<T, O> = ChunkBatcher<UpsertChunk<T, O>>;
306
307/// The paged flavor's stash: the paged columnar merge batcher, consolidating
308/// like [`UpsertChunkBatcher`] but storing each chain entry as a `Column`
309/// routed through the storage-owned pager, which pages cold chains out of
310/// RSS.
311type UpsertPagedBatcher<T, O> = ColumnMergeBatcher<UpsertKey, T, UpsertDiff<O>>;
312
313/// The chunker that sorts and consolidates raw input into the `Column` chunks
314/// both stash batchers consume.
315type UpsertChunker<T, O> = ColumnChunker<UpsertUpdate<T, O>>;
316
317/// The operator's data-output handle. A fueled `Vec` builder so the drain can
318/// `give_fueled` each emitted update and yield to timely under large snapshot
319/// drains instead of monopolizing the worker.
320type UpsertOutputHandle<T> =
321    AsyncOutputHandle<T, FueledBuilder<CapacityContainerBuilder<Vec<(UpsertValue, T, Diff)>>>>;
322
323// The persist-feedback arrangement stores `(UpsertKey, Row)` pairs in columnar
324// chunk batches ([`FeedbackSpine`]). `UpsertValue` is
325// `Result<Row, Box<UpsertError>>`, so we fold both arms into a single `Row`
326// with a leading tag column so they share the value column.
327
328/// Encode an [`UpsertValue`] as a `Row` with a leading tag column so both `Ok`
329/// and `Err` payloads round-trip through `Row` byte storage.
330///
331/// Used on the render path. `pub` only so [`crate::fuzz_exports`] can re-export
332/// it under the `fuzzing` feature for the storage fuzz crate. The enclosing
333/// module is crate-private, so it is not otherwise reachable. Not a stable
334/// public API.
335#[doc(hidden)]
336pub fn upsert_value_to_row(value: &UpsertValue) -> Row {
337    let mut row = Row::default();
338    let mut packer = row.packer();
339    match value {
340        Ok(ok) => {
341            packer.push(Datum::UInt8(0));
342            packer.extend(ok.iter());
343        }
344        Err(err) => {
345            packer.push(Datum::UInt8(1));
346            let bytes =
347                bincode::serialize(err.as_ref()).expect("UpsertError is serializable via bincode");
348            packer.push(Datum::Bytes(&bytes));
349        }
350    }
351    row
352}
353
354/// Heap-size estimate for an emitted [`UpsertValue`], used to drive
355/// `give_fueled` yielding on the output edge.
356fn upsert_value_byte_len(value: &UpsertValue) -> usize {
357    match value {
358        Ok(row) => row.byte_len(),
359        Err(err) => std::mem::size_of_val(err.as_ref()),
360    }
361}
362
363/// Decode an [`UpsertValue`] produced by [`upsert_value_to_row`] back from the
364/// `DatumSeq` view returned by a `ValRowSpine` cursor.
365///
366/// Exists only for the storage fuzz crate, so it is gated behind the `fuzzing`
367/// feature. Not a stable public API.
368#[cfg(feature = "fuzzing")]
369#[doc(hidden)]
370pub fn datum_seq_to_upsert_value(seq: DatumSeq<'_>) -> UpsertValue {
371    decode_upsert_value(seq)
372}
373
374/// Decode an [`UpsertValue`] produced by [`upsert_value_to_row`] from any
375/// datum iterator, whether over a columnar `Row` reference or an owned `Row`.
376fn decode_upsert_value<'a>(mut iter: impl Iterator<Item = Datum<'a>>) -> UpsertValue {
377    let tag = match iter.next() {
378        Some(Datum::UInt8(tag)) => tag,
379        other => panic!("upsert value missing UInt8 tag, got {:?}", other),
380    };
381    match tag {
382        0 => {
383            let mut row = Row::default();
384            row.packer().extend(iter);
385            Ok(row)
386        }
387        1 => {
388            let bytes = match iter.next() {
389                Some(Datum::Bytes(b)) => b,
390                other => panic!("upsert error tag missing Bytes payload, got {:?}", other),
391            };
392            let err: UpsertError =
393                bincode::deserialize(bytes).expect("UpsertError bincode round-trip");
394            Err(Box::new(err))
395        }
396        tag => panic!("unknown upsert value tag {tag}"),
397    }
398}
399
400/// Transforms a stream of upserts (key-value updates) into a differential
401/// collection.
402///
403/// Persist feedback is arranged into a differential trace (DD manages the
404/// spine lifecycle). Source input is stashed with a custom `UpsertDiff`
405/// Semigroup that deduplicates by keeping the highest FromTime per (key, time).
406///
407/// Has two inputs:
408///   1. **Source input** — upsert commands from the external source.
409///   2. **Persist input** — feedback of the operator's own output, read back
410///      from persist.  Arranged into a trace for prior-value lookups.
411///
412/// `flavor` selects the stash and feedback-arrangement representation; see
413/// the module docs.
414#[allow(clippy::disallowed_methods)]
415pub fn upsert_inner<'scope, T, FromTime>(
416    flavor: UpsertStashFlavor,
417    input: VecCollection<'scope, T, (UpsertKey, Option<UpsertValue>, FromTime), Diff>,
418    key_indices: Vec<usize>,
419    resume_upper: Antichain<T>,
420    persist_ok: VecCollection<'scope, T, Row, Diff>,
421    persist_err: VecCollection<'scope, T, DataflowError, Diff>,
422    persist_token: Option<Vec<PressOnDropButton>>,
423    upsert_metrics: UpsertMetrics,
424    source_config: crate::source::SourceExportCreationConfig,
425) -> (
426    VecCollection<'scope, T, Result<Row, DataflowError>, Diff>,
427    StreamVec<'scope, T, (Option<GlobalId>, HealthStatusUpdate)>,
428    StreamVec<'scope, T, Infallible>,
429    PressOnDropButton,
430)
431where
432    T: Timestamp + TotalOrder + Sync,
433    T: Refines<mz_repr::Timestamp> + differential_dataflow::lattice::Lattice,
434    T: columnation::Columnation,
435    T: columnar::Columnar + Default,
436    for<'a> columnar::Ref<'a, T>: Copy + Ord,
437    FromTime: Debug + timely::ExchangeData + Clone + Ord + Sync,
438    FromTime: UpsertSourceTime,
439{
440    // The feedback keying and encoding are flavor-independent; the
441    // arrangement they feed is not, so each arm names its own arrange types
442    // and hands the arrangement to the shared operator loop.
443    let encoded = encode_feedback(
444        persist_ok,
445        persist_err,
446        key_indices,
447        source_config.source_statistics.clone(),
448    );
449    match flavor {
450        UpsertStashFlavor::Chunked => {
451            // Chains and sealed batches alike are `FeedbackChunk`s whose
452            // bodies spill to the buffer pool, behind the same process spill
453            // gate as the source stash.
454            let persist_arranged = arrange_core::<
455                _,
456                _,
457                ChunkChunker<(UpsertKey, Row), T, Diff>,
458                ChunkBatcher<FeedbackChunk<T>>,
459                ChunkBuilder<FeedbackChunk<T>>,
460                FeedbackSpine<T>,
461            >(encoded, Pipeline, "Persist feedback");
462            build_upsert_operator::<ChunkedArm, _, _>(
463                input,
464                resume_upper,
465                persist_arranged,
466                persist_token,
467                upsert_metrics,
468                source_config,
469            )
470        }
471        UpsertStashFlavor::Paged => {
472            // The batcher routes its spine input through the storage-owned
473            // pager, paging cold feedback chains out of RSS, while
474            // `ValRowSpine` keeps keys in a columnation arena (`UpsertKey`
475            // is fixed-size `[u8; 32]`) and values as packed `Row` bytes in
476            // a `DatumContainer`.
477            let persist_arranged = arrange_core::<
478                _,
479                _,
480                PagedChunker<((UpsertKey, Row), T, Diff)>,
481                UpsertFeedbackBatcher<T>,
482                ValRowColPagedBuilder<UpsertKey, T, Diff>,
483                FundedValRowSpine<UpsertKey, T, Diff>,
484            >(encoded, Pipeline, "Persist feedback");
485            build_upsert_operator::<PagedArm, _, _>(
486                input,
487                resume_upper,
488                persist_arranged,
489                persist_token,
490                upsert_metrics,
491                source_config,
492            )
493        }
494    }
495}
496
497/// Key the persist feedback by [`UpsertKey`], record source statistics, and
498/// encode `(UpsertKey, UpsertValue)` as `(UpsertKey, Row)` `Column` chunks,
499/// the input both flavors' feedback arrangements consume. Built with
500/// `Pipeline` downstream of an `UpsertKey::hashed` exchange, so the
501/// arrangement keeps that locality.
502fn encode_feedback<'scope, T>(
503    persist_ok: VecCollection<'scope, T, Row, Diff>,
504    persist_err: VecCollection<'scope, T, DataflowError, Diff>,
505    key_indices: Vec<usize>,
506    source_statistics: SourceStatistics,
507) -> Stream<'scope, T, Column<((UpsertKey, Row), T, Diff)>>
508where
509    T: Timestamp + TotalOrder + Sync,
510    T: Refines<mz_repr::Timestamp> + differential_dataflow::lattice::Lattice,
511    T: columnation::Columnation,
512    T: columnar::Columnar + Default,
513    for<'a> columnar::Ref<'a, T>: Copy + Ord,
514{
515    let persist_keyed = crate::upsert::key_persist_feedback(persist_ok, persist_err, key_indices);
516    let persist_keyed = persist_keyed
517        .inner
518        // The arrangement already implicitly exchanges by key, so this is redundant, but we want to
519        // do it earlier so that we can inspect the stream properly for source statistics.
520        .exchange(move |((key, _), _, _)| UpsertKey::hashed(key))
521        .as_collection()
522        .inspect(move |((_, row), _, diff)| {
523            source_statistics.update_records_indexed_by(diff.into_inner());
524            source_statistics.update_bytes_indexed_by(
525                row.as_ref().map_or(0, |r| r.byte_len().try_into().unwrap()) * diff.into_inner(),
526            );
527        });
528    persist_keyed
529        .inner
530        .unary::<ColumnBuilder<((UpsertKey, Row), T, Diff)>, _, _, _>(
531            Pipeline,
532            "Persist feedback encode",
533            |_, _| {
534                move |input, output| {
535                    input.for_each(|time, data| {
536                        let mut session = output.session_with_builder(&time);
537                        for ((key, value), ts, diff) in data.drain(..) {
538                            let row = upsert_value_to_row(&value);
539                            session.give(((&key, &row), &ts, &diff));
540                        }
541                    });
542                }
543            },
544        )
545}
546
547/// The flavor-independent upsert-v2 operator loop, generic over an
548/// [`UpsertStashArm`].
549///
550/// Consumes the source input and the arranged persist feedback (constructed
551/// per flavor by [`upsert_inner`]) and drives the ingest / seal / drain /
552/// capability loop described in the module docs, calling through the arm at
553/// the few points where the flavors diverge.
554#[allow(clippy::disallowed_methods)]
555fn build_upsert_operator<'scope, A, T, FromTime>(
556    input: VecCollection<'scope, T, (UpsertKey, Option<UpsertValue>, FromTime), Diff>,
557    resume_upper: Antichain<T>,
558    persist_arranged: Arranged<'scope, TraceAgent<A::Spine>>,
559    persist_token: Option<Vec<PressOnDropButton>>,
560    upsert_metrics: UpsertMetrics,
561    source_config: crate::source::SourceExportCreationConfig,
562) -> (
563    VecCollection<'scope, T, Result<Row, DataflowError>, Diff>,
564    StreamVec<'scope, T, (Option<GlobalId>, HealthStatusUpdate)>,
565    StreamVec<'scope, T, Infallible>,
566    PressOnDropButton,
567)
568where
569    A: UpsertStashArm<T, FromTime::Order>,
570    T: Timestamp + TotalOrder + Sync,
571    T: Refines<mz_repr::Timestamp> + differential_dataflow::lattice::Lattice,
572    T: columnation::Columnation,
573    T: columnar::Columnar + Default,
574    for<'a> columnar::Ref<'a, T>: Copy + Ord,
575    FromTime: Debug + timely::ExchangeData + Clone + Ord + Sync,
576    FromTime: UpsertSourceTime,
577{
578    let mut persist_trace = persist_arranged.trace.clone();
579
580    // Probe the persist arrangement's stream for frontier tracking.
581    // This replaces receiving the batch stream as an input — we just
582    // read the probe frontier to know when persist has caught up.
583    use timely::dataflow::operators::Probe;
584    let (persist_probe, _persist_probe_stream) = persist_arranged.stream.probe();
585
586    // Build the async processing operator.
587    let mut builder = AsyncOperatorBuilder::new("Upsert V2".to_string(), input.scope());
588
589    let (output_handle, output) = builder
590        .new_output::<FueledBuilder<CapacityContainerBuilder<Vec<(UpsertValue, T, Diff)>>>>();
591    let (_snapshot_handle, snapshot_stream) =
592        builder.new_output::<CapacityContainerBuilder<Vec<Infallible>>>();
593    let (_health_output, health_stream) = builder
594        .new_output::<CapacityContainerBuilder<Vec<(Option<GlobalId>, HealthStatusUpdate)>>>();
595
596    let mut input = builder.new_input_for(
597        input.inner,
598        Exchange::new(move |((key, _, _), _, _)| UpsertKey::hashed(key)),
599        &output_handle,
600    );
601
602    // We still need the persist stream as an input so the operator wakes
603    // when the persist arrangement produces batches (frontier advances).
604    // We read the actual frontier from the probe though.
605    let mut persist_wakeup = builder.new_disconnected_input(_persist_probe_stream, Pipeline);
606
607    let shutdown_button = builder.build(move |caps| async move {
608        // Hold the persist source tokens for the operator's lifetime so the
609        // feedback shard stays open until shutdown.
610        let _persist_token = persist_token;
611
612        let [output_cap, snapshot_cap, _health_cap]: [_; 3] = caps.try_into().unwrap();
613        drop(output_cap);
614        let mut snapshot_cap = CapabilitySet::from_elem(snapshot_cap);
615
616        let mut hydrating = true;
617
618        // Source stash. The batcher maintains geometrically-sized sorted
619        // chains and consolidates via the UpsertDiff Semigroup as data is
620        // pushed in, bounding memory to O(unique key-time pairs) even during
621        // large initial snapshots. How cold stash state leaves RSS is
622        // flavor-specific; see the arm impls.
623        let mut batcher = A::new_batcher();
624        // The chunker sorts and consolidates raw input into the `Column` chunks
625        // the batcher consumes.
626        let mut chunker: UpsertChunker<T, FromTime::Order> = Default::default();
627        // Scratch buffer for accumulating source events before flushing to
628        // the batcher. Drained on each iteration via the chunker.
629        let mut push_buffer: Vec<UpsertUpdate<T, FromTime::Order>> = Vec::new();
630
631        // Capability held at the minimum time of any buffered data. When
632        // Some, the operator may still produce output; when None, the
633        // batcher is empty.
634        let mut stash_cap: Option<Capability<T>> = None;
635        let mut input_upper = Antichain::from_elem(Timestamp::minimum());
636
637        let snapshot_start = std::time::Instant::now();
638        let mut prev_persist_upper = Antichain::from_elem(Timestamp::minimum());
639
640        // Accumulators for rehydration metrics, set as gauges when rehydration completes.
641        let mut rehydration_total: u64 = 0;
642        let mut rehydration_updates: u64 = 0;
643
644        // Main operator loop. Each iteration performs four steps:
645        //   Step 1: Ingest source data into the batcher.
646        //   Step 2: Read the persist frontier and update rehydration state.
647        //   Step 3: Seal the batcher, drain eligible entries, push back the rest.
648        //   Step 4: Manage the output capability.
649        loop {
650            // Block until woken by source input or a persist frontier advance.
651            tokio::select! {
652                _ = input.ready() => {}
653                _ = persist_wakeup.ready() => {
654                    while let Some(event) = persist_wakeup.next_sync() {
655                        if let AsyncEvent::Data(_, batches) = event {
656                            for batch in batches {
657                                mz_timely_util::columnar::chunk::metrics::record_batch(batch.len());
658                            }
659                        }
660                    }
661                }
662            }
663
664            // Step 1: Ingest source data.
665            // Read all available source events, wrap each value in an
666            // UpsertDiff (carrying FromTime for dedup), and buffer them.
667            // Events before the resume_upper are dropped (already persisted).
668            while let Some(event) = input.next_sync() {
669                match event {
670                    AsyncEvent::Data(cap, data) => {
671                        let mut pushed_any = false;
672                        for ((key, value, from_time), ts, diff) in data {
673                            assert!(diff.is_positive(), "invalid upsert input");
674                            if PartialOrder::less_equal(&input_upper, &resume_upper)
675                                && !resume_upper.less_equal(&ts)
676                            {
677                                continue;
678                            }
679                            let value = value.as_ref().map(upsert_value_to_row);
680                            let from_time = from_time.upsert_order();
681                            push_buffer.push((key, ts, UpsertDiff { from_time, value }));
682                            pushed_any = true;
683                        }
684                        // Track the minimum capability across all buffered data
685                        // so we can emit output at the correct times.
686                        if pushed_any {
687                            stash_cap = Some(match stash_cap {
688                                Some(prev) if cap.time() < prev.time() => cap,
689                                Some(prev) => prev,
690                                None => cap,
691                            });
692                        }
693                    }
694                    AsyncEvent::Progress(upper) => {
695                        if PartialOrder::less_than(&upper, &resume_upper) {
696                            continue;
697                        }
698                        input_upper = upper;
699                    }
700                }
701            }
702
703            // Flush buffered events through the chunker into the batcher. This
704            // triggers the chunker + geometric chain merging, which consolidates
705            // entries for the same (key, time) via the UpsertDiff Semigroup.
706            A::flush(&mut push_buffer, &mut chunker, &mut batcher);
707
708            // Step 2: Read persist frontier.
709            // The persist probe tells us which output times have been
710            // committed back through the feedback loop. This determines:
711            //   - Whether rehydration is complete (persist >= resume_upper).
712            //   - Which source entries are eligible for processing (their
713            //     time must equal persist_upper so the feedback trace holds
714            //     the correct prior state).
715            //   - How far to compact the persist trace.
716            let persist_upper = persist_probe.with_frontier(|f| f.to_owned());
717
718            if persist_upper != prev_persist_upper {
719                let last_rehydration_chunk =
720                    hydrating && PartialOrder::less_equal(&resume_upper, &persist_upper);
721
722                if last_rehydration_chunk {
723                    hydrating = false;
724                    upsert_metrics
725                        .rehydration_latency
726                        .set(snapshot_start.elapsed().as_secs_f64());
727                    upsert_metrics.rehydration_total.set(rehydration_total);
728                    upsert_metrics.rehydration_updates.set(rehydration_updates);
729                    tracing::info!(
730                        worker_id = %source_config.worker_id,
731                        source_id = %source_config.id,
732                        "upsert finished rehydration",
733                    );
734                    snapshot_cap.downgrade(&[]);
735                }
736
737                let _ = snapshot_cap.try_downgrade(persist_upper.iter());
738
739                // Compact the trace so the spine can merge old batches.
740                persist_trace.set_logical_compaction(persist_upper.borrow());
741                persist_trace.set_physical_compaction(persist_upper.borrow());
742
743                prev_persist_upper = persist_upper.clone();
744            }
745
746            // Step 3: Seal & drain.
747            // Seal the batcher at input_upper to extract all source-finalized
748            // entries as sorted, consolidated chunks. The seal merges all
749            // internal chains (O(N) linear merge of sorted data) and splits
750            // by time: entries at ts < input_upper are extracted, the rest
751            // stay in the batcher.
752            //
753            // Extracted entries are partitioned into:
754            //   - Eligible (ts == persist_upper): processed now via a
755            //     prior-value lookup in the persist trace.
756            //   - Ineligible (persist_upper < ts < input_upper): persist
757            //     hasn't caught up yet, so pushed back into the batcher.
758            //
759            // We skip the seal entirely unless an eligible entry is at all
760            // possible. `seal` performs an O(N) merge of all chains
761            // regardless of how much it extracts, so calling it when nothing
762            // can be processed makes the operator quadratic in the number of
763            // wakeups (a real pathology during upstream snapshots and during
764            // rehydration when the source races ahead of persist).
765            //
766            // For an entry at `ts` to be eligible we need
767            // `ts == persist_upper && ts < input_upper`. The necessary
768            // preconditions, expressible without scanning the batcher:
769            //   1. `cap.time() <= persist_upper`. Since `cap.time()` is
770            //      maintained as a lower bound on `min(ts in batcher)`, if
771            //      `cap.time() > persist_upper` then every buffered ts is
772            //      strictly above persist_upper and none can equal it.
773            //   2. `persist_upper < input_upper`. Otherwise no `ts` that
774            //      satisfies `ts == persist_upper` can also satisfy
775            //      `ts < input_upper`.
776            //
777            // This naturally covers both the post-hydration source-snapshot
778            // case (cap == persist == input → condition 2 fails) and the
779            // rehydration-with-source-ahead case (cap > persist → condition
780            // 1 fails). It also no-ops correctly when persist has shut down
781            // (empty persist_upper makes condition 2 vacuously false).
782            if let Some(cap) = stash_cap.as_mut()
783                && !persist_upper.less_than(cap.time())
784                && PartialOrder::less_than(&persist_upper, &input_upper)
785            {
786                // Step 1 already consolidated `push_buffer` through the chunker
787                // (which readies a complete chunk per `push_into`), so the
788                // chunker holds nothing pending here and we can seal directly.
789                let (sealed, _description) = batcher.seal(input_upper.clone());
790                // Frontier of data remaining in the batcher (ts >= input_upper).
791                let remaining_frontier = batcher.frontier().to_owned();
792
793                let mut ineligible = Vec::new();
794                // The drain emits eligible output directly through
795                // `output_handle` (fueled), so there is no intermediate
796                // output buffer to drain afterward.
797                let drain_stats = A::drain(
798                    sealed,
799                    &mut ineligible,
800                    &output_handle,
801                    &*cap,
802                    &persist_upper,
803                    &mut persist_trace,
804                    source_config.worker_id,
805                    source_config.id,
806                )
807                .await;
808
809                upsert_metrics.multi_get_size.inc_by(drain_stats.eligible);
810                upsert_metrics
811                    .multi_get_result_count
812                    .inc_by(drain_stats.result_count);
813                upsert_metrics
814                    .multi_put_size
815                    .inc_by(drain_stats.output_count);
816                upsert_metrics.upsert_inserts.inc_by(drain_stats.inserts);
817                upsert_metrics.upsert_updates.inc_by(drain_stats.updates);
818                upsert_metrics.upsert_deletes.inc_by(drain_stats.deletes);
819
820                if hydrating {
821                    rehydration_total += drain_stats.inserts;
822                    rehydration_updates += drain_stats.eligible;
823                }
824
825                // Step 4: Capability management.
826                // Downgrade the output capability to the minimum time of any
827                // remaining data: either entries still in the batcher (above
828                // input_upper) or ineligible entries being pushed back.
829                let min_ineligible_ts = ineligible.iter().map(|(_, ts, _)| ts).min().cloned();
830                A::flush(&mut ineligible, &mut chunker, &mut batcher);
831
832                // `Option::min` alone would be wrong here, `None` sorts low.
833                // Chain the candidates and take the min over present ones.
834                let min_ts = remaining_frontier
835                    .elements()
836                    .first()
837                    .into_iter()
838                    .chain(min_ineligible_ts.as_ref())
839                    .min();
840                match min_ts {
841                    Some(min_ts) => cap.downgrade(min_ts),
842                    // Batcher is completely empty. Drop the capability so
843                    // downstream operators can make progress.
844                    None => stash_cap = None,
845                }
846            }
847
848            if input_upper.is_empty() {
849                break;
850            }
851        }
852    });
853
854    (
855        output
856            .as_collection()
857            .map(|result: UpsertValue| match result {
858                Ok(ok) => Ok(ok),
859                Err(err) => Err(DataflowError::from(EnvelopeError::Upsert(*err))),
860            }),
861        health_stream,
862        snapshot_stream,
863        shutdown_button.press_on_drop(),
864    )
865}
866
867/// The flavor-specific pieces of the upsert-v2 operator: the stash batcher,
868/// the feedback arrangement's spine, and how the drain reads prior state.
869/// [`build_upsert_operator`] holds the flavor-independent operator loop and
870/// calls through this trait at the few points where the flavors diverge.
871trait UpsertStashArm<T, O>
872where
873    T: Timestamp + Lattice + columnar::Columnar + Default,
874    for<'a> columnar::Ref<'a, T>: Copy + Ord,
875    O: columnar::Columnar + Default + Ord + Clone + Send + Sync + 'static,
876    for<'a> columnar::Ref<'a, O>: Ord + Copy,
877{
878    /// The feedback arrangement's spine. `'static` because the operator
879    /// future owns a trace agent for it.
880    type Spine: TraceReader<Time = T> + 'static;
881    /// The source-stash batcher. `'static` because the operator future owns
882    /// it.
883    type Batcher: Batcher<Time = T> + 'static;
884
885    /// A new stash batcher for one source dataflow.
886    fn new_batcher() -> Self::Batcher;
887
888    /// Push one sorted, consolidated chunk body into the batcher, in the
889    /// batcher's chunk representation.
890    fn push_chunk(batcher: &mut Self::Batcher, chunk: ColumnBody<UpsertUpdate<T, O>>);
891
892    /// Consolidate `updates` through `chunker` into chunk bodies and push
893    /// them into `batcher`, emptying `updates` (keeping its capacity). The
894    /// chunker readies a fully-consolidated chunk per `push_into`, so the
895    /// `extract` loop drains everything it produced.
896    fn flush(
897        updates: &mut Vec<UpsertUpdate<T, O>>,
898        chunker: &mut UpsertChunker<T, O>,
899        batcher: &mut Self::Batcher,
900    ) {
901        use timely::container::{ContainerBuilder as _, PushInto as _};
902        if updates.is_empty() {
903            return;
904        }
905        let mut raw: Column<UpsertUpdate<T, O>> = Default::default();
906        for update in updates.drain(..) {
907            raw.push_into(&update);
908        }
909        chunker.push_into(&mut raw);
910        while let Some(chunk) = chunker.extract() {
911            Self::push_chunk(batcher, std::mem::take(chunk));
912        }
913    }
914
915    /// Classify one sealed stash against `persist_upper` and emit eligible
916    /// output; see [`DrainStats`].
917    async fn drain(
918        sealed: Vec<<Self::Batcher as Batcher>::Output>,
919        ineligible: &mut Vec<UpsertUpdate<T, O>>,
920        output_handle: &UpsertOutputHandle<T>,
921        output_cap: &Capability<T>,
922        persist_upper: &Antichain<T>,
923        trace: &mut TraceAgent<Self::Spine>,
924        worker_id: usize,
925        source_id: GlobalId,
926    ) -> DrainStats;
927}
928
929/// The chunked flavor: [`UpsertChunkBatcher`] stash, [`FeedbackSpine`]
930/// feedback arrangement, bulk-probe drain.
931///
932/// NOTE: the seal's chain merge loads the bodies of the chunks it actually
933/// merges, while untouched survivors keep their spilled bodies. The
934/// partition against the upper passes chunks whose resident time bounds fall
935/// entirely on one side through without loading them, so only chunks the
936/// upper splits round-trip the pool codec. The drain then loads sealed
937/// chunks one at a time, so a large drain (a frontier advance releasing a
938/// snapshot's worth of stash at once) holds at most one chunk resident
939/// rather than the whole backlog.
940struct ChunkedArm;
941
942impl<T, O> UpsertStashArm<T, O> for ChunkedArm
943where
944    T: Timestamp + TotalOrder + Lattice + Sync,
945    T: columnation::Columnation + columnar::Columnar + Default,
946    for<'a> columnar::Ref<'a, T>: Copy + Ord,
947    O: columnar::Columnar + Default + Ord + Clone + Send + Sync + 'static,
948    for<'a> columnar::Ref<'a, O>: Ord + Copy,
949{
950    type Spine = FeedbackSpine<T>;
951    type Batcher = UpsertChunkBatcher<T, O>;
952
953    fn new_batcher() -> Self::Batcher {
954        Batcher::new(None, 0)
955    }
956
957    fn push_chunk(batcher: &mut Self::Batcher, chunk: ColumnBody<UpsertUpdate<T, O>>) {
958        batcher.push_into(ColumnChunk::from_body(chunk));
959    }
960
961    async fn drain(
962        sealed: Vec<UpsertChunk<T, O>>,
963        ineligible: &mut Vec<UpsertUpdate<T, O>>,
964        output_handle: &UpsertOutputHandle<T>,
965        output_cap: &Capability<T>,
966        persist_upper: &Antichain<T>,
967        trace: &mut TraceAgent<Self::Spine>,
968        worker_id: usize,
969        source_id: GlobalId,
970    ) -> DrainStats {
971        drain_sealed_input_chunked(
972            sealed.into_iter().map(ColumnChunk::into_body),
973            ineligible,
974            output_handle,
975            output_cap,
976            persist_upper,
977            trace,
978            worker_id,
979            source_id,
980        )
981        .await
982    }
983}
984
985/// The paged flavor: [`UpsertPagedBatcher`] stash, `ValRowSpine` feedback
986/// arrangement, cursor drain. Cold chains page out of RSS through the
987/// storage-owned column pager, captured once per batcher at construction.
988struct PagedArm;
989
990impl<T, O> UpsertStashArm<T, O> for PagedArm
991where
992    T: Timestamp + TotalOrder + Lattice + Sync,
993    T: columnation::Columnation + columnar::Columnar + Default,
994    for<'a> columnar::Ref<'a, T>: Copy + Ord,
995    O: columnar::Columnar + Default + Ord + Clone + Send + Sync + 'static,
996    for<'a> columnar::Ref<'a, O>: Ord + Copy,
997{
998    type Spine = FundedValRowSpine<UpsertKey, T, Diff>;
999    type Batcher = UpsertPagedBatcher<T, O>;
1000
1001    fn new_batcher() -> Self::Batcher {
1002        let mut batcher: UpsertPagedBatcher<T, O> = Batcher::new(None, 0);
1003        batcher.set_pager(crate::upsert::upsert_stash_pager::pager());
1004        batcher
1005    }
1006
1007    fn push_chunk(batcher: &mut Self::Batcher, chunk: ColumnBody<UpsertUpdate<T, O>>) {
1008        // The pager's chains are edge containers, so the body goes back onto one.
1009        batcher.push_into(Column::from(chunk));
1010    }
1011
1012    async fn drain(
1013        sealed: Vec<Column<UpsertUpdate<T, O>>>,
1014        ineligible: &mut Vec<UpsertUpdate<T, O>>,
1015        output_handle: &UpsertOutputHandle<T>,
1016        output_cap: &Capability<T>,
1017        persist_upper: &Antichain<T>,
1018        trace: &mut TraceAgent<Self::Spine>,
1019        worker_id: usize,
1020        source_id: GlobalId,
1021    ) -> DrainStats {
1022        drain_sealed_input_paged(
1023            sealed,
1024            ineligible,
1025            output_handle,
1026            output_cap,
1027            persist_upper,
1028            trace,
1029            worker_id,
1030            source_id,
1031        )
1032        .await
1033    }
1034}
1035
1036/// Where a stashed entry's time falls relative to the feedback frontier.
1037enum TimeClass {
1038    /// `ts < persist_upper`: already persisted, dropped by the drain.
1039    AlreadyPersisted,
1040    /// `ts == persist_upper`: processed now.
1041    Eligible,
1042    /// `ts > persist_upper`: re-stashed until persist catches up.
1043    Ineligible,
1044}
1045
1046/// Classify `ts` against `persist_upper`: the single spelling of the drain's
1047/// eligibility test. The chunked drain's probe-collection pass and its
1048/// classification pass, and the paged drain's cursor walk, all call this, so
1049/// the probe set and the classification cannot disagree. Under the
1050/// operator's total order, `Eligible` means `ts` equals the frontier's one
1051/// element.
1052fn classify_time<T: PartialOrder>(persist_upper: &Antichain<T>, ts: &T) -> TimeClass {
1053    if !persist_upper.less_equal(ts) {
1054        TimeClass::AlreadyPersisted
1055    } else if persist_upper.less_than(ts) {
1056        TimeClass::Ineligible
1057    } else {
1058        TimeClass::Eligible
1059    }
1060}
1061
1062/// Counts from a single call to [`drain_sealed_input_chunked`] or
1063/// [`drain_sealed_input_paged`], used to update metrics.
1064struct DrainStats {
1065    /// Number of eligible entries probed against the feedback trace.
1066    eligible: u64,
1067    /// Number of probed entries for which a prior value was found.
1068    result_count: u64,
1069    /// New value written with no prior value (insert).
1070    inserts: u64,
1071    /// New value written over an existing value (update).
1072    updates: u64,
1073    /// Tombstone (None) applied to an existing value (delete).
1074    deletes: u64,
1075    /// Total output records emitted (retractions + insertions).
1076    output_count: u64,
1077}
1078
1079/// Process sealed chunks from the batcher, classifying each entry by its
1080/// timestamp relative to `persist_upper`:
1081///
1082///   * `ts == persist_upper`: eligible for processing now (bulk probe of the
1083///     feedback trace + output).
1084///   * `ts >  persist_upper`: not yet processable. Returned in `ineligible`
1085///     for re-stashing until the feedback frontier catches up to it.
1086///   * `ts <  persist_upper`: already persisted by some writer and not
1087///     relevant anymore, so DROPPED. The downstream persist_sink would filter
1088///     such updates out anyway since the shard upper is further ahead, and
1089///     our state is already up-to-date to `persist_upper` so we could not
1090///     emit correct retractions for it. Re-stashing it would strand the data
1091///     forever (`persist_upper` only advances, so `ts == persist_upper` can
1092///     never again hold) and pin the operator's output frontier below the
1093///     shard upper. This mirrors v1's `relevant = persist_upper.less_equal(ts)`.
1094///
1095/// The sealed chunks are already sorted and consolidated by the merge
1096/// batcher, so each chunk's eligible keys form a sorted, deduplicated probe
1097/// set, and the prior state comes back through one bulk `extract_into` pass
1098/// per trace batch. Resident fence metadata selects the touched trace
1099/// chunks, and only those bodies are read back, copy-out, per chunk probed.
1100/// Sealed chunks are pulled from the iterator one at a time and dropped
1101/// before the next is requested, so at most one loaded stash chunk (plus the
1102/// probe hits for its keys) is resident regardless of drain size. Only the
1103/// re-stashed ineligible set is materialized.
1104async fn drain_sealed_input_chunked<T, O>(
1105    sealed: impl Iterator<Item = ColumnBody<UpsertUpdate<T, O>>>,
1106    ineligible: &mut Vec<UpsertUpdate<T, O>>,
1107    output_handle: &UpsertOutputHandle<T>,
1108    output_cap: &Capability<T>,
1109    persist_upper: &Antichain<T>,
1110    trace: &mut TraceAgent<FeedbackSpine<T>>,
1111    worker_id: usize,
1112    source_id: GlobalId,
1113) -> DrainStats
1114where
1115    T: Timestamp + TotalOrder + Lattice + Sync,
1116    T: columnation::Columnation + columnar::Columnar + Default,
1117    for<'a> columnar::Ref<'a, T>: Copy + Ord,
1118    O: columnar::Columnar,
1119{
1120    let mut eligible_count: u64 = 0;
1121    let mut result_count: u64 = 0;
1122    let mut output_count: u64 = 0;
1123    let mut inserts: u64 = 0;
1124    let mut updates: u64 = 0;
1125    let mut deletes: u64 = 0;
1126
1127    // The batches this drain reads against. `Rc` clones of the trace's
1128    // sealed chunk batches: chunk bodies stay spilled and are read back
1129    // copy-out, per probed chunk, inside `extract_into`.
1130    let batches = trace
1131        .batches_through(Antichain::new().borrow())
1132        .expect("complete batch set for the feedback trace; is it closed?");
1133
1134    // Eligible keys are probed against the trace in windows of this many
1135    // distinct keys, bounding what one pass stages resident: probe hits carry
1136    // full values, so an unwindowed pass over a byte-graded chunk of small
1137    // records could stage tens of thousands of values at once.
1138    const PROBE_WINDOW: usize = 1024;
1139
1140    for chunk in sealed {
1141        use columnar::{Index, Len};
1142        let view = chunk.borrow();
1143        let total = view.len();
1144        let mut start = 0;
1145        while start < total {
1146            // Pass 1: this window's sorted, deduplicated probe keys. The
1147            // chunk is sorted by (key, time), so a window is a contiguous
1148            // record range, probes come out sorted, and dedup is a neighbor
1149            // test. The window closes where its PROBE_WINDOW + 1st distinct
1150            // eligible key would begin.
1151            let mut probe_col = <UpsertKey as columnar::Columnar>::Container::default();
1152            let mut probe_count = 0usize;
1153            let mut end = total;
1154            {
1155                use columnar::Push;
1156                let mut last_probe: Option<&UpsertKey> = None;
1157                for index in start..total {
1158                    let (key, ts, _diff) = view.get(index);
1159                    let ts = <T as columnar::Columnar>::into_owned(ts);
1160                    if matches!(classify_time(persist_upper, &ts), TimeClass::Eligible) {
1161                        if last_probe != Some(key) {
1162                            if probe_count == PROBE_WINDOW {
1163                                end = index;
1164                                break;
1165                            }
1166                            probe_col.push(key);
1167                            probe_count += 1;
1168                            last_probe = Some(key);
1169                        }
1170                    }
1171                }
1172            }
1173
1174            // Pass 2: bulk-probe the trace's batches through the chunk
1175            // batches' `UnloadChunk` surface and consolidate the hits into
1176            // the prior value per key. Hits arrive per batch, so equal
1177            // `(key, val)` pairs from different batches are non-adjacent.
1178            // Sort before folding.
1179            let mut old_values: std::collections::BTreeMap<UpsertKey, UpsertValue> =
1180                std::collections::BTreeMap::new();
1181            if probe_count > 0 {
1182                use columnar::Borrow;
1183                let mut staging = <FeedbackUpdate<T> as columnar::Columnar>::Container::default();
1184                for batch in &batches {
1185                    batch.extract_into(probe_col.borrow(), &mut staging);
1186                }
1187                let staged = staging.borrow();
1188                let mut hits: Vec<_> = (0..staged.len())
1189                    .map(|i| {
1190                        let ((key, val), _time, diff) = staged.get(i);
1191                        (key, val, <Diff as columnar::Columnar>::into_owned(diff))
1192                    })
1193                    .collect();
1194                hits.sort_by(|a, b| (a.0, a.1).cmp(&(b.0, b.1)));
1195                let mut i = 0;
1196                while i < hits.len() {
1197                    let (key, val, _) = hits[i];
1198                    let mut count = Diff::ZERO;
1199                    let mut j = i;
1200                    while j < hits.len() && hits[j].0 == key && hits[j].1 == val {
1201                        count += hits[j].2;
1202                        j += 1;
1203                    }
1204                    if count.is_positive() {
1205                        assert!(
1206                            count == 1.into(),
1207                            "unexpected multiple entries for the same key in persist trace"
1208                        );
1209                        let prev = old_values.insert(*key, decode_upsert_value(val.iter()));
1210                        assert!(
1211                            prev.is_none(),
1212                            "unexpected multiple values for the same key in persist trace"
1213                        );
1214                    }
1215                    i = j;
1216                }
1217            }
1218
1219            // Pass 3: classify and emit this window's records.
1220            for index in start..end {
1221                let (key, ts, diff) = view.get(index);
1222                let ts = <T as columnar::Columnar>::into_owned(ts);
1223                match classify_time(persist_upper, &ts) {
1224                    TimeClass::AlreadyPersisted => continue,
1225                    TimeClass::Ineligible => {
1226                        // Re-stash for later (owned).
1227                        ineligible.push((
1228                            *key,
1229                            ts,
1230                            <UpsertDiff<O> as columnar::Columnar>::into_owned(diff),
1231                        ));
1232                        continue;
1233                    }
1234                    TimeClass::Eligible => {}
1235                }
1236
1237                // ts == persist_upper: eligible. The chunk holds one entry per
1238                // (key, time) and eligibility pins the time, so this key appears
1239                // at most once and its prior value can move out of the map.
1240                eligible_count += 1;
1241                let old_value = old_values.remove(key);
1242
1243                if old_value.is_some() {
1244                    result_count += 1;
1245                }
1246
1247                match diff.value {
1248                    Some(row) => {
1249                        if let Some(old_val) = old_value {
1250                            let size = upsert_value_byte_len(&old_val);
1251                            output_handle
1252                                .give_fueled(
1253                                    output_cap,
1254                                    (old_val, ts.clone(), Diff::MINUS_ONE),
1255                                    size,
1256                                )
1257                                .await;
1258                            output_count += 1;
1259                            updates += 1;
1260                        } else {
1261                            inserts += 1;
1262                        }
1263                        let new_val = decode_upsert_value(row.iter());
1264                        let size = upsert_value_byte_len(&new_val);
1265                        output_handle
1266                            .give_fueled(output_cap, (new_val, ts, Diff::ONE), size)
1267                            .await;
1268                        output_count += 1;
1269                    }
1270                    None => {
1271                        if let Some(old_val) = old_value {
1272                            let size = upsert_value_byte_len(&old_val);
1273                            output_handle
1274                                .give_fueled(output_cap, (old_val, ts, Diff::MINUS_ONE), size)
1275                                .await;
1276                            output_count += 1;
1277                            deletes += 1;
1278                        }
1279                    }
1280                }
1281            }
1282            start = end;
1283        }
1284    }
1285
1286    tracing::debug!(
1287        worker_id = %worker_id,
1288        source_id = %source_id,
1289        ineligible = ineligible.len(),
1290        eligible = eligible_count,
1291        "drained stash",
1292    );
1293
1294    DrainStats {
1295        eligible: eligible_count,
1296        result_count,
1297        inserts,
1298        updates,
1299        deletes,
1300        output_count,
1301    }
1302}
1303
1304/// [`drain_sealed_input_chunked`]'s counterpart for the paged flavor,
1305/// classifying entries the same way but reading prior state through a trace
1306/// cursor.
1307///
1308/// The sealed chunks are already sorted and consolidated by the merge
1309/// batcher, so the trace cursor walks forward through keys in order and
1310/// seeks amortize. Entries are walked by reference rather than collecting
1311/// the eligible set into an owned Vec, and eligible values are emitted
1312/// straight from the column's `RowRef` with no owned `UpsertDiff` copy. Only
1313/// the re-stashed ineligible set is materialized.
1314async fn drain_sealed_input_paged<T, O>(
1315    sealed: Vec<Column<UpsertUpdate<T, O>>>,
1316    ineligible: &mut Vec<UpsertUpdate<T, O>>,
1317    output_handle: &UpsertOutputHandle<T>,
1318    output_cap: &Capability<T>,
1319    persist_upper: &Antichain<T>,
1320    trace: &mut TraceAgent<FundedValRowSpine<UpsertKey, T, Diff>>,
1321    worker_id: usize,
1322    source_id: GlobalId,
1323) -> DrainStats
1324where
1325    T: Timestamp + TotalOrder + Lattice + Sync,
1326    T: columnation::Columnation + columnar::Columnar,
1327    O: columnar::Columnar,
1328{
1329    use columnar::Index as _;
1330
1331    // Classify each entry by its timestamp relative to `persist_upper`:
1332    //
1333    //   * `ts == persist_upper`: eligible for processing now.
1334    //   * `ts >  persist_upper`: not yet processable; re-stashed (ineligible)
1335    //     until the feedback frontier catches up to it.
1336    //   * `ts <  persist_upper`: already persisted by some writer and not
1337    //     relevant anymore. We DROP it. The downstream persist_sink would
1338    //     filter such updates out anyway since the shard upper is further
1339    //     ahead, and our state is already up-to-date to `persist_upper` so we
1340    //     could not emit correct retractions for it. Re-stashing it would
1341    //     strand the data forever (`persist_upper` only advances, so
1342    //     `ts == persist_upper` can never again hold) and pin the operator's
1343    //     output frontier below the shard upper. This mirrors v1's
1344    //     `relevant = persist_upper.less_equal(ts)`.
1345    let mut eligible_count: u64 = 0;
1346    let mut result_count: u64 = 0;
1347    let mut output_count: u64 = 0;
1348    let mut inserts: u64 = 0;
1349    let mut updates: u64 = 0;
1350    let mut deletes: u64 = 0;
1351
1352    let (mut cursor, storage) = trace.cursor();
1353
1354    for chunk in &sealed {
1355        for (key, ts, diff) in chunk.borrow().into_index_iter() {
1356            let ts = <T as columnar::Columnar>::into_owned(ts);
1357            match classify_time(persist_upper, &ts) {
1358                TimeClass::AlreadyPersisted => continue,
1359                TimeClass::Ineligible => {
1360                    // Re-stash for later (owned).
1361                    ineligible.push((
1362                        *key,
1363                        ts,
1364                        <UpsertDiff<O> as columnar::Columnar>::into_owned(diff),
1365                    ));
1366                    continue;
1367                }
1368                TimeClass::Eligible => {}
1369            }
1370
1371            // ts == persist_upper: eligible. Look up the prior value for this
1372            // key in the persist trace and emit the retraction / insertion. The
1373            // spine stores keys in a columnation arena, so we seek by the
1374            // column's borrowed `&UpsertKey` directly.
1375            eligible_count += 1;
1376            cursor.seek_key(&storage, key);
1377            let old_value = match cursor.get_key(&storage) {
1378                Some(found) if found == key => {
1379                    let mut result = None;
1380                    while let Some(val) = cursor.get_val(&storage) {
1381                        let mut count = Diff::ZERO;
1382                        cursor.map_times(&storage, |_time, d| {
1383                            count += d.clone();
1384                        });
1385                        if count.is_positive() {
1386                            assert!(
1387                                count == 1.into(),
1388                                "unexpected multiple entries for the same key in persist trace"
1389                            );
1390                            assert!(
1391                                result.is_none(),
1392                                "unexpected multiple values for the same key in persist trace"
1393                            );
1394                            result = Some(decode_upsert_value(val));
1395                        }
1396                        cursor.step_val(&storage);
1397                    }
1398                    result
1399                }
1400                _ => None,
1401            };
1402
1403            if old_value.is_some() {
1404                result_count += 1;
1405            }
1406
1407            match diff.value {
1408                Some(row) => {
1409                    if let Some(old_val) = old_value {
1410                        let size = upsert_value_byte_len(&old_val);
1411                        output_handle
1412                            .give_fueled(output_cap, (old_val, ts.clone(), Diff::MINUS_ONE), size)
1413                            .await;
1414                        output_count += 1;
1415                        updates += 1;
1416                    } else {
1417                        inserts += 1;
1418                    }
1419                    let new_val = decode_upsert_value(row.iter());
1420                    let size = upsert_value_byte_len(&new_val);
1421                    output_handle
1422                        .give_fueled(output_cap, (new_val, ts, Diff::ONE), size)
1423                        .await;
1424                    output_count += 1;
1425                }
1426                None => {
1427                    if let Some(old_val) = old_value {
1428                        let size = upsert_value_byte_len(&old_val);
1429                        output_handle
1430                            .give_fueled(output_cap, (old_val, ts, Diff::MINUS_ONE), size)
1431                            .await;
1432                        output_count += 1;
1433                        deletes += 1;
1434                    }
1435                }
1436            }
1437        }
1438    }
1439
1440    tracing::debug!(
1441        worker_id = %worker_id,
1442        source_id = %source_id,
1443        ineligible = ineligible.len(),
1444        eligible = eligible_count,
1445        "drained stash",
1446    );
1447
1448    DrainStats {
1449        eligible: eligible_count,
1450        result_count,
1451        inserts,
1452        updates,
1453        deletes,
1454        output_count,
1455    }
1456}
1457
1458#[cfg(test)]
1459mod test {
1460    //! No test drives errors through the persist feedback, so every harness below closes
1461    //! the error input by dropping its handle.
1462
1463    use mz_ore::metrics::MetricsRegistry;
1464    use mz_persist_types::ShardId;
1465    use mz_repr::{Datum, Timestamp as MzTimestamp};
1466    use mz_storage_operators::persist_source::Subtime;
1467    use mz_storage_types::sources::SourceEnvelope;
1468    use mz_storage_types::sources::envelope::{KeyEnvelope, UpsertEnvelope, UpsertStyle};
1469    use timely::dataflow::operators::capture::Extract;
1470    use timely::dataflow::operators::{Capture, Input};
1471    use timely::progress::Timestamp;
1472
1473    use crate::metrics::StorageMetrics;
1474    use crate::metrics::upsert::UpsertMetricDefs;
1475    use crate::source::SourceExportCreationConfig;
1476    use crate::statistics::{SourceStatistics, SourceStatisticsMetricDefs};
1477
1478    use super::*;
1479
1480    // The tests drive the operator with a plain integer `FromTime` standing in
1481    // for a Kafka offset; project it to itself so dedup orders by it directly.
1482    impl UpsertSourceTime for i32 {
1483        type Order = i32;
1484        fn upsert_order(&self) -> i32 {
1485            *self
1486        }
1487    }
1488
1489    type Ts = (MzTimestamp, Subtime);
1490
1491    fn new_ts(ts: u64) -> Ts {
1492        (MzTimestamp::new(ts), Subtime::minimum())
1493    }
1494
1495    fn key(k: i64) -> UpsertKey {
1496        UpsertKey::from_key(Ok(&Row::pack_slice(&[Datum::Int64(k)])))
1497    }
1498
1499    fn row(k: i64, v: i64) -> Row {
1500        Row::pack_slice(&[Datum::Int64(k), Datum::Int64(v)])
1501    }
1502
1503    // Runs the test body once per stash flavor and asserts the two flavors
1504    // produce identical (consolidated) output, so every scenario covers both
1505    // operator arms. Returns one flavor's output for the caller's own
1506    // expected-value assertion.
1507    macro_rules! upsert_test {
1508        (|$input:ident, $persist:ident, $worker:ident| $body:block) => {{
1509            let run = |flavor: UpsertStashFlavor| {
1510                let output_handle = timely::execute_directly(move |$worker| {
1511                    let (mut $input, mut $persist, output_handle) = $worker
1512                        .dataflow::<MzTimestamp, _, _>(|scope| {
1513                            scope.scoped::<Ts, _, _>("upsert", |scope| {
1514                                let (input_handle, input) = scope.new_input();
1515                                let (persist_handle, persist_input) = scope.new_input();
1516                                let (_persist_err_handle, persist_err_input) = scope.new_input();
1517                                let source_id = GlobalId::User(0);
1518
1519                                let reg = MetricsRegistry::new();
1520                                let upsert_defs = UpsertMetricDefs::register_with(&reg);
1521                                let upsert_metrics =
1522                                    UpsertMetrics::new(&upsert_defs, source_id, 0, None);
1523
1524                                let reg2 = MetricsRegistry::new();
1525                                let storage_metrics = StorageMetrics::register_with(&reg2);
1526
1527                                let reg3 = MetricsRegistry::new();
1528                                let stats_defs =
1529                                    SourceStatisticsMetricDefs::register_with(&reg3);
1530                                let envelope = SourceEnvelope::Upsert(UpsertEnvelope {
1531                                    source_arity: 2,
1532                                    style: UpsertStyle::Default(KeyEnvelope::Flattened),
1533                                    key_indices: vec![0],
1534                                });
1535                                let source_statistics = SourceStatistics::new(
1536                                    source_id, 0, &stats_defs, source_id, &ShardId::new(),
1537                                    envelope, Antichain::from_elem(Timestamp::minimum()),
1538                                );
1539                                let source_config = SourceExportCreationConfig {
1540                                    id: source_id,
1541                                    worker_id: 0,
1542                                    metrics: storage_metrics,
1543                                    source_statistics,
1544                                };
1545
1546                                let (output, _, _, button) = upsert_inner(
1547                                    flavor,
1548                                    input.as_collection(),
1549                                    vec![0],
1550                                    Antichain::from_elem(Timestamp::minimum()),
1551                                    persist_input.as_collection(),
1552                                    persist_err_input.as_collection(),
1553                                    None,
1554                                    upsert_metrics,
1555                                    source_config,
1556                                );
1557                                std::mem::forget(button);
1558                                (input_handle, persist_handle, output.inner.capture())
1559                            })
1560                        });
1561
1562                    $body
1563
1564                    output_handle
1565                });
1566
1567                let mut actual: Vec<_> = output_handle
1568                    .extract()
1569                    .into_iter()
1570                    .flat_map(|(_cap, container)| container)
1571                    .collect();
1572                differential_dataflow::consolidation::consolidate_updates(&mut actual);
1573                actual
1574            };
1575
1576            let paged = run(UpsertStashFlavor::Paged);
1577            let chunked = run(UpsertStashFlavor::Chunked);
1578            assert_eq!(paged, chunked, "stash flavors must produce equal output");
1579            chunked
1580        }};
1581    }
1582
1583    #[mz_ore::test]
1584    #[cfg_attr(miri, ignore)]
1585    fn gh_9160_repro() {
1586        let actual = upsert_test!(|input, persist, worker| {
1587            let key0 = key(0);
1588            let key1 = key(1);
1589            let value1 = row(0, 0);
1590            let value3 = row(0, 1);
1591            let value4 = row(0, 2);
1592
1593            input.send(((key0, Some(Ok(value1.clone())), 1), new_ts(0), Diff::ONE));
1594            input.advance_to(new_ts(2));
1595            worker.step();
1596
1597            persist.send((value1, new_ts(0), Diff::ONE));
1598            persist.advance_to(new_ts(1));
1599            worker.step();
1600
1601            input.send_batch(&mut vec![
1602                ((key1, None, 2), new_ts(2), Diff::ONE),
1603                ((key0, Some(Ok(value3)), 3), new_ts(3), Diff::ONE),
1604            ]);
1605            input.advance_to(new_ts(3));
1606            input.send_batch(&mut vec![(
1607                (key0, Some(Ok(value4)), 4),
1608                new_ts(3),
1609                Diff::ONE,
1610            )]);
1611            input.advance_to(new_ts(4));
1612            worker.step();
1613
1614            persist.advance_to(new_ts(3));
1615            worker.step();
1616        });
1617
1618        let value1 = row(0, 0);
1619        let value4 = row(0, 2);
1620        let expected: Vec<(Result<Row, DataflowError>, _, _)> = vec![
1621            (Ok(value1.clone()), new_ts(0), Diff::ONE),
1622            (Ok(value1), new_ts(3), Diff::MINUS_ONE),
1623            (Ok(value4), new_ts(3), Diff::ONE),
1624        ];
1625        assert_eq!(actual, expected);
1626    }
1627
1628    #[mz_ore::test]
1629    #[cfg_attr(miri, ignore)]
1630    fn out_of_order_keys_across_timestamps() {
1631        let actual = upsert_test!(|input, persist, worker| {
1632            let key_high = key(99);
1633            let key_low = key(1);
1634            let val_a = row(99, 1);
1635            let val_b = row(1, 2);
1636
1637            input.send(((key_high, Some(Ok(val_a.clone())), 1), new_ts(0), Diff::ONE));
1638            input.advance_to(new_ts(1));
1639            worker.step();
1640            persist.send((val_a.clone(), new_ts(0), Diff::ONE));
1641            persist.advance_to(new_ts(1));
1642            worker.step();
1643
1644            input.send(((key_low, Some(Ok(val_b.clone())), 2), new_ts(1), Diff::ONE));
1645            input.advance_to(new_ts(2));
1646            worker.step();
1647            persist.send((val_b.clone(), new_ts(1), Diff::ONE));
1648            persist.advance_to(new_ts(2));
1649            worker.step();
1650
1651            let val_a2 = row(99, 10);
1652            let val_b2 = row(1, 20);
1653            input.send_batch(&mut vec![
1654                (
1655                    (key_high, Some(Ok(val_a2.clone())), 3),
1656                    new_ts(2),
1657                    Diff::ONE,
1658                ),
1659                ((key_low, Some(Ok(val_b2.clone())), 4), new_ts(2), Diff::ONE),
1660            ]);
1661            input.advance_to(new_ts(3));
1662            worker.step();
1663            persist.advance_to(new_ts(3));
1664            worker.step();
1665        });
1666
1667        let val_a = row(99, 1);
1668        let val_b = row(1, 2);
1669        let val_a2 = row(99, 10);
1670        let val_b2 = row(1, 20);
1671        let expected: Vec<(Result<Row, DataflowError>, _, _)> = vec![
1672            (Ok(val_b.clone()), new_ts(1), Diff::ONE),
1673            (Ok(val_b), new_ts(2), Diff::MINUS_ONE),
1674            (Ok(val_b2), new_ts(2), Diff::ONE),
1675            (Ok(val_a.clone()), new_ts(0), Diff::ONE),
1676            (Ok(val_a), new_ts(2), Diff::MINUS_ONE),
1677            (Ok(val_a2), new_ts(2), Diff::ONE),
1678        ];
1679        let mut actual_sorted = actual;
1680        let mut expected_sorted = expected;
1681        actual_sorted.sort();
1682        expected_sorted.sort();
1683        assert_eq!(actual_sorted, expected_sorted);
1684    }
1685
1686    #[mz_ore::test]
1687    #[cfg_attr(miri, ignore)]
1688    fn rehydration_then_update() {
1689        let actual = upsert_test!(|input, persist, worker| {
1690            let k = key(42);
1691            let old_val = row(42, 100);
1692            let new_val = row(42, 200);
1693
1694            persist.send((old_val, new_ts(0), Diff::ONE));
1695            persist.advance_to(new_ts(1));
1696            worker.step();
1697
1698            input.send(((k, Some(Ok(new_val)), 1), new_ts(1), Diff::ONE));
1699            input.advance_to(new_ts(2));
1700            worker.step();
1701            persist.advance_to(new_ts(2));
1702            worker.step();
1703        });
1704
1705        let old_val = row(42, 100);
1706        let new_val = row(42, 200);
1707        let expected: Vec<(Result<Row, DataflowError>, _, _)> = vec![
1708            (Ok(old_val), new_ts(1), Diff::MINUS_ONE),
1709            (Ok(new_val), new_ts(1), Diff::ONE),
1710        ];
1711        assert_eq!(actual, expected);
1712    }
1713
1714    #[mz_ore::test]
1715    #[cfg_attr(miri, ignore)]
1716    fn drain_crosses_probe_window() {
1717        // More distinct keys than one drain probe window holds, all eligible
1718        // at the same timestamp. The records fit one sealed chunk (the ship
1719        // threshold is ~2 MiB), so the drain must close a probe window
1720        // mid-chunk and carry keys correctly across the boundary.
1721        const KEYS: i64 = 1500;
1722        let actual = upsert_test!(|input, persist, worker| {
1723            for k in 0..KEYS {
1724                persist.send((row(k, k), new_ts(0), Diff::ONE));
1725            }
1726            persist.advance_to(new_ts(1));
1727            worker.step();
1728
1729            for k in 0..KEYS {
1730                input.send(((key(k), Some(Ok(row(k, k + 1))), 1), new_ts(1), Diff::ONE));
1731            }
1732            input.advance_to(new_ts(2));
1733            worker.step();
1734            persist.advance_to(new_ts(2));
1735            worker.step();
1736        });
1737
1738        let mut expected: Vec<(Result<Row, DataflowError>, _, _)> = Vec::new();
1739        for k in 0..KEYS {
1740            expected.push((Ok(row(k, k)), new_ts(1), Diff::MINUS_ONE));
1741            expected.push((Ok(row(k, k + 1)), new_ts(1), Diff::ONE));
1742        }
1743        let mut actual_sorted = actual;
1744        actual_sorted.sort();
1745        expected.sort();
1746        assert_eq!(actual_sorted, expected);
1747    }
1748
1749    /// The probe-window scenario with committed chunks forced through a
1750    /// private buffer pool, so the chunked flavor's seal and drain read
1751    /// spilled bodies back through the pool codec rather than resident
1752    /// memory. The override is thread-scoped and `execute_directly` runs the
1753    /// worker on this thread, so it reaches the operator's batchers. The
1754    /// paged flavor (which the harness also runs) routes through the column
1755    /// pager rather than the chunk override, so it stays resident and serves
1756    /// as the reference.
1757    #[mz_ore::test]
1758    #[cfg_attr(miri, ignore)]
1759    fn drain_reads_spilled_chunks() {
1760        use mz_ore::pool::Pool;
1761        use mz_timely_util::columnar::chunk::set_spill_override;
1762
1763        let pool = Pool::new().expect("pool creation");
1764        set_spill_override(Some(pool.clone()));
1765
1766        const KEYS: i64 = 1500;
1767        let actual = upsert_test!(|input, persist, worker| {
1768            for k in 0..KEYS {
1769                persist.send((row(k, k), new_ts(0), Diff::ONE));
1770            }
1771            persist.advance_to(new_ts(1));
1772            worker.step();
1773
1774            for k in 0..KEYS {
1775                input.send(((key(k), Some(Ok(row(k, k + 1))), 1), new_ts(1), Diff::ONE));
1776            }
1777            input.advance_to(new_ts(2));
1778            worker.step();
1779            persist.advance_to(new_ts(2));
1780            worker.step();
1781        });
1782
1783        set_spill_override(None);
1784        assert!(
1785            pool.stats().inserts > 0,
1786            "chunks should have spilled through the pool"
1787        );
1788
1789        let mut expected: Vec<(Result<Row, DataflowError>, _, _)> = Vec::new();
1790        for k in 0..KEYS {
1791            expected.push((Ok(row(k, k)), new_ts(1), Diff::MINUS_ONE));
1792            expected.push((Ok(row(k, k + 1)), new_ts(1), Diff::ONE));
1793        }
1794        let mut actual_sorted = actual;
1795        actual_sorted.sort();
1796        expected.sort();
1797        assert_eq!(actual_sorted, expected);
1798    }
1799
1800    #[mz_ore::test]
1801    #[cfg_attr(miri, ignore)]
1802    fn delete_existing_key() {
1803        let actual = upsert_test!(|input, persist, worker| {
1804            let k = key(7);
1805            let val = row(7, 77);
1806
1807            input.send(((k, Some(Ok(val.clone())), 1), new_ts(0), Diff::ONE));
1808            input.advance_to(new_ts(1));
1809            worker.step();
1810            persist.send((val, new_ts(0), Diff::ONE));
1811            persist.advance_to(new_ts(1));
1812            worker.step();
1813
1814            input.send(((k, None, 2), new_ts(1), Diff::ONE));
1815            input.advance_to(new_ts(2));
1816            worker.step();
1817            persist.advance_to(new_ts(2));
1818            worker.step();
1819        });
1820
1821        let val = row(7, 77);
1822        let expected: Vec<(Result<Row, DataflowError>, _, _)> = vec![
1823            (Ok(val.clone()), new_ts(0), Diff::ONE),
1824            (Ok(val), new_ts(1), Diff::MINUS_ONE),
1825        ];
1826        assert_eq!(actual, expected);
1827    }
1828
1829    #[mz_ore::test]
1830    #[cfg_attr(miri, ignore)]
1831    fn multi_batch_rehydration() {
1832        let actual = upsert_test!(|input, persist, worker| {
1833            let k = key(5);
1834            let old_val = row(5, 10);
1835            let new_val = row(5, 20);
1836            let updated_val = row(5, 30);
1837
1838            persist.send((old_val.clone(), new_ts(0), Diff::ONE));
1839            persist.send((old_val, new_ts(0), Diff::MINUS_ONE));
1840            persist.send((new_val, new_ts(0), Diff::ONE));
1841            persist.advance_to(new_ts(1));
1842            worker.step();
1843
1844            input.send(((k, Some(Ok(updated_val)), 1), new_ts(1), Diff::ONE));
1845            input.advance_to(new_ts(2));
1846            worker.step();
1847            persist.advance_to(new_ts(2));
1848            worker.step();
1849        });
1850
1851        let new_val = row(5, 20);
1852        let updated_val = row(5, 30);
1853        let expected: Vec<(Result<Row, DataflowError>, _, _)> = vec![
1854            (Ok(new_val), new_ts(1), Diff::MINUS_ONE),
1855            (Ok(updated_val), new_ts(1), Diff::ONE),
1856        ];
1857        assert_eq!(actual, expected);
1858    }
1859
1860    #[mz_ore::test]
1861    #[cfg_attr(miri, ignore)]
1862    fn delete_nonexistent_key() {
1863        let actual = upsert_test!(|input, persist, worker| {
1864            let k = key(99);
1865
1866            persist.advance_to(new_ts(1));
1867            worker.step();
1868
1869            input.send(((k, None, 1), new_ts(1), Diff::ONE));
1870            input.advance_to(new_ts(2));
1871            worker.step();
1872            persist.advance_to(new_ts(2));
1873            worker.step();
1874        });
1875
1876        assert!(actual.is_empty(), "expected empty output, got: {actual:?}");
1877    }
1878
1879    #[mz_ore::test]
1880    #[cfg_attr(miri, ignore)]
1881    fn reinsert_after_delete() {
1882        let actual = upsert_test!(|input, persist, worker| {
1883            let k = key(3);
1884            let val_a = row(3, 10);
1885            let val_b = row(3, 20);
1886
1887            input.send(((k, Some(Ok(val_a.clone())), 1), new_ts(0), Diff::ONE));
1888            input.advance_to(new_ts(1));
1889            worker.step();
1890            persist.send((val_a.clone(), new_ts(0), Diff::ONE));
1891            persist.advance_to(new_ts(1));
1892            worker.step();
1893
1894            input.send(((k, None, 2), new_ts(1), Diff::ONE));
1895            input.advance_to(new_ts(2));
1896            worker.step();
1897            persist.send((val_a, new_ts(1), Diff::MINUS_ONE));
1898            persist.advance_to(new_ts(2));
1899            worker.step();
1900
1901            input.send(((k, Some(Ok(val_b.clone())), 3), new_ts(2), Diff::ONE));
1902            input.advance_to(new_ts(3));
1903            worker.step();
1904            persist.advance_to(new_ts(3));
1905            worker.step();
1906        });
1907
1908        let val_a = row(3, 10);
1909        let val_b = row(3, 20);
1910        let mut expected: Vec<(Result<Row, DataflowError>, _, _)> = vec![
1911            (Ok(val_a.clone()), new_ts(0), Diff::ONE),
1912            (Ok(val_a), new_ts(1), Diff::MINUS_ONE),
1913            (Ok(val_b), new_ts(2), Diff::ONE),
1914        ];
1915        expected.sort();
1916        let mut actual = actual;
1917        actual.sort();
1918        assert_eq!(actual, expected);
1919    }
1920
1921    #[mz_ore::test]
1922    #[cfg_attr(miri, ignore)]
1923    fn idempotent_update() {
1924        let actual = upsert_test!(|input, persist, worker| {
1925            let k = key(11);
1926            let val = row(11, 50);
1927
1928            input.send(((k, Some(Ok(val.clone())), 1), new_ts(0), Diff::ONE));
1929            input.advance_to(new_ts(1));
1930            worker.step();
1931            persist.send((val.clone(), new_ts(0), Diff::ONE));
1932            persist.advance_to(new_ts(1));
1933            worker.step();
1934
1935            input.send(((k, Some(Ok(val.clone())), 2), new_ts(1), Diff::ONE));
1936            input.advance_to(new_ts(2));
1937            worker.step();
1938            persist.advance_to(new_ts(2));
1939            worker.step();
1940        });
1941
1942        let val = row(11, 50);
1943        let expected: Vec<(Result<Row, DataflowError>, _, _)> =
1944            vec![(Ok(val), new_ts(0), Diff::ONE)];
1945        assert_eq!(actual, expected);
1946    }
1947
1948    /// Operator-level repro of the 0dt read-only-handoff stranding bug.
1949    ///
1950    /// Models a lagging replacement generation: the external (old) writer has
1951    /// already advanced the shard — and therefore the feedback `persist_upper`
1952    /// — to `T = 10`, while the operator itself has emitted nothing. The
1953    /// lagging replacement now produces source data at timestamps BELOW that
1954    /// upper (`ts = 5, 7`), i.e. data the external writer has already persisted.
1955    ///
1956    /// The drain DROPS such already-persisted data (it satisfies
1957    /// neither `ts == persist_upper` nor `ts > persist_upper`), mirroring v1's
1958    /// `relevant = persist_upper.less_equal(ts)`. Were it instead re-stashed,
1959    /// the data would be stranded forever — `persist_upper` only advances, so
1960    /// `ts == persist_upper` could never again hold — and `min_ineligible_ts`
1961    /// would pin the operator's output capability at `ts = 5`, BELOW the shard
1962    /// upper, where it would stay for good. Dropping it lets the frontier
1963    /// advance freely.
1964    #[mz_ore::test]
1965    #[cfg_attr(miri, ignore)]
1966    fn lagging_replacement_below_upper_strands_data() {
1967        for flavor in [UpsertStashFlavor::Paged, UpsertStashFlavor::Chunked] {
1968            let (frontier, emitted) = run_below_upper_scenario_v2(flavor);
1969
1970            // The below-upper data is discarded (no output) and the output
1971            // frontier is not pinned below the shard upper (10); it advances
1972            // to the input upper (11), matching v1's behavior.
1973            assert!(
1974                emitted.is_empty(),
1975                "below-upper data should be dropped, not emitted; got {emitted:?} ({flavor:?})"
1976            );
1977            assert_eq!(
1978                frontier,
1979                vec![new_ts(11)],
1980                "v2 output frontier should advance to the input upper, not pin below \
1981                 persist_upper ({flavor:?})"
1982            );
1983            assert!(
1984                frontier[0] >= new_ts(10),
1985                "v2 output frontier {frontier:?} should reach at least persist_upper (10) \
1986                 ({flavor:?})"
1987            );
1988        }
1989    }
1990
1991    /// Shared driver for the lagging-replacement scenario against v2. Returns
1992    /// `(output_frontier, consolidated_emitted_updates)`.
1993    fn run_below_upper_scenario_v2(
1994        flavor: UpsertStashFlavor,
1995    ) -> (Vec<Ts>, Vec<(Result<Row, DataflowError>, Ts, Diff)>) {
1996        use timely::dataflow::operators::Probe;
1997
1998        let (frontier, capture) = timely::execute_directly(move |worker| {
1999            let (mut input, mut persist, probe, capture) =
2000                worker.dataflow::<MzTimestamp, _, _>(|scope| {
2001                    scope.scoped::<Ts, _, _>("upsert", |scope| {
2002                        let (input_handle, input) = scope.new_input();
2003                        let (persist_handle, persist_input) = scope.new_input();
2004                        let (_persist_err_handle, persist_err_input) = scope.new_input();
2005                        let source_id = GlobalId::User(0);
2006
2007                        let reg = MetricsRegistry::new();
2008                        let upsert_defs = UpsertMetricDefs::register_with(&reg);
2009                        let upsert_metrics = UpsertMetrics::new(&upsert_defs, source_id, 0, None);
2010
2011                        let reg2 = MetricsRegistry::new();
2012                        let storage_metrics = StorageMetrics::register_with(&reg2);
2013
2014                        let reg3 = MetricsRegistry::new();
2015                        let stats_defs = SourceStatisticsMetricDefs::register_with(&reg3);
2016                        let envelope = SourceEnvelope::Upsert(UpsertEnvelope {
2017                            source_arity: 2,
2018                            style: UpsertStyle::Default(KeyEnvelope::Flattened),
2019                            key_indices: vec![0],
2020                        });
2021                        let source_statistics = SourceStatistics::new(
2022                            source_id,
2023                            0,
2024                            &stats_defs,
2025                            source_id,
2026                            &ShardId::new(),
2027                            envelope,
2028                            Antichain::from_elem(Timestamp::minimum()),
2029                        );
2030                        let source_config = SourceExportCreationConfig {
2031                            id: source_id,
2032                            worker_id: 0,
2033                            metrics: storage_metrics,
2034                            source_statistics,
2035                        };
2036
2037                        let (output, _, _, button) = upsert_inner(
2038                            flavor,
2039                            input.as_collection(),
2040                            vec![0],
2041                            Antichain::from_elem(Timestamp::minimum()),
2042                            persist_input.as_collection(),
2043                            persist_err_input.as_collection(),
2044                            None,
2045                            upsert_metrics,
2046                            source_config,
2047                        );
2048                        std::mem::forget(button);
2049                        let (probe, stream) = output.inner.probe();
2050                        (input_handle, persist_handle, probe, stream.capture())
2051                    })
2052                });
2053
2054            // The external writer has advanced the shard (feedback persist_upper)
2055            // to T = 10 WITHOUT the operator emitting anything itself.
2056            persist.advance_to(new_ts(10));
2057            for _ in 0..20 {
2058                worker.step();
2059            }
2060
2061            // The lagging replacement produces source data at ts BELOW the
2062            // current persist_upper (5 and 7 while persist_upper = 10).
2063            input.send(((key(0), Some(Ok(row(0, 1))), 1), new_ts(5), Diff::ONE));
2064            input.send(((key(1), Some(Ok(row(1, 2))), 2), new_ts(7), Diff::ONE));
2065            input.advance_to(new_ts(11));
2066            for _ in 0..20 {
2067                worker.step();
2068            }
2069
2070            (probe.with_frontier(|f| f.to_vec()), capture)
2071        });
2072
2073        let mut emitted: Vec<_> = capture
2074            .extract()
2075            .into_iter()
2076            .flat_map(|(_cap, c)| c)
2077            .collect();
2078        differential_dataflow::consolidation::consolidate_updates(&mut emitted);
2079        (frontier, emitted)
2080    }
2081}