Skip to main content

mz_storage/render/
persist_sink.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//! Render an operator that persists a source collection.
11//!
12//! ## Implementation
13//!
14//! This module defines the `persist_sink` operator, that writes
15//! a collection produced by source rendering into a persist shard.
16//!
17//! It attempts to use all workers to write data to persist, and uses
18//! single-instance workers to coordinate work. The below diagram
19//! is an overview how it it shaped. There is more information
20//! in the doc comments of the top-level functions of this module.
21//!
22//!```text
23//!
24//!                                       ,------------.
25//!                                       | source     |
26//!                                       | collection |
27//!                                       +---+--------+
28//!                                       /   |
29//!                                      /    |
30//!                                     /     |
31//!                                    /      |
32//!                                   /       |
33//!                                  /        |
34//!                                 /         |
35//!                                /          |
36//!                               /     ,-+-----------------------.
37//!                              /      | mint_batch_descriptions |
38//!                             /       | one arbitrary worker    |
39//!                            |        +-,--,--------+----+------+
40//!                           ,----------´.-´         |     \
41//!                       _.-´ |       .-´            |      \
42//!                   _.-´     |    .-´               |       \
43//!                .-´  .------+----|-------+---------|--------\-----.
44//!               /    /            |       |         |         \     \
45//!        ,--------------.   ,-----------------.     |     ,-----------------.
46//!        | write_batches|   |  write_batches  |     |     |  write_batches  |
47//!        | worker 0     |   | worker 1        |     |     | worker N        |
48//!        +-----+--------+   +-+---------------+     |     +--+--------------+
49//!               \              \                    |        /
50//!                `-.            `,                  |       /
51//!                   `-._          `-.               |      /
52//!                       `-._         `-.            |     /
53//!                           `---------. `-.         |    /
54//!                                     +`---`---+-------------,
55//!                                     | append_batches       |
56//!                                     | one arbitrary worker |
57//!                                     +------+---------------+
58//!```
59//!
60//! `mint_batch_descriptions` also takes the ingestion's remap upper as a disconnected input, and
61//! broadcasts a committed ceiling to the writers, which `append_batches` does not see. Both edges
62//! are left out of the diagram to avoid clutter.
63//!
64//! ## Similarities with `mz_compute::sink::persist_sink`
65//!
66//! This module has many similarities with the compute version of
67//! the same concept, and in fact, is entirely derived from it.
68//!
69//! Compute requires that its `persist_sink` is _self-correcting_;
70//! that is, it corrects what the collection in persist
71//! accumulates to if the collection has values changed at
72//! previous timestamps. It does this by continually comparing
73//! the input stream with the collection as read back from persist.
74//!
75//! Source collections, while definite, cannot be reliably by
76//! re-produced once written down, which means compute's
77//! `persist_sink`'s self-correction mechanism would need to be
78//! skipped on operator startup, and would cause unnecessary read
79//! load on persist.
80//!
81//! Additionally, persisting sources requires we use bounded
82//! amounts of memory, even if a single timestamp represents
83//! a huge amount of data. This is not (currently) possible
84//! to guarantee while also performing self-correction.
85//!
86//! Because of this, we have ripped out the self-correction
87//! mechanism, and aggressively simplified the sub-operators.
88//! Some, particularly `append_batches` could be merged with
89//! the compute version, but that requires some amount of
90//! onerous refactoring that we have chosen to skip for now.
91//!
92// TODO(guswynn): merge at least the `append_batches` operator`
93
94use std::cmp::Ordering;
95use std::collections::{BTreeMap, VecDeque};
96use std::fmt::Debug;
97use std::ops::AddAssign;
98use std::rc::Rc;
99use std::sync::Arc;
100use std::time::Duration;
101
102use differential_dataflow::difference::Monoid;
103use differential_dataflow::lattice::Lattice;
104use differential_dataflow::{AsCollection, Hashable, VecCollection};
105use futures::{StreamExt, future};
106use itertools::Itertools;
107use mz_ore::cast::CastFrom;
108use mz_ore::collections::HashMap;
109use mz_persist_client::Diagnostics;
110use mz_persist_client::batch::{Batch, BatchBuilder, ProtoBatch};
111use mz_persist_client::cache::PersistClientCache;
112use mz_persist_client::error::UpperMismatch;
113use mz_persist_types::codec_impls::UnitSchema;
114use mz_persist_types::{Codec, Codec64};
115use mz_repr::{Diff, GlobalId, Row};
116use mz_storage_types::controller::CollectionMetadata;
117use mz_storage_types::errors::DataflowError;
118use mz_storage_types::sources::SourceData;
119use mz_storage_types::{StorageDiff, dyncfgs};
120use mz_timely_util::builder_async::{
121    Event, OperatorBuilder as AsyncOperatorBuilder, PressOnDropButton,
122};
123use serde::{Deserialize, Serialize};
124use timely::PartialOrder;
125use timely::container::CapacityContainerBuilder;
126use timely::dataflow::channels::pact::{Exchange, Pipeline};
127use timely::dataflow::operators::vec::Broadcast;
128use timely::dataflow::operators::{Capability, CapabilitySet, InspectCore};
129use timely::dataflow::{Scope, Stream, StreamVec};
130use timely::progress::{Antichain, Timestamp};
131use tokio::sync::Semaphore;
132use tracing::trace;
133
134use crate::metrics::source::SourcePersistSinkMetrics;
135use crate::statistics::SourceStatistics;
136use crate::storage_state::StorageState;
137
138/// Metrics about batches.
139#[derive(Clone, Debug, Default, Deserialize, Serialize)]
140struct BatchMetrics {
141    inserts: u64,
142    retractions: u64,
143    error_inserts: u64,
144    error_retractions: u64,
145}
146
147impl AddAssign<&BatchMetrics> for BatchMetrics {
148    fn add_assign(&mut self, rhs: &BatchMetrics) {
149        let BatchMetrics {
150            inserts: self_inserts,
151            retractions: self_retractions,
152            error_inserts: self_error_inserts,
153            error_retractions: self_error_retractions,
154        } = self;
155        let BatchMetrics {
156            inserts: rhs_inserts,
157            retractions: rhs_retractions,
158            error_inserts: rhs_error_inserts,
159            error_retractions: rhs_error_retractions,
160        } = rhs;
161        *self_inserts += rhs_inserts;
162        *self_retractions += rhs_retractions;
163        *self_error_inserts += rhs_error_inserts;
164        *self_error_retractions += rhs_error_retractions;
165    }
166}
167
168/// Manages batches and metrics.
169struct BatchBuilderAndMetadata<K, V, T, D>
170where
171    K: Codec,
172    V: Codec,
173    T: Timestamp + Lattice + Codec64,
174{
175    builder: BatchBuilder<K, V, T, D>,
176    /// Largest update timestamp staged so far, `None` while empty.
177    ///
178    /// `append_batches` needs this to decide, after an `UpperMismatch`, whether a batch lies
179    /// entirely below a raised append lower. A batch completely below the append lower is
180    /// deleted. A batch whose data straddles an append lower has its bounds adjusted instead.
181    data_max_ts: Option<T>,
182    metrics: BatchMetrics,
183}
184
185impl<K, V, T, D> BatchBuilderAndMetadata<K, V, T, D>
186where
187    K: Codec + Debug,
188    V: Codec + Debug,
189    T: Timestamp + Lattice + Codec64,
190    D: Monoid + Codec64,
191{
192    /// Creates a new batch. Updates at any timestamp at or beyond the builder's lower may be
193    /// added, in any order.
194    fn new(builder: BatchBuilder<K, V, T, D>) -> Self {
195        BatchBuilderAndMetadata {
196            builder,
197            data_max_ts: None,
198            metrics: Default::default(),
199        }
200    }
201
202    /// Adds an update to the batch.
203    async fn add(&mut self, k: &K, v: &V, t: &T, d: &D) {
204        self.data_max_ts = Some(match self.data_max_ts.take() {
205            Some(max) => max.join(t),
206            None => t.clone(),
207        });
208
209        self.builder.add(k, v, t, d).await.expect("invalid usage");
210    }
211
212    /// Finishes the batch, registering it under `lower` and `upper`.
213    ///
214    /// Panics if no update was ever added, since an empty batch has no largest timestamp. Callers
215    /// open a builder on the first update rather than up front, so reaching this is a bug.
216    async fn finish(self, lower: Antichain<T>, upper: Antichain<T>) -> HollowBatchAndMetadata<T> {
217        let data_max_ts = self.data_max_ts.expect("finishing an empty builder");
218        // `BatchBuilder::finish` rejects an update at or beyond `upper`, so a builder that was
219        // handed updates outside the description it is being finished under fails here rather
220        // than producing a batch whose parts reach past their registered bounds.
221        let batch = self
222            .builder
223            .finish(upper.clone())
224            .await
225            .expect("invalid usage");
226        HollowBatchAndMetadata {
227            lower,
228            upper,
229            data_max_ts,
230            batch: batch.into_transmittable_batch(),
231            metrics: self.metrics,
232        }
233    }
234}
235
236/// A batch or data + metrics moved from `write_batches` to `append_batches`.
237#[derive(Clone, Debug, Deserialize, Serialize)]
238#[serde(bound(
239    serialize = "T: Timestamp + Codec64",
240    deserialize = "T: Timestamp + Codec64"
241))]
242struct HollowBatchAndMetadata<T> {
243    lower: Antichain<T>,
244    upper: Antichain<T>,
245    data_max_ts: T,
246    batch: ProtoBatch,
247    metrics: BatchMetrics,
248}
249
250/// Holds finished batches for `append_batches`.
251#[derive(Debug, Default)]
252struct BatchSet {
253    finished: Vec<FinishedBatch>,
254    batch_metrics: BatchMetrics,
255}
256
257#[derive(Debug)]
258struct FinishedBatch {
259    batch: Batch<SourceData, (), mz_repr::Timestamp, StorageDiff>,
260    data_max_ts: mz_repr::Timestamp,
261}
262
263/// The batch builder the source sink writes with.
264type SourceBatchBuilder = BatchBuilderAndMetadata<SourceData, (), mz_repr::Timestamp, StorageDiff>;
265
266/// How far past the remap upper the minter commits its ceiling while a snapshot pins the
267/// frontier, in milliseconds, the unit of `mz_repr::Timestamp`. `None` turns committing ahead off.
268///
269/// `timestamp_interval` is the floor to avoid having data outrun the ceiling. Remap emits a new
270/// binding and downgrades its capability together. Rows emitted by reclock under the new binding
271/// are racing the new ceiling to the batch writers. At one timestamp interval the previous ceiling
272/// already covers the current timestamp's rows. Rows still outrun it when consecutive probes are
273/// further apart than the lookahead, which costs a batch rather than correctness.
274fn description_lookahead(lookahead: Duration, timestamp_interval: Duration) -> Option<u64> {
275    // Saturates rather than panics, since the lookahead is operator-supplied. A lookahead that
276    // large overflows the ceiling in `next_mint`, which then commits nothing.
277    let to_millis = |d: Duration| u64::try_from(d.as_millis()).unwrap_or(u64::MAX);
278    let lookahead = to_millis(lookahead);
279    // The floor is applied past the zero test, which is the only value that turns committing
280    // ahead off.
281    (lookahead > 0).then(|| lookahead.max(to_millis(timestamp_interval)))
282}
283
284/// Adds one update to `builder`, keeping the batch metrics in step.
285async fn stage_update(
286    builder: &mut SourceBatchBuilder,
287    row: Result<Row, DataflowError>,
288    ts: mz_repr::Timestamp,
289    diff: Diff,
290) {
291    let is_value = row.is_ok();
292
293    builder
294        .add(&SourceData(row), &(), &ts, &diff.into_inner())
295        .await;
296
297    // Note that we assume `diff` is either +1 or -1 here, being anything else is a logic bug we
298    // can't handle at the metric layer. We also assume this addition doesn't overflow.
299    match (is_value, diff.is_positive()) {
300        (true, true) => builder.metrics.inserts += diff.unsigned_abs(),
301        (true, false) => builder.metrics.retractions += diff.unsigned_abs(),
302        (false, true) => builder.metrics.error_inserts += diff.unsigned_abs(),
303        (false, false) => builder.metrics.error_retractions += diff.unsigned_abs(),
304    }
305}
306
307/// Continuously writes the `desired_stream` into persist
308/// This is done via a multi-stage operator graph:
309///
310/// 1. `mint_batch_descriptions` emits new batch descriptions whenever the
311///    frontier of `desired_collection` advances. A batch description is
312///    a pair of `(lower, upper)` that tells write operators
313///    which updates to write and in the end tells the append operator
314///    what frontiers to use when calling `append`/`compare_and_append`.
315///    This is a single-worker operator.
316/// 2. `write_batches` writes the `desired_collection` to persist as
317///    batches and sends those batches along.
318///    This does not yet append the batches to the persist shard, the update are
319///    only uploaded/prepared to be appended to a shard. Also: we only write
320///    updates for batch descriptions that we learned about from
321///    `mint_batch_descriptions`.
322/// 3. `append_batches` takes as input the minted batch descriptions and written
323///    batches. Whenever the frontiers sufficiently advance, we take a batch
324///    description and all the batches that belong to it and append it to the
325///    persist shard.
326///
327/// This operator assumes that the `desired_collection` comes pre-sharded.
328///
329/// Note that `mint_batch_descriptions` inspects the frontier of
330/// `desired_collection`, and passes the data through to `write_batches`.
331/// This is done to avoid a clone of the underlying data so that both
332/// operators can have the collection as input.
333///
334/// `snapshot_time` is the as_of, the time the collection's snapshot lands at, given only for an
335/// export that snapshots in this incarnation. The frontier sits there for the length of the
336/// snapshot, which is what gives the sink something to group. `remap_upper` carries the
337/// ingestion's remap upper as its frontier, which paces the grouping. See
338/// [`mint_batch_descriptions`].
339pub(crate) fn render<'scope>(
340    scope: Scope<'scope, mz_repr::Timestamp>,
341    collection_id: GlobalId,
342    target: CollectionMetadata,
343    desired_collection: VecCollection<'scope, mz_repr::Timestamp, Result<Row, DataflowError>, Diff>,
344    storage_state: &StorageState,
345    metrics: SourcePersistSinkMetrics,
346    busy_signal: Arc<Semaphore>,
347    snapshot_time: Option<Antichain<mz_repr::Timestamp>>,
348    timestamp_interval: Duration,
349    remap_upper: StreamVec<'scope, mz_repr::Timestamp, ()>,
350) -> (
351    StreamVec<'scope, mz_repr::Timestamp, ()>,
352    StreamVec<'scope, mz_repr::Timestamp, Rc<anyhow::Error>>,
353    Vec<PressOnDropButton>,
354) {
355    let persist_clients = Arc::clone(&storage_state.persist_clients);
356
357    let operator_name = format!("persist_sink({})", collection_id);
358
359    let config_set = storage_state.storage_configuration.config_set();
360    let lookahead = description_lookahead(
361        dyncfgs::STORAGE_PERSIST_SINK_DESCRIPTION_LOOKAHEAD.get(config_set),
362        timestamp_interval,
363    );
364
365    let (batch_descriptions, commitments, passthrough_desired_stream, mint_token) =
366        mint_batch_descriptions(
367            scope,
368            collection_id,
369            &operator_name,
370            &target,
371            desired_collection,
372            remap_upper,
373            Arc::clone(&persist_clients),
374            lookahead,
375            snapshot_time,
376        );
377
378    let source_statistics = storage_state
379        .aggregated_statistics
380        .get_source(&collection_id)
381        .expect("statistics initialized")
382        .clone();
383
384    let (written_batches, write_token) = write_batches(
385        scope,
386        collection_id.clone(),
387        &operator_name,
388        &target,
389        batch_descriptions.clone(),
390        commitments,
391        passthrough_desired_stream.as_collection(),
392        Arc::clone(&persist_clients),
393        source_statistics,
394        Arc::clone(&busy_signal),
395    );
396
397    let (upper_stream, append_errors, append_token) = append_batches(
398        scope,
399        collection_id.clone(),
400        operator_name,
401        &target,
402        batch_descriptions,
403        written_batches,
404        persist_clients,
405        storage_state,
406        metrics,
407        Arc::clone(&busy_signal),
408    );
409
410    (
411        upper_stream,
412        append_errors,
413        vec![mint_token, write_token, append_token],
414    )
415}
416
417/// Whenever the frontier advances, this mints a new batch description (lower
418/// and upper) that writers should use for writing the next set of batches to
419/// persist.
420///
421/// With a `lookahead`, and while the frontier sits at `snapshot_time`, it also commits to a
422/// ceiling that far past the data and broadcasts it on the second output. A
423/// ceiling is not a description: it gives the writers a bound to group updates under before the
424/// frontier certifies anything, and it binds this operator, which mints nothing below an
425/// outstanding ceiling. So the whole snapshot and the catch-up behind it become one description,
426/// emitted when the frontier reaches the ceiling. See [`next_mint`].
427///
428/// Only one of the workers does this, meaning there will only be one
429/// description in the stream, even in case of multiple timely workers. Use
430/// `broadcast()` to, ahem, broadcast, the one description to all downstream
431/// write operators/workers.
432///
433/// `remap_upper` carries the ingestion's remap upper as its frontier, which paces the
434/// commitments. Reclocking stamps every update below that upper before the update exists, so a
435/// ceiling ahead of it leads every row that can still arrive, however long this export's own data
436/// has been quiet.
437fn mint_batch_descriptions<'scope>(
438    scope: Scope<'scope, mz_repr::Timestamp>,
439    collection_id: GlobalId,
440    operator_name: &str,
441    target: &CollectionMetadata,
442    desired_collection: VecCollection<'scope, mz_repr::Timestamp, Result<Row, DataflowError>, Diff>,
443    remap_upper: StreamVec<'scope, mz_repr::Timestamp, ()>,
444    persist_clients: Arc<PersistClientCache>,
445    lookahead: Option<u64>,
446    snapshot_time: Option<Antichain<mz_repr::Timestamp>>,
447) -> (
448    StreamVec<
449        'scope,
450        mz_repr::Timestamp,
451        (Antichain<mz_repr::Timestamp>, Antichain<mz_repr::Timestamp>),
452    >,
453    StreamVec<'scope, mz_repr::Timestamp, Commitment>,
454    StreamVec<'scope, mz_repr::Timestamp, (Result<Row, DataflowError>, mz_repr::Timestamp, Diff)>,
455    PressOnDropButton,
456) {
457    let persist_location = target.persist_location.clone();
458    let shard_id = target.data_shard;
459    let target_relation_desc = target.relation_desc.clone();
460
461    // Only one worker is responsible for determining batch descriptions. All
462    // workers must write batches with the same description, to ensure that they
463    // can be combined into one batch that gets appended to Consensus state.
464    let hashed_id = collection_id.hashed();
465    let active_worker = usize::cast_from(hashed_id) % scope.peers() == scope.index();
466
467    // Only the "active" operator will mint batches. All other workers have an
468    // empty frontier. It's necessary to insert all of these into
469    // `compute_state.sink_write_frontier` below so we properly clear out
470    // default frontiers of non-active workers.
471
472    let mut mint_op = AsyncOperatorBuilder::new(
473        format!("{} mint_batch_descriptions", operator_name),
474        scope.clone(),
475    );
476
477    let (output, output_stream) = mint_op.new_output::<CapacityContainerBuilder<Vec<_>>>();
478    let (ceiling_output, ceiling_output_stream) =
479        mint_op.new_output::<CapacityContainerBuilder<Vec<_>>>();
480    let (data_output, data_output_stream) =
481        mint_op.new_output::<CapacityContainerBuilder<Vec<_>>>();
482
483    // The description, ceiling and data-passthrough outputs are all driven by this input, so
484    // they use a standard input connection.
485    let mut desired_input = mint_op.new_input_for_many(
486        desired_collection.inner,
487        Pipeline,
488        [&output, &ceiling_output, &data_output],
489    );
490
491    // The remap upper only influences the bounds of descriptions, it doesn't drive the outputs, so
492    // this input is disconnected. The stream is already broadcast, so every worker holds its
493    // frontier and the active one reads it in place.
494    let mut remap_input = mint_op.new_disconnected_input(remap_upper, Pipeline);
495
496    let shutdown_button = mint_op.build(move |capabilities| async move {
497        // Non-active workers should just pass the data through.
498        if !active_worker {
499            // The description and ceiling outputs are entirely driven by the active worker, so we
500            // drop their capabilities here. The data-passthrough output just uses the data
501            // capabilities.
502            drop(capabilities);
503            while let Some(event) = desired_input.next().await {
504                match event {
505                    Event::Data([_output_cap, _ceiling_cap, data_output_cap], mut data) => {
506                        data_output.give_container(&data_output_cap, &mut data);
507                    }
508                    Event::Progress(_) => {}
509                }
510            }
511            return;
512        }
513        // The data-passthrough output uses the data capabilities, so we drop its capability here.
514        let [desc_cap, ceiling_cap, _]: [_; 3] =
515            capabilities.try_into().expect("one capability per output");
516        let mut cap_set = CapabilitySet::from_elem(desc_cap);
517        let mut ceiling_cap_set = CapabilitySet::from_elem(ceiling_cap);
518
519        // Initialize this operators's `upper` to the `upper` of the persist shard we are writing
520        // to. Data from the source not beyond this time will be dropped, as it has already
521        // been persisted.
522        // In the future, sources will avoid passing through data not beyond this upper
523        let mut current_upper = {
524            // TODO(aljoscha): We need to figure out what to do with error
525            // results from these calls.
526            let persist_client = persist_clients
527                .open(persist_location)
528                .await
529                .expect("could not open persist client");
530
531            let mut write = persist_client
532                .open_writer::<SourceData, (), mz_repr::Timestamp, StorageDiff>(
533                    shard_id,
534                    Arc::new(target_relation_desc),
535                    Arc::new(UnitSchema),
536                    Diagnostics {
537                        shard_name: collection_id.to_string(),
538                        handle_purpose: format!(
539                            "storage::persist_sink::mint_batch_descriptions {}",
540                            collection_id
541                        ),
542                    },
543                )
544                .await
545                .expect("could not open persist shard");
546
547            // TODO: this sink currently cannot tolerate a stale upper... which is bad because the
548            // upper can become stale as soon as it is read. (For example, if another concurrent
549            // instance of the sink has updated it.) Fetching a recent upper helps to mitigate this,
550            // but ideally we would just skip ahead if we discover that our upper is stale.
551            let upper = write.fetch_recent_upper().await.clone();
552            // explicitly expire the once-used write handle.
553            write.expire().await;
554            upper
555        };
556
557        // The current input frontier.
558        let mut desired_frontier = Antichain::from_elem(mz_repr::Timestamp::minimum());
559
560        // The ingestion's remap upper, which drives the ceiling.
561        let mut remap_upper = Antichain::from_elem(mz_repr::Timestamp::minimum());
562
563        // The outstanding ceiling, if any. While one is held nothing below it is minted, which is
564        // what makes it binding.
565        let mut committed: Option<mz_repr::Timestamp> = None;
566
567        // The wait sits at the bottom of the loop, so a mint pass runs once before the first
568        // event as well as after every event. A fresh export's snapshot lands at the minimum and
569        // its frontier never moves during it, so no event announces the pin. Its first ceiling
570        // has to go out ahead of the snapshot's rows rather than alongside them.
571        loop {
572            // `desired_frontier` starts at the minimum. A fresh export's as_of is the minimum too,
573            // so it reads as snapshotting from the first activation. A later as_of waits for the
574            // progress statement that moves the frontier there. A row can arrive ahead of that
575            // statement, and committing on it would anchor the ceiling at the shard upper instead
576            // of the snapshot's time.
577            let snapshot_in_progress = snapshot_time
578                .as_ref()
579                .is_some_and(|time| *time == desired_frontier);
580
581            while let Some(mint) = next_mint(
582                &current_upper,
583                &desired_frontier,
584                &remap_upper,
585                committed,
586                lookahead.filter(|_| snapshot_in_progress),
587            ) {
588                let lower = current_upper
589                    .as_option()
590                    .copied()
591                    .expect("a non-empty current upper, or nothing is minted");
592
593                let upper = match mint {
594                    Mint::Ceiling(ceiling) => {
595                        let cap = ceiling_cap_set
596                            .try_delayed(&lower)
597                            .expect("ceiling capability holds the current upper");
598                        trace!(
599                            "persist_sink {collection_id}/{shard_id}: \
600                                committing ceiling: {:?}",
601                            ceiling
602                        );
603                        ceiling_output.give(&cap, Commitment { lower, ceiling });
604                        committed = Some(ceiling);
605                        continue;
606                    }
607                    Mint::Description(upper) => upper,
608                };
609
610                let batch_description = (current_upper.to_owned(), upper.to_owned());
611
612                let cap = cap_set
613                    .try_delayed(&lower)
614                    .ok_or_else(|| {
615                        format!(
616                            "minter cannot delay {:?} to {:?}. \
617                                Likely because we already emitted a \
618                                batch description and delayed.",
619                            cap_set, lower
620                        )
621                    })
622                    .unwrap();
623
624                trace!(
625                    "persist_sink {collection_id}/{shard_id}: \
626                        new batch_description: {:?}",
627                    batch_description
628                );
629
630                output.give(&cap, batch_description);
631
632                // We downgrade our capability to the batch
633                // description upper, as there will never be
634                // any overlapping descriptions.
635                trace!(
636                    "persist_sink {collection_id}/{shard_id}: \
637                        downgrading to {:?}",
638                    upper
639                );
640                cap_set.downgrade(upper.iter());
641                ceiling_cap_set.downgrade(upper.iter());
642
643                // The description that retires a ceiling covers everything the writers grouped
644                // under it, so the next one starts fresh.
645                committed = None;
646                current_upper = upper;
647            }
648
649            tokio::select! {
650                event = desired_input.next() => match event {
651                    Some(Event::Data([_output_cap, _ceiling_cap, data_output_cap], mut data)) => {
652                        // Just passthrough the data.
653                        data_output.give_container(&data_output_cap, &mut data);
654                    }
655                    Some(Event::Progress(frontier)) => desired_frontier = frontier,
656                    // Input is exhausted, so we can shut down.
657                    None => return,
658                },
659                // During the snapshot, the frontier is pinned and the remap upper drives the
660                // ceiling.
661                _ = remap_input.ready() => {},
662            }
663
664            // The stream carries no data, only its frontier.
665            while let Some(event) = remap_input.next_sync() {
666                if let Event::Progress(frontier) = event {
667                    remap_upper = frontier;
668                }
669            }
670        }
671    });
672
673    (
674        output_stream,
675        ceiling_output_stream,
676        data_output_stream,
677        shutdown_button.press_on_drop(),
678    )
679}
680
681/// A bound the minter commits ahead of the frontier, broadcast to the writers.
682///
683/// `lower` is the lower of the description that will end at or past `ceiling`, which the writers
684/// need before that description exists, since a builder declares its lower when it opens.
685#[derive(Clone, Copy, Debug, PartialEq, Eq, Deserialize, Serialize)]
686struct Commitment {
687    lower: mz_repr::Timestamp,
688    ceiling: mz_repr::Timestamp,
689}
690
691/// What [`mint_batch_descriptions`] does with the frontier and the data it has seen.
692#[derive(Debug, Clone, PartialEq, Eq)]
693enum Mint {
694    /// Emit a description ending here, the unit `append_batches` appends.
695    Description(Antichain<mz_repr::Timestamp>),
696    /// Commit to this ceiling and broadcast it to the writers. Not a description: it gives them a
697    /// bound to group under, and binds the minter to mint nothing below it.
698    Ceiling(mz_repr::Timestamp),
699}
700
701/// The next thing the minter emits, or `None` when there is nothing to do.
702///
703/// A description is derived from the frontier, which certifies that everything below it has
704/// arrived. An outstanding `committed` ceiling suppresses that until the frontier reaches the
705/// ceiling, which is what collapses a whole snapshot and the catch-up behind it into one
706/// description.
707///
708/// * `current_upper`:
709///     The upper of the last minted description, initialized from the shard.
710///     Any data before this has already been processed.
711/// * `desired_frontier`:
712///     Data input frontier.
713///     This frontier advances when the reclocked data stream advances.
714/// * `remap_upper`:
715///     The frontier of the remap shard, which advances on each successful probe.
716/// * `committed` (Optional):
717///     This is the current value of the ceiling. If set, the upper of a minted description will be
718///     at or beyond this.
719/// * `lookahead` (Optional):
720///     If present, `desired_frontier` is pinned to `as_of` (e.g. due to snapshot), so we advance
721///     `committed` as `remap_upper` + `lookahead` to consolidate in-flight data into one shared
722///     batch.
723fn next_mint(
724    current_upper: &Antichain<mz_repr::Timestamp>,
725    desired_frontier: &Antichain<mz_repr::Timestamp>,
726    remap_upper: &Antichain<mz_repr::Timestamp>,
727    committed: Option<mz_repr::Timestamp>,
728    lookahead: Option<u64>,
729) -> Option<Mint> {
730    let frontier_reached = |ts: mz_repr::Timestamp| {
731        PartialOrder::less_equal(&Antichain::from_elem(ts), desired_frontier)
732    };
733    if committed.is_none_or(frontier_reached)
734        && PartialOrder::less_than(current_upper, desired_frontier)
735    {
736        return Some(Mint::Description(desired_frontier.clone()));
737    }
738
739    // A bunch of checks that will return None (meaning there's no ceiling).
740    // If feature is not enabled or not snapshotting...
741    let lookahead = lookahead?;
742
743    // If the data shard has advanced to empty frontier...
744    let lower = *current_upper.as_option()?;
745
746    // If the remap shard has advanced to the empty frontier, or the add overflows because of
747    // a large lookahead value (see `description_lookahead`).
748    let ceiling = remap_upper.as_option()?.checked_add(lookahead)?;
749
750    // lower can be beyond remap_upper + lookahead any time the current_upper jumps ahead of
751    // of the remap_upper the operator is tracking.
752    (lower < ceiling && committed.is_none_or(|c| c < ceiling)).then_some(Mint::Ceiling(ceiling))
753}
754
755/// Writes `desired_collection` to persist, but only for updates
756/// that fall into batch a description that we get via `batch_descriptions`.
757/// This forwards a `HollowBatch` (with additional metadata)
758/// for any batch of updates that was written.
759///
760/// Every update below the ceiling `mint_batch_descriptions` commits goes into one open builder,
761/// whatever its timestamp, so a pinned frontier costs one batch rather than one per timestamp. An
762/// update that outruns the ceiling, or arrives while none is committed, has no bound to be grouped
763/// under and writes a batch of its own timestamp, which is what the sink does for every update when
764/// nothing is ever committed ahead of the frontier.
765///
766/// This operator assumes that the `desired_collection` comes pre-sharded.
767///
768/// This also and updates various metrics.
769fn write_batches<'scope>(
770    scope: Scope<'scope, mz_repr::Timestamp>,
771    collection_id: GlobalId,
772    operator_name: &str,
773    target: &CollectionMetadata,
774    batch_descriptions: Stream<
775        'scope,
776        mz_repr::Timestamp,
777        Vec<(Antichain<mz_repr::Timestamp>, Antichain<mz_repr::Timestamp>)>,
778    >,
779    commitments: StreamVec<'scope, mz_repr::Timestamp, Commitment>,
780    desired_collection: VecCollection<'scope, mz_repr::Timestamp, Result<Row, DataflowError>, Diff>,
781    persist_clients: Arc<PersistClientCache>,
782    source_statistics: SourceStatistics,
783    busy_signal: Arc<Semaphore>,
784) -> (
785    StreamVec<'scope, mz_repr::Timestamp, HollowBatchAndMetadata<mz_repr::Timestamp>>,
786    PressOnDropButton,
787) {
788    let worker_index = scope.index();
789
790    let persist_location = target.persist_location.clone();
791    let shard_id = target.data_shard;
792    let target_relation_desc = target.relation_desc.clone();
793
794    let mut write_op =
795        AsyncOperatorBuilder::new(format!("{} write_batches", operator_name), scope.clone());
796
797    let (output, output_stream) = write_op.new_output::<CapacityContainerBuilder<Vec<_>>>();
798
799    let mut descriptions_input =
800        write_op.new_input_for(batch_descriptions.broadcast(), Pipeline, &output);
801    // Commitments only route updates into builders, so this input is disconnected: holding the
802    // output back on it would stall every batch behind the minter's own progress.
803    let mut commitments_input = write_op.new_disconnected_input(commitments.broadcast(), Pipeline);
804    let mut desired_input = write_op.new_disconnected_input(desired_collection.inner, Pipeline);
805
806    // This operator accepts the current and desired update streams for a `persist` shard.
807    // It attempts to write out updates, starting from the current's upper frontier, that
808    // will cause the changes of desired to be committed to persist, _but only those also past the
809    // upper_.
810
811    let shutdown_button = write_op.build(move |_capabilities| async move {
812        // Builders for timestamps no ceiling covers, keyed by timestamp.
813        //
814        // A batch builder cannot be split, so an update may only join updates at other timestamps
815        // once a bound they all fall below is known. Until then a timestamp gets a builder to
816        // itself, which is safe without knowing the descriptions because a description covers a
817        // timestamp entirely or not at all. It is finished only once its description is ready,
818        // since rows at its timestamp can still be staged after that description is processed.
819        let mut uncovered_builders: BTreeMap<mz_repr::Timestamp, SourceBatchBuilder> =
820            BTreeMap::new();
821
822        // The outstanding commitment: the bounds the open builder takes updates in, and the lower
823        // it declares. The ceiling is raised to the description's own upper once that arrives, so
824        // rows landing between arrival and readiness still join. `None` until the first commitment,
825        // which leaves every timestamp writing its own batch.
826        let mut commitment: Option<Commitment> = None;
827
828        // The one open builder, holding every update inside the commitment regardless of
829        // timestamp. It reaches `persist_blob_target_size` and spills its parts to blob, so what it
830        // holds resident is one unflushed part however long the frontier stays pinned. Only ever
831        // one, because the minter mints nothing below an outstanding ceiling, so there is exactly
832        // one description in flight for it to be finished under.
833        let mut open_builder: Option<SourceBatchBuilder> = None;
834
835        // Contains descriptions of batches for which we know that we can
836        // write data. We got these from the "centralized" operator that
837        // determines batch descriptions for all writers.
838        //
839        // `Antichain` does not implement `Ord`, so we cannot use a `BTreeMap`. We need to search
840        // through the map, so we cannot use the `mz_ore` wrapper either.
841        #[allow(clippy::disallowed_types)]
842        let mut in_flight_batches = std::collections::HashMap::<
843            (Antichain<mz_repr::Timestamp>, Antichain<mz_repr::Timestamp>),
844            Capability<mz_repr::Timestamp>,
845        >::new();
846
847        // TODO(aljoscha): We need to figure out what to do with error results from these calls.
848        let persist_client = persist_clients
849            .open(persist_location)
850            .await
851            .expect("could not open persist client");
852
853        let write = persist_client
854            .open_writer::<SourceData, (), mz_repr::Timestamp, StorageDiff>(
855                shard_id,
856                Arc::new(target_relation_desc),
857                Arc::new(UnitSchema),
858                Diagnostics {
859                    shard_name: collection_id.to_string(),
860                    handle_purpose: format!(
861                        "storage::persist_sink::write_batches {}",
862                        collection_id
863                    ),
864                },
865            )
866            .await
867            .expect("could not open persist shard");
868
869        // The current input frontiers.
870        let mut batch_descriptions_frontier = Antichain::from_elem(Timestamp::minimum());
871        let mut desired_frontier = Antichain::from_elem(Timestamp::minimum());
872
873        // The frontiers of the inputs we have processed, used to avoid redoing work
874        let mut processed_desired_frontier = Antichain::from_elem(Timestamp::minimum());
875        let mut processed_descriptions_frontier = Antichain::from_elem(Timestamp::minimum());
876
877        // A "safe" choice for the lower of new batches we are creating.
878        let mut operator_batch_lower = Antichain::from_elem(Timestamp::minimum());
879
880        while !(batch_descriptions_frontier.is_empty() && desired_frontier.is_empty()) {
881            // Wait for either inputs to become ready
882            tokio::select! {
883                _ = descriptions_input.ready() => {},
884                _ = commitments_input.ready() => {},
885                _ = desired_input.ready() => {},
886            }
887
888            // Collect ready work from all three inputs. Commitments and descriptions are processed
889            // before the data of the same round so that an update whose bound arrived alongside it
890            // can go straight into the open builder instead of writing a batch of its own.
891            let ready_commitments =
892                std::iter::from_fn(|| commitments_input.next_sync()).collect_vec();
893            let ready_descriptions =
894                std::iter::from_fn(|| descriptions_input.next_sync()).collect_vec();
895            let ready_events = std::iter::from_fn(|| desired_input.next_sync()).collect_vec();
896
897            // We now start the async work for the input we received. Until we finish the dataflow
898            // should be marked as busy.
899            let permit = busy_signal.acquire().await;
900
901            for event in ready_commitments {
902                let Event::Data(_cap, data) = event else {
903                    continue;
904                };
905                for next in data {
906                    // The minter only commits while the source is snapshotting, which pins the
907                    // frontier. The description that retires the commitment moves the frontier for
908                    // good. Every commitment during that snapshot shares one lower.
909                    let held = commitment.get_or_insert(next);
910                    assert_eq!(
911                        held.lower,
912                        next.lower,
913                        "persist_sink {collection_id}/{shard_id}: commitment {next:?} while {held:?} is outstanding"
914                    );
915                    held.ceiling = held.ceiling.max(next.ceiling);
916                }
917            }
918
919            for event in ready_descriptions {
920                match event {
921                    Event::Data(cap, data) => {
922                        // Ingest new batch descriptions.
923                        for description in data {
924                            if collection_id.is_user() {
925                                trace!(
926                                    "persist_sink {collection_id}/{shard_id}: \
927                                        write_batches: \
928                                        new_description: {:?}, \
929                                        desired_frontier: {:?}, \
930                                        batch_descriptions_frontier: {:?}",
931                                    description, desired_frontier, batch_descriptions_frontier,
932                                );
933                            }
934
935                            let (lower, upper) = (&description.0, &description.1);
936                            let lower_ts = *lower
937                                .as_option()
938                                .expect("minted descriptions have a single-element lower");
939
940                            // The description that retires a commitment ends at or past the
941                            // ceiling, so rows landing between the description arriving and the
942                            // description becoming ready can still join the builder that will
943                            // finish under that description.
944                            if let Some(held) = commitment.filter(|held| held.lower == lower_ts)
945                                && let Some(upper_ts) = upper.as_option()
946                            {
947                                commitment = Some(Commitment {
948                                    ceiling: held.ceiling.max(*upper_ts),
949                                    ..held
950                                });
951                            }
952
953                            match in_flight_batches.entry(description) {
954                                std::collections::hash_map::Entry::Vacant(v) => {
955                                    // This _should_ be `.retain`, but rust
956                                    // currently thinks we can't use `cap`
957                                    // as an owned value when using the
958                                    // match guard `Some(event)`
959                                    v.insert(cap.delayed(cap.time()));
960                                }
961                                std::collections::hash_map::Entry::Occupied(o) => {
962                                    let (description, _) = o.remove_entry();
963                                    panic!(
964                                        "write_batches: sink {} got more than one \
965                                            batch for description {:?}, in-flight: {:?}",
966                                        collection_id, description, in_flight_batches
967                                    );
968                                }
969                            }
970                        }
971                    }
972                    Event::Progress(frontier) => {
973                        batch_descriptions_frontier = frontier;
974                    }
975                }
976            }
977
978            for event in ready_events {
979                match event {
980                    Event::Data(_cap, data) => {
981                        // Extract desired rows as positive contributions to `correction`.
982                        if collection_id.is_user() && !data.is_empty() {
983                            trace!(
984                                "persist_sink {collection_id}/{shard_id}: \
985                                    updates: {:?}, \
986                                    in-flight-batches: {:?}, \
987                                    desired_frontier: {:?}, \
988                                    batch_descriptions_frontier: {:?}",
989                                data,
990                                in_flight_batches,
991                                desired_frontier,
992                                batch_descriptions_frontier,
993                            );
994                        }
995
996                        for (row, ts, diff) in data {
997                            if write.upper().less_equal(&ts) {
998                                // Every description this operator has emitted was covered by the
999                                // desired frontier at the time, so no update below
1000                                // `operator_batch_lower` can still be in flight. An update that
1001                                // arrives anyway belongs to a description that is already gone: it
1002                                // matches no later description, so its batch would never be
1003                                // appended and the update would be lost unnoticed. Not a
1004                                // `debug_assert!`, which compiles out of the optimized and release
1005                                // profiles and would leave the loss silent everywhere it matters.
1006                                assert!(
1007                                    operator_batch_lower.less_equal(&ts),
1008                                    "persist_sink {collection_id}/{shard_id}: update at {ts:?} \
1009                                    arrived below the emitted batch lower {operator_batch_lower:?}",
1010                                );
1011
1012                                let inside =
1013                                    commitment.filter(|held| held.lower <= ts && ts < held.ceiling);
1014                                let builder = if let Some(held) = inside {
1015                                    // The description retiring the commitment is known to contain
1016                                    // it, so the update joins the one open builder whatever its
1017                                    // timestamp.
1018                                    open_builder.get_or_insert_with(|| {
1019                                        BatchBuilderAndMetadata::new(
1020                                            write.builder(Antichain::from_elem(held.lower)),
1021                                        )
1022                                    })
1023                                } else {
1024                                    // Nothing to group under, so the only lower this builder can
1025                                    // declare is the operator's own, the one lower at or below
1026                                    // every description that could come to cover it. That
1027                                    // declaration is what registers the batch truncated once it is
1028                                    // appended under a description's narrower bounds.
1029                                    uncovered_builders.entry(ts).or_insert_with(|| {
1030                                        BatchBuilderAndMetadata::new(
1031                                            write.builder(operator_batch_lower.clone()),
1032                                        )
1033                                    })
1034                                };
1035                                stage_update(builder, row, ts, diff).await;
1036                                source_statistics.inc_updates_staged_by(1);
1037                            }
1038                        }
1039                    }
1040                    Event::Progress(frontier) => {
1041                        desired_frontier = frontier;
1042                    }
1043                }
1044            }
1045
1046            // We may have the opportunity to commit updates, if either frontier
1047            // has moved
1048            if PartialOrder::less_equal(&processed_desired_frontier, &desired_frontier)
1049                || PartialOrder::less_equal(
1050                    &processed_descriptions_frontier,
1051                    &batch_descriptions_frontier,
1052                )
1053            {
1054                trace!(
1055                    "persist_sink {collection_id}/{shard_id}: \
1056                        CAN emit: \
1057                        processed_desired_frontier: {:?}, \
1058                        processed_descriptions_frontier: {:?}, \
1059                        desired_frontier: {:?}, \
1060                        batch_descriptions_frontier: {:?}",
1061                    processed_desired_frontier,
1062                    processed_descriptions_frontier,
1063                    desired_frontier,
1064                    batch_descriptions_frontier,
1065                );
1066
1067                trace!(
1068                    "persist_sink {collection_id}/{shard_id}: \
1069                        in-flight batches: {:?}, \
1070                        batch_descriptions_frontier: {:?}, \
1071                        desired_frontier: {:?}",
1072                    in_flight_batches, batch_descriptions_frontier, desired_frontier,
1073                );
1074
1075                // We can write updates for a given batch description when
1076                // a) the batch is not beyond `batch_descriptions_frontier`,
1077                // and b) we know that we have seen all updates that would
1078                // fall into the batch, from `desired_frontier`.
1079                let ready_batches = in_flight_batches
1080                    .keys()
1081                    .filter(|(lower, upper)| {
1082                        !PartialOrder::less_equal(&batch_descriptions_frontier, lower)
1083                            && !PartialOrder::less_than(&desired_frontier, upper)
1084                    })
1085                    .cloned()
1086                    .collect::<Vec<_>>();
1087
1088                trace!(
1089                    "persist_sink {collection_id}/{shard_id}: \
1090                        ready batches: {:?}",
1091                    ready_batches,
1092                );
1093
1094                for batch_description in ready_batches {
1095                    let cap = in_flight_batches.remove(&batch_description).unwrap();
1096
1097                    if collection_id.is_user() {
1098                        trace!(
1099                            "persist_sink {collection_id}/{shard_id}: \
1100                                emitting done batch: {:?}, cap: {:?}",
1101                            batch_description, cap
1102                        );
1103                    }
1104
1105                    let (batch_lower, batch_upper) = batch_description;
1106                    let lower = *batch_lower
1107                        .as_option()
1108                        .expect("minted descriptions have a single-element lower");
1109
1110                    let covered: Vec<_> = uncovered_builders
1111                        .keys()
1112                        .copied()
1113                        .filter(|ts| batch_lower.less_equal(ts) && !batch_upper.less_equal(ts))
1114                        .collect();
1115                    let mut batch_tokens = Vec::with_capacity(covered.len() + 1);
1116                    for ts in covered {
1117                        let builder = uncovered_builders.remove(&ts).expect("just looked up");
1118                        batch_tokens.push(
1119                            builder
1120                                .finish(batch_lower.clone(), batch_upper.clone())
1121                                .await,
1122                        );
1123                    }
1124
1125                    // If snapshotting, and the as_of = T, where T is greater than the minimum,
1126                    // the minter emits a description covering `(0, T)` and `Commitment{lower:T}`.
1127                    // The open builder must wait for the appropriate description.
1128                    if let Some(held) = commitment.filter(|held| held.lower == lower)
1129                        && let Some(builder) = open_builder.take()
1130                    {
1131                        assert!(
1132                            !batch_upper.less_than(&held.ceiling),
1133                            "persist_sink {collection_id}/{shard_id}: description upper {batch_upper:?} is below commitment ceiling {commitment:?}",
1134                        );
1135                        commitment = None;
1136                        if collection_id.is_user() {
1137                            trace!(
1138                                "persist_sink {collection_id}/{shard_id}: \
1139                                    wrote batch from worker {}: ({:?}, {:?}), containing {:?}",
1140                                worker_index, batch_lower, batch_upper, builder.metrics
1141                            );
1142                        }
1143
1144                        batch_tokens.push(
1145                            builder
1146                                .finish(batch_lower.clone(), batch_upper.clone())
1147                                .await,
1148                        );
1149                    }
1150
1151                    // The next "safe" lower for batches is the meet (max) of all the emitted
1152                    // batches. These uppers all are not beyond the `desired_frontier`, which
1153                    // means all updates received by this operator will be beyond this lower.
1154                    // Additionally, the `mint_batch_descriptions` operator ensures that
1155                    // later-received batch descriptions will start beyond these uppers as
1156                    // well.
1157                    //
1158                    // It is impossible to emit a batch description that is
1159                    // beyond a not-yet emitted description in `in_flight_batches`, as
1160                    // a that description would also have been chosen as ready above.
1161                    operator_batch_lower = operator_batch_lower.join(&batch_upper);
1162
1163                    output.give_container(&cap, &mut batch_tokens);
1164
1165                    processed_desired_frontier.clone_from(&desired_frontier);
1166                    processed_descriptions_frontier.clone_from(&batch_descriptions_frontier);
1167                }
1168            } else {
1169                trace!(
1170                    "persist_sink {collection_id}/{shard_id}: \
1171                        cannot emit: processed_desired_frontier: {:?}, \
1172                        processed_descriptions_frontier: {:?}, \
1173                        desired_frontier: {:?}",
1174                    processed_desired_frontier, processed_descriptions_frontier, desired_frontier
1175                );
1176            }
1177            drop(permit);
1178        }
1179    });
1180
1181    // Use `InspectCore::inspect_container` instead of `Inspect::inspect`.
1182    // `Inspect` carries a `where for<'a> &'a C: IntoIterator` bound, and on
1183    // macOS the solver can satisfy that bound by chasing objc2's
1184    // `&Retained<T>: IntoIterator` blanket impl into an endless
1185    // `Retained<Retained<…>>` chain, overflowing the recursion limit.
1186    // `InspectCore` has no such bound, so the cascade never starts. We
1187    // iterate the container by hand to recover the per-item callback.
1188    let output_stream = if collection_id.is_user() {
1189        InspectCore::inspect_container(output_stream, |event| {
1190            if let Ok((_, data)) = event {
1191                for d in data {
1192                    trace!("batch: {:?}", d);
1193                }
1194            }
1195        })
1196    } else {
1197        output_stream
1198    };
1199
1200    (output_stream, shutdown_button.press_on_drop())
1201}
1202
1203/// Fuses written batches together and appends them to persist using one
1204/// `compare_and_append` call. Writing only happens for batch descriptions where
1205/// we know that no future batches will arrive, that is, for those batch
1206/// descriptions that are not beyond the frontier of both the
1207/// `batch_descriptions` and `batches` inputs.
1208///
1209/// This also keeps the shared frontier that is stored in `compute_state` in
1210/// sync with the upper of the persist shard, and updates various metrics
1211/// and statistics objects.
1212fn append_batches<'scope>(
1213    scope: Scope<'scope, mz_repr::Timestamp>,
1214    collection_id: GlobalId,
1215    operator_name: String,
1216    target: &CollectionMetadata,
1217    batch_descriptions: Stream<
1218        'scope,
1219        mz_repr::Timestamp,
1220        Vec<(Antichain<mz_repr::Timestamp>, Antichain<mz_repr::Timestamp>)>,
1221    >,
1222    batches: StreamVec<'scope, mz_repr::Timestamp, HollowBatchAndMetadata<mz_repr::Timestamp>>,
1223    persist_clients: Arc<PersistClientCache>,
1224    storage_state: &StorageState,
1225    metrics: SourcePersistSinkMetrics,
1226    busy_signal: Arc<Semaphore>,
1227) -> (
1228    StreamVec<'scope, mz_repr::Timestamp, ()>,
1229    StreamVec<'scope, mz_repr::Timestamp, Rc<anyhow::Error>>,
1230    PressOnDropButton,
1231) {
1232    let persist_location = target.persist_location.clone();
1233    let shard_id = target.data_shard;
1234    let target_relation_desc = target.relation_desc.clone();
1235
1236    // We can only be lenient with concurrent modifications when we know that
1237    // this source pipeline is using the feedback upsert operator, which works
1238    // correctly when multiple instances of an ingestion pipeline produce
1239    // different updates, because of concurrency/non-determinism.
1240    let use_continual_feedback_upsert = dyncfgs::STORAGE_USE_CONTINUAL_FEEDBACK_UPSERT
1241        .get(storage_state.storage_configuration.config_set());
1242    let bail_on_concurrent_modification = !use_continual_feedback_upsert;
1243
1244    let mut read_only_rx = storage_state.read_only_rx.clone();
1245
1246    let operator_name = format!("{} append_batches", operator_name);
1247    let mut append_op = AsyncOperatorBuilder::new(operator_name, scope.clone());
1248
1249    let hashed_id = collection_id.hashed();
1250    let active_worker = usize::cast_from(hashed_id) % scope.peers() == scope.index();
1251    let worker_id = scope.index();
1252
1253    // Both of these inputs are disconnected from the output capabilities of this operator, as
1254    // any output of this operator is entirely driven by the `compare_and_append`s. Currently
1255    // this operator has no outputs, but they may be added in the future, when merging with
1256    // the compute `persist_sink`.
1257    let mut descriptions_input =
1258        append_op.new_disconnected_input(batch_descriptions, Exchange::new(move |_| hashed_id));
1259    let mut batches_input =
1260        append_op.new_disconnected_input(batches, Exchange::new(move |_| hashed_id));
1261
1262    let current_upper = Rc::clone(&storage_state.source_uppers[&collection_id]);
1263    if !active_worker {
1264        // This worker is not writing, so make sure it's "taken out" of the
1265        // calculation by advancing to the empty frontier.
1266        current_upper.borrow_mut().clear();
1267    }
1268
1269    let source_statistics = storage_state
1270        .aggregated_statistics
1271        .get_source(&collection_id)
1272        .expect("statistics initialized")
1273        .clone();
1274
1275    // An output whose frontier tracks the last successful compare and append of this operator
1276    let (_upper_output, upper_stream) = append_op.new_output::<CapacityContainerBuilder<Vec<_>>>();
1277
1278    // This operator accepts the batch descriptions and tokens that represent
1279    // written batches. Written batches get appended to persist when we learn
1280    // from our input frontiers that we have seen all batches for a given batch
1281    // description.
1282
1283    let (shutdown_button, errors) = append_op.build_fallible(move |caps| Box::pin(async move {
1284        let [upper_cap_set]: &mut [_; 1] = caps.try_into().unwrap();
1285
1286        // This may SEEM unnecessary, but metrics contains extra
1287        // `DeleteOnDrop`-wrapped fields that will NOT be moved into this
1288        // closure otherwise, dropping and destroying
1289        // those metrics. This is because rust now only moves the
1290        // explicitly-referenced fields into closures.
1291        let metrics = metrics;
1292
1293        // Contains descriptions of batches for which we know that we can
1294        // write data. We got these from the "centralized" operator that
1295        // determines batch descriptions for all writers.
1296        //
1297        // `Antichain` does not implement `Ord`, so we cannot use a `BTreeSet`. We need to search
1298        // through the set, so we cannot use the `mz_ore` wrapper either.
1299        #[allow(clippy::disallowed_types)]
1300        let mut in_flight_descriptions = std::collections::HashSet::<(
1301            Antichain<mz_repr::Timestamp>,
1302            Antichain<mz_repr::Timestamp>,
1303        )>::new();
1304
1305        // In flight batches that haven't been `compare_and_append`'d yet, plus metrics about
1306        // the batch.
1307        let mut in_flight_batches = HashMap::<
1308            (Antichain<mz_repr::Timestamp>, Antichain<mz_repr::Timestamp>),
1309            BatchSet,
1310        >::new();
1311
1312        source_statistics.initialize_rehydration_latency_ms();
1313        if !active_worker {
1314            // The non-active workers report that they are done snapshotting and hydrating.
1315            let empty_frontier = Antichain::new();
1316            source_statistics.initialize_snapshot_committed(&empty_frontier);
1317            source_statistics.update_rehydration_latency_ms(&empty_frontier);
1318            return Ok(());
1319        }
1320
1321        let persist_client = persist_clients
1322            .open(persist_location)
1323            .await?;
1324
1325        let mut write = persist_client
1326            .open_writer::<SourceData, (), mz_repr::Timestamp, StorageDiff>(
1327                shard_id,
1328                Arc::new(target_relation_desc),
1329                Arc::new(UnitSchema),
1330                Diagnostics {
1331                    shard_name:collection_id.to_string(),
1332                    handle_purpose: format!("persist_sink::append_batches {}", collection_id)
1333                },
1334            )
1335            .await?;
1336
1337        // Initialize this sink's `upper` to the `upper` of the persist shard we are writing
1338        // to. Data from the source not beyond this time will be dropped, as it has already
1339        // been persisted.
1340        // In the future, sources will avoid passing through data not beyond this upper
1341        // VERY IMPORTANT: Only the active write worker must change the
1342        // shared upper. All other workers have already cleared this
1343        // upper above.
1344        current_upper.borrow_mut().clone_from(write.upper());
1345        upper_cap_set.downgrade(current_upper.borrow().iter());
1346        source_statistics.initialize_snapshot_committed(write.upper());
1347
1348        // The current input frontiers.
1349        let mut batch_description_frontier = Antichain::from_elem(Timestamp::minimum());
1350        let mut batches_frontier = Antichain::from_elem(Timestamp::minimum());
1351
1352        loop {
1353            tokio::select! {
1354                Some(event) = descriptions_input.next() => {
1355                    match event {
1356                        Event::Data(_cap, data) => {
1357                            // Ingest new batch descriptions.
1358                            for batch_description in data {
1359                                if collection_id.is_user() {
1360                                    trace!(
1361                                        "persist_sink {collection_id}/{shard_id}: \
1362                                            append_batches: sink {}, \
1363                                            new description: {:?}, \
1364                                            batch_description_frontier: {:?}",
1365                                        collection_id,
1366                                        batch_description,
1367                                        batch_description_frontier
1368                                    );
1369                                }
1370
1371                                // This line has to be broken up, or
1372                                // rustfmt fails in the whole function :(
1373                                let is_new = in_flight_descriptions.insert(
1374                                    batch_description.clone()
1375                                );
1376
1377                                assert!(
1378                                    is_new,
1379                                    "append_batches: sink {} got more than one batch \
1380                                        for a given description in-flight: {:?}",
1381                                    collection_id, in_flight_batches
1382                                );
1383                            }
1384
1385                            continue;
1386                        }
1387                        Event::Progress(frontier) => {
1388                            batch_description_frontier = frontier;
1389                        }
1390                    }
1391                }
1392                Some(event) = batches_input.next() => {
1393                    match event {
1394                        Event::Data(_cap, data) => {
1395                            for batch in data {
1396                                let batch_description = (batch.lower.clone(), batch.upper.clone());
1397
1398                                let batches = in_flight_batches
1399                                    .entry(batch_description)
1400                                    .or_default();
1401
1402                                batches.finished.push(FinishedBatch {
1403                                    batch: write.batch_from_transmittable_batch(batch.batch),
1404                                    data_max_ts: batch.data_max_ts,
1405                                });
1406                                batches.batch_metrics += &batch.metrics;
1407                            }
1408                            continue;
1409                        }
1410                        Event::Progress(frontier) => {
1411                            batches_frontier = frontier;
1412                        }
1413                    }
1414                }
1415                else => {
1416                    // All inputs are exhausted, so we can shut down.
1417                    return Ok(());
1418                }
1419            };
1420
1421            // Peel off any batches that are not beyond the frontier
1422            // anymore.
1423            //
1424            // It is correct to consider batches that are not beyond the
1425            // `batches_frontier` because it is held back by the writer
1426            // operator as long as a) the `batch_description_frontier` did
1427            // not advance and b) as long as the `desired_frontier` has not
1428            // advanced to the `upper` of a given batch description.
1429
1430            let mut done_batches = in_flight_descriptions
1431                .iter()
1432                .filter(|(lower, _upper)| !PartialOrder::less_equal(&batches_frontier, lower))
1433                .cloned()
1434                .collect::<Vec<_>>();
1435
1436            trace!(
1437                "persist_sink {collection_id}/{shard_id}: \
1438                    append_batches: in_flight: {:?}, \
1439                    done: {:?}, \
1440                    batch_frontier: {:?}, \
1441                    batch_description_frontier: {:?}",
1442                in_flight_descriptions,
1443                done_batches,
1444                batches_frontier,
1445                batch_description_frontier
1446            );
1447
1448            // Append batches in order, to ensure that their `lower` and
1449            // `upper` line up.
1450            done_batches.sort_by(|a, b| {
1451                if PartialOrder::less_than(a, b) {
1452                    Ordering::Less
1453                } else if PartialOrder::less_than(b, a) {
1454                    Ordering::Greater
1455                } else {
1456                    Ordering::Equal
1457                }
1458            });
1459
1460            let validate_part_bounds_on_write = write.validate_part_bounds_on_write();
1461            let mut todo = VecDeque::new();
1462
1463            if validate_part_bounds_on_write {
1464                // Persist will expect each batch's bounds to match the append-time bounds; write them separately.
1465                for done_batch_metadata in done_batches.drain(..) {
1466                    in_flight_descriptions.remove(&done_batch_metadata);
1467                    let batch_set = in_flight_batches
1468                        .remove(&done_batch_metadata)
1469                        .unwrap_or_default();
1470                    todo.push_back((done_batch_metadata, batch_set));
1471                }
1472            } else {
1473                // Persist should allow batches to be written as part of a single append even when the bounds don't
1474                // match exactly; group all eligible batches together.
1475                let mut combined_batch_metadata = None;
1476                let mut combined_batch_set = BatchSet::default();
1477                for done_batch_metadata in done_batches.drain(..) {
1478                    in_flight_descriptions.remove(&done_batch_metadata);
1479                    let mut batch_set = in_flight_batches
1480                        .remove(&done_batch_metadata)
1481                        .unwrap_or_default();
1482                    match combined_batch_metadata.as_mut() {
1483                        Some((_, upper)) => *upper = done_batch_metadata.1,
1484                        None => combined_batch_metadata = Some(done_batch_metadata),
1485                    }
1486                    combined_batch_set.batch_metrics += &batch_set.batch_metrics;
1487                    combined_batch_set.finished.append(&mut batch_set.finished);
1488                }
1489                if let Some(done_batch_metadata) = combined_batch_metadata {
1490                    todo.push_back((done_batch_metadata, combined_batch_set))
1491                }
1492            };
1493
1494            while let Some((done_batch_metadata, batch_set)) = todo.pop_front() {
1495                in_flight_descriptions.remove(&done_batch_metadata);
1496
1497                let mut batches = batch_set.finished;
1498
1499                trace!(
1500                    "persist_sink {collection_id}/{shard_id}: \
1501                        done batch: {:?}, {:?}",
1502                    done_batch_metadata,
1503                    batches
1504                );
1505
1506                let (batch_lower, batch_upper) = done_batch_metadata;
1507
1508                let batch_metrics = batch_set.batch_metrics;
1509
1510                let mut to_append = batches.iter_mut().map(|b| &mut b.batch).collect::<Vec<_>>();
1511
1512                let result = {
1513                    let maybe_err = if *read_only_rx.borrow() {
1514
1515                        // We have to wait for either us coming out of read-only
1516                        // mode or someone else applying a write that covers our
1517                        // batch.
1518                        //
1519                        // If we didn't wait for the latter here, and just go
1520                        // around the loop again, we might miss a moment where
1521                        // _we_ have to write down a batch. For example when our
1522                        // input frontier advances to a state where we can
1523                        // write, and the read-write instance sees the same
1524                        // update but then crashes before it can append a batch.
1525
1526                        let maybe_err = loop {
1527                            if collection_id.is_user() {
1528                                tracing::debug!(
1529                                    %worker_id,
1530                                    %collection_id,
1531                                    %shard_id,
1532                                    ?batch_lower,
1533                                    ?batch_upper,
1534                                    ?current_upper,
1535                                    "persist_sink is in read-only mode, waiting until we come out of it or the shard upper advances"
1536                                );
1537                            }
1538
1539                            // We don't try to be smart here, and for example
1540                            // use `wait_for_upper_past()`. We'd have to use a
1541                            // select!, which would require cancel safety of
1542                            // `wait_for_upper_past()`, which it doesn't
1543                            // advertise.
1544                            let _ = tokio::time::timeout(
1545                                Duration::from_secs(1),
1546                                read_only_rx.changed(),
1547                            )
1548                            .await;
1549
1550                            if !*read_only_rx.borrow() {
1551                                if collection_id.is_user() {
1552                                    tracing::debug!(
1553                                        %worker_id,
1554                                        %collection_id,
1555                                        %shard_id,
1556                                        ?batch_lower,
1557                                        ?batch_upper,
1558                                        ?current_upper,
1559                                        "persist_sink has come out of read-only mode"
1560                                    );
1561                                }
1562
1563                                // It's okay to write now.
1564                                break Ok(());
1565                            }
1566
1567                            let current_upper = write.fetch_recent_upper().await;
1568
1569                            if PartialOrder::less_than(&batch_upper, current_upper) {
1570                                // We synthesize an `UpperMismatch` so that we can go
1571                                // through the same logic below for trimming down our
1572                                // batches.
1573                                //
1574                                // Notably, we are not trying to be smart, and teach the
1575                                // write operator about read-only mode. Writing down
1576                                // those batches does not append anything to the persist
1577                                // shard, and it would be a hassle to figure out in the
1578                                // write workers how to trim down batches in read-only
1579                                // mode, when the shard upper advances.
1580                                //
1581                                // Right here, in the logic below, we have all we need
1582                                // for figuring out how to trim our batches.
1583
1584                                if collection_id.is_user() {
1585                                    tracing::debug!(
1586                                        %worker_id,
1587                                        %collection_id,
1588                                        %shard_id,
1589                                        ?batch_lower,
1590                                        ?batch_upper,
1591                                        ?current_upper,
1592                                        "persist_sink not appending in read-only mode"
1593                                    );
1594                                }
1595
1596                                break Err(UpperMismatch {
1597                                    current: current_upper.clone(),
1598                                    expected: batch_lower.clone()}
1599                                );
1600                            }
1601                        };
1602
1603                        maybe_err
1604                    } else {
1605                        // It's okay to proceed with the write.
1606                        Ok(())
1607                    };
1608
1609                    match maybe_err {
1610                        Ok(()) => {
1611                            let _permit = busy_signal.acquire().await;
1612
1613                            write.compare_and_append_batch(
1614                                &mut to_append[..],
1615                                batch_lower.clone(),
1616                                batch_upper.clone(),
1617                                validate_part_bounds_on_write,
1618                            )
1619                            .await
1620                            .expect("Invalid usage")
1621                        },
1622                        Err(e) => {
1623                            // We forward the synthesize error message, so that
1624                            // we go though the batch cleanup logic below.
1625                            Err(e)
1626                        }
1627                    }
1628                };
1629
1630
1631                // These metrics are independent of whether it was _us_ or
1632                // _someone_ that managed to commit a batch that advanced the
1633                // upper.
1634                source_statistics.update_snapshot_committed(&batch_upper);
1635                source_statistics.update_rehydration_latency_ms(&batch_upper);
1636                metrics
1637                    .progress
1638                    .set(mz_persist_client::metrics::encode_ts_metric(&batch_upper));
1639
1640                if collection_id.is_user() {
1641                    trace!(
1642                        "persist_sink {collection_id}/{shard_id}: \
1643                            append result for batch ({:?} -> {:?}): {:?}",
1644                        batch_lower,
1645                        batch_upper,
1646                        result
1647                    );
1648                }
1649
1650                match result {
1651                    Ok(()) => {
1652                        // Only update these metrics when we know that _we_ were
1653                        // successful.
1654                        let committed =
1655                            batch_metrics.inserts + batch_metrics.retractions;
1656                        source_statistics
1657                            .inc_updates_committed_by(committed);
1658                        metrics.processed_batches.inc();
1659                        metrics.row_inserts.inc_by(batch_metrics.inserts);
1660                        metrics.row_retractions.inc_by(batch_metrics.retractions);
1661                        metrics.error_inserts.inc_by(batch_metrics.error_inserts);
1662                        metrics
1663                            .error_retractions
1664                            .inc_by(batch_metrics.error_retractions);
1665
1666                        current_upper.borrow_mut().clone_from(&batch_upper);
1667                        upper_cap_set.downgrade(current_upper.borrow().iter());
1668                    }
1669                    Err(mismatch) => {
1670                        // We tried to to a non-contiguous append, that won't work.
1671                        if PartialOrder::less_than(&mismatch.current, &batch_lower) {
1672                            // Best-effort attempt to delete unneeded batches.
1673                            future::join_all(batches.into_iter().map(|b| b.batch.delete())).await;
1674
1675                            // We always bail when this happens, regardless of
1676                            // `bail_on_concurrent_modification`.
1677                            tracing::warn!(
1678                                "persist_sink({}): invalid upper! \
1679                                    Tried to append batch ({:?} -> {:?}) but upper \
1680                                    is {:?}. This is surpising and likely indicates \
1681                                    a bug in the persist sink, but we'll restart the \
1682                                    dataflow and try again.",
1683                                collection_id, batch_lower, batch_upper, mismatch.current,
1684                            );
1685                            anyhow::bail!("collection concurrently modified. Ingestion dataflow will be restarted");
1686                        } else if PartialOrder::less_than(&mismatch.current, &batch_upper) {
1687                            // The shard's upper was ahead of our batch's lower
1688                            // but not ahead of our upper. Cut down the
1689                            // description by advancing its lower to the current
1690                            // shard upper and try again. IMPORTANT: We can only
1691                            // advance the lower, meaning we cut updates away,
1692                            // we must not "extend" the batch by changing to a
1693                            // lower that is not beyond the current lower. This
1694                            // invariant is checked by the first if branch: if
1695                            // `!(current_upper < lower)` then it holds that
1696                            // `lower <= current_upper`.
1697
1698                            // First, construct a new batch description with the
1699                            // lower advanced to the current shard upper.
1700                            let new_batch_lower = mismatch.current.clone();
1701                            let new_done_batch_metadata =
1702                                (new_batch_lower.clone(), batch_upper.clone());
1703
1704                            // Re-append every batch that still holds something we owe, under the
1705                            // narrowed description. A batch may hold data on both sides of the new
1706                            // lower: persist registers it truncated and filters the updates
1707                            // outside the registered bounds on read, so the ones the concurrent
1708                            // writer already committed do not come back. A batch entirely below
1709                            // the new lower owes nothing and is deleted instead, to keep parts
1710                            // that would be truncated away in full out of shard state.
1711                            let mut batch_delete_futures = vec![];
1712                            let mut new_batch_set = BatchSet::default();
1713                            for batch in batches {
1714                                if new_batch_lower.less_equal(&batch.data_max_ts) {
1715                                    new_batch_set.finished.push(batch);
1716                                } else {
1717                                    batch_delete_futures.push(batch.batch.delete());
1718                                }
1719                            }
1720
1721                            // Re-add the new batch to the list of batches to process.
1722                            todo.push_front((new_done_batch_metadata, new_batch_set));
1723
1724                            // Best-effort attempt to delete unneeded batches.
1725                            future::join_all(batch_delete_futures).await;
1726                        } else {
1727                            // Best-effort attempt to delete unneeded batches.
1728                            future::join_all(batches.into_iter().map(|b| b.batch.delete())).await;
1729                        }
1730
1731                        if bail_on_concurrent_modification {
1732                            tracing::warn!(
1733                                "persist_sink({}): invalid upper! \
1734                                    Tried to append batch ({:?} -> {:?}) but upper \
1735                                    is {:?}. This is not a problem, it just means \
1736                                    someone else was faster than us. We will try \
1737                                    again with a new batch description.",
1738                                collection_id, batch_lower, batch_upper, mismatch.current,
1739                            );
1740                            anyhow::bail!("collection concurrently modified. Ingestion dataflow will be restarted");
1741                        }
1742                    }
1743                }
1744            }
1745        }
1746    }));
1747
1748    (upper_stream, errors, shutdown_button.press_on_drop())
1749}
1750
1751#[cfg(test)]
1752mod tests {
1753    use std::cell::RefCell;
1754    use std::str::FromStr;
1755
1756    use mz_build_info::DUMMY_BUILD_INFO;
1757    use mz_dyncfg::{ConfigUpdates, ConfigVal};
1758    use mz_ore::metrics::MetricsRegistry;
1759    use mz_ore::now::SYSTEM_TIME;
1760    use mz_ore::url::SensitiveUrl;
1761    use mz_persist_client::PersistLocation;
1762    use mz_persist_client::cfg::PersistConfig;
1763    use mz_persist_client::rpc::PubSubClientConnection;
1764    use mz_persist_types::ShardId;
1765    use mz_repr::{Datum, RelationDesc, SqlScalarType};
1766    use mz_storage_types::sources::SourceEnvelope;
1767    use mz_storage_types::sources::envelope::{KeyEnvelope, NoneEnvelope};
1768    use timely::dataflow::operators::Input;
1769
1770    use crate::statistics::SourceStatisticsMetricDefs;
1771
1772    use super::*;
1773
1774    fn ts(t: u64) -> mz_repr::Timestamp {
1775        t.into()
1776    }
1777
1778    fn frontier(t: u64) -> Antichain<mz_repr::Timestamp> {
1779        Antichain::from_elem(ts(t))
1780    }
1781
1782    /// One step of a `write_batches` script.
1783    #[derive(Clone)]
1784    enum Step {
1785        /// Deliver a batch description, as `mint_batch_descriptions` would.
1786        Description(u64, u64),
1787        /// Deliver a commitment, as `mint_batch_descriptions` would.
1788        Commit(u64, u64),
1789        /// Deliver `count` updates at time `at`.
1790        Updates(u64, usize),
1791        /// Advance both input frontiers.
1792        AdvanceTo(u64),
1793    }
1794
1795    /// What a batch emitted by `write_batches` carries, flattened for assertions.
1796    #[derive(Clone, Debug, PartialEq, Eq, PartialOrd, Ord)]
1797    struct EmittedBatch {
1798        lower: u64,
1799        upper: u64,
1800        data_max_ts: u64,
1801        inserts: u64,
1802    }
1803
1804    /// Runs a timely worker to completion on a blocking thread.
1805    ///
1806    /// The test body is a single poll of the runtime's `block_on` future, so it runs under one
1807    /// tokio cooperative budget. A worker driven inline spends that budget on the operators'
1808    /// `select!` and semaphore polls, and once it is gone every such poll returns `Pending` and
1809    /// re-wakes itself, parking the operator for good with no error. Blocking threads have no
1810    /// budget.
1811    async fn run_worker<T: Send + 'static>(
1812        worker: impl FnOnce(&mut timely::worker::Worker) -> T + Send + Sync + 'static,
1813    ) -> T {
1814        mz_ore::task::spawn_blocking(
1815            || "persist_sink_test_worker",
1816            move || timely::execute_directly(worker),
1817        )
1818        .await
1819    }
1820
1821    /// Drives `write_batches` through `script` and returns the batches it emitted, along with a
1822    /// handle to the shard so callers can append them and read the result back.
1823    async fn run_write_batches(
1824        target: CollectionMetadata,
1825        persist_clients: Arc<PersistClientCache>,
1826        script: Vec<Step>,
1827    ) -> Vec<(EmittedBatch, ProtoBatch)> {
1828        run_worker(move |worker| {
1829            // `ProtoBatch` is not `Ord`, so the captured stream is summarized on the way out
1830            // rather than going through `Capture`.
1831            let emitted = Rc::new(RefCell::new(Vec::new()));
1832
1833            let (mut descs_input, mut ceilings_input, mut data_input, button) =
1834                worker.dataflow::<mz_repr::Timestamp, _, _>(|scope| {
1835                    let (descs_input, descs) = scope.new_input();
1836                    let (ceilings_input, ceilings) = scope.new_input();
1837                    let (data_input, data) = scope.new_input();
1838
1839                    let source_id = GlobalId::User(0);
1840                    let stats_defs =
1841                        SourceStatisticsMetricDefs::register_with(&MetricsRegistry::new());
1842                    let source_statistics = SourceStatistics::new(
1843                        source_id,
1844                        0,
1845                        &stats_defs,
1846                        source_id,
1847                        &target.data_shard,
1848                        SourceEnvelope::None(NoneEnvelope {
1849                            key_envelope: KeyEnvelope::None,
1850                            key_arity: 0,
1851                        }),
1852                        Antichain::from_elem(Timestamp::minimum()),
1853                    );
1854
1855                    let (batches, button) = write_batches(
1856                        scope,
1857                        source_id,
1858                        "test",
1859                        &target,
1860                        descs,
1861                        ceilings,
1862                        data.as_collection(),
1863                        persist_clients,
1864                        source_statistics,
1865                        Arc::new(Semaphore::new(Semaphore::MAX_PERMITS)),
1866                    );
1867                    let sink = Rc::clone(&emitted);
1868                    InspectCore::inspect_container(batches, move |event| {
1869                        if let Ok((_, data)) = event {
1870                            for b in data {
1871                                sink.borrow_mut().push((
1872                                    EmittedBatch {
1873                                        lower: b.lower.as_option().expect("single lower").into(),
1874                                        upper: b.upper.as_option().expect("single upper").into(),
1875                                        data_max_ts: b.data_max_ts.into(),
1876                                        inserts: b.metrics.inserts,
1877                                    },
1878                                    b.batch.clone(),
1879                                ));
1880                            }
1881                        }
1882                    });
1883
1884                    (descs_input, ceilings_input, data_input, button)
1885                });
1886
1887            // We want the operator to finish processing before we advance the script,
1888            // but when the operator is waiting for persist,
1889            // a plain timely `step` will find no work to do and return immediately.
1890            // So we repeatedly `step_or_park`, where `park`ing the thread forces a delay,
1891            // allowing the operator to finish waiting for persist and do its work on the next `step`.
1892            fn pump(worker: &mut timely::worker::Worker) {
1893                // no solid reason for this number, it's a selection that seems high enough to
1894                // work reliably, but not create a ton of delay (~32ms parked).
1895                for _ in 0..32 {
1896                    worker.step_or_park(Some(Duration::from_millis(1)));
1897                }
1898            }
1899
1900            // Twice, so the operator is past opening its persist handles before the script runs.
1901            pump(worker);
1902            pump(worker);
1903
1904            for step in script {
1905                match step {
1906                    Step::Description(lower, upper) => {
1907                        descs_input.send((frontier(lower), frontier(upper)));
1908                    }
1909                    Step::Commit(lower, ceiling) => ceilings_input.send(Commitment {
1910                        lower: ts(lower),
1911                        ceiling: ts(ceiling),
1912                    }),
1913                    Step::Updates(at, count) => {
1914                        for i in 0..i64::try_from(count).expect("small count") {
1915                            let row = Row::pack_slice(&[Datum::Int64(i)]);
1916                            data_input.send((Ok(row), ts(at), Diff::ONE));
1917                        }
1918                    }
1919                    Step::AdvanceTo(t) => {
1920                        descs_input.advance_to(ts(t));
1921                        ceilings_input.advance_to(ts(t));
1922                        data_input.advance_to(ts(t));
1923                    }
1924                }
1925                // NOTE: `send` buffers until its container fills, so without a flush every step
1926                // before the next `advance_to` would reach the operator together, in one round.
1927                descs_input.flush();
1928                ceilings_input.flush();
1929                data_input.flush();
1930                pump(worker);
1931            }
1932
1933            descs_input.close();
1934            ceilings_input.close();
1935            data_input.close();
1936            for _ in 0..1_000 {
1937                if !worker.step_or_park(Some(Duration::from_millis(1))) {
1938                    break;
1939                }
1940            }
1941
1942            drop(button);
1943            while worker.step() {}
1944
1945            let mut emitted = emitted.borrow().clone();
1946            emitted.sort_by(|a, b| a.0.cmp(&b.0));
1947            emitted
1948        })
1949        .await
1950    }
1951
1952    fn test_target() -> CollectionMetadata {
1953        CollectionMetadata {
1954            persist_location: PersistLocation {
1955                blob_uri: SensitiveUrl::from_str("mem://").expect("invalid URL"),
1956                consensus_uri: SensitiveUrl::from_str("mem://").expect("invalid URL"),
1957            },
1958            data_shard: ShardId::new(),
1959            relation_desc: RelationDesc::builder()
1960                .with_column("a", SqlScalarType::Int64.nullable(false))
1961                .finish(),
1962            txns_shard: None,
1963        }
1964    }
1965
1966    /// Turn on part bounds validation so _append_ checks the bounds the sink writes.
1967    /// Both settings default off in code but are turned on in production.
1968    fn test_persist_clients() -> Arc<PersistClientCache> {
1969        let persist_cfg =
1970            PersistConfig::new_default_configs(&DUMMY_BUILD_INFO, SYSTEM_TIME.clone());
1971        let mut updates = ConfigUpdates::default();
1972        updates.add_dynamic(
1973            "persist_validate_part_bounds_on_write",
1974            ConfigVal::Bool(true),
1975        );
1976        updates.add_dynamic(
1977            "persist_validate_part_bounds_on_read",
1978            ConfigVal::Bool(true),
1979        );
1980        updates.apply(&persist_cfg.configs);
1981        Arc::new(PersistClientCache::new(
1982            persist_cfg,
1983            &MetricsRegistry::new(),
1984            |_, _| PubSubClientConnection::noop(),
1985        ))
1986    }
1987
1988    /// A single `compare_and_append` over `[lower, upper)` carrying every emitted batch.
1989    fn one_append(
1990        emitted: Vec<(EmittedBatch, ProtoBatch)>,
1991        lower: u64,
1992        upper: u64,
1993    ) -> Vec<(u64, u64, Vec<ProtoBatch>)> {
1994        vec![(lower, upper, emitted.into_iter().map(|(_, p)| p).collect())]
1995    }
1996
1997    /// One `compare_and_append` per description the batches were written for, ascending by lower.
1998    fn append_per_description(
1999        emitted: Vec<(EmittedBatch, ProtoBatch)>,
2000    ) -> Vec<(u64, u64, Vec<ProtoBatch>)> {
2001        let mut by_desc: BTreeMap<(u64, u64), Vec<ProtoBatch>> = BTreeMap::new();
2002        for (batch, proto) in emitted {
2003            by_desc
2004                .entry((batch.lower, batch.upper))
2005                .or_default()
2006                .push(proto);
2007        }
2008        by_desc
2009            .into_iter()
2010            .map(|((lower, upper), protos)| (lower, upper, protos))
2011            .collect()
2012    }
2013
2014    /// Applies each entry in `appends` as one `compare_and_append` over `[lower, upper)`, in order,
2015    /// then reads the shard back as of `as_of` and returns the summed diffs.
2016    ///
2017    /// Batches written for different descriptions need separate entries, because persist rejects a
2018    /// batch whose upper is below the append upper. A `lower` above a batch's own lower registers
2019    /// it truncated, which is what the sink relies on when a concurrent writer has already claimed
2020    /// part of the range.
2021    ///
2022    /// Part bounds validation is what catches a batch whose parts reach outside their registered
2023    /// bounds, so the tests append for real rather than stopping at what `write_batches` emitted.
2024    async fn append_and_read_back(
2025        target: &CollectionMetadata,
2026        persist_clients: &PersistClientCache,
2027        appends: Vec<(u64, u64, Vec<ProtoBatch>)>,
2028        as_of: u64,
2029    ) -> i64 {
2030        let persist_client = persist_clients
2031            .open(target.persist_location.clone())
2032            .await
2033            .expect("could not open persist client");
2034        let mut write = persist_client
2035            .open_writer::<SourceData, (), mz_repr::Timestamp, StorageDiff>(
2036                target.data_shard,
2037                Arc::new(target.relation_desc.clone()),
2038                Arc::new(UnitSchema),
2039                Diagnostics::for_tests(),
2040            )
2041            .await
2042            .expect("could not open persist shard");
2043
2044        assert!(
2045            write.validate_part_bounds_on_write(),
2046            "part bounds validation is off, so this append proves nothing about batch bounds"
2047        );
2048
2049        for (lower, upper, protos) in appends {
2050            let mut batches: Vec<_> = protos
2051                .into_iter()
2052                .map(|proto| write.batch_from_transmittable_batch(proto))
2053                .collect();
2054            let mut to_append: Vec<_> = batches.iter_mut().collect();
2055            write
2056                .compare_and_append_batch(
2057                    &mut to_append[..],
2058                    frontier(lower),
2059                    frontier(upper),
2060                    true,
2061                )
2062                .await
2063                .expect("invalid usage")
2064                .expect("upper mismatch");
2065
2066            assert_eq!(write.fetch_recent_upper().await, &frontier(upper));
2067        }
2068
2069        let mut read = persist_client
2070            .open_leased_reader::<SourceData, (), mz_repr::Timestamp, StorageDiff>(
2071                target.data_shard,
2072                Arc::new(target.relation_desc.clone()),
2073                Arc::new(UnitSchema),
2074                Diagnostics::for_tests(),
2075                true,
2076            )
2077            .await
2078            .expect("invalid usage");
2079        let contents = read
2080            .snapshot_and_fetch(frontier(as_of))
2081            .await
2082            .expect("since <= as_of");
2083
2084        contents.iter().map(|(_, _, d)| *d).sum()
2085    }
2086
2087    /// Several descriptions can become ready in the same pass. Each is written under its own
2088    /// bounds, so every batch holds exactly the updates its description covers however that ready
2089    /// set happens to be ordered.
2090    ///
2091    /// NOTE: `in_flight_batches` is a `HashMap`, so the ready set comes out in no particular order.
2092    /// Enough descriptions are used here that an all-ascending pass is unlikely.
2093    #[mz_ore::test(tokio::test(flavor = "multi_thread"))]
2094    #[cfg_attr(miri, ignore)] // unsupported operation: returning ready events from epoll_wait
2095    async fn write_batches_handles_descriptions_ready_in_one_pass() {
2096        const DESCRIPTIONS: u64 = 6;
2097        const DONE: u64 = DESCRIPTIONS * 2;
2098
2099        let persist_clients = test_persist_clients();
2100        let target = test_target();
2101
2102        // One update inside each of the tiling descriptions [0,2), [2,4), ... None of them is ready
2103        // until the frontier passes every upper, so they all come due together.
2104        let mut script = vec![];
2105        for i in 0..DESCRIPTIONS {
2106            script.push(Step::Updates(i * 2 + 1, 1));
2107        }
2108        for i in 0..DESCRIPTIONS {
2109            script.push(Step::Description(i * 2, i * 2 + 2));
2110        }
2111        script.push(Step::AdvanceTo(DONE));
2112
2113        let emitted = run_write_batches(target.clone(), Arc::clone(&persist_clients), script).await;
2114
2115        assert_eq!(
2116            emitted.len(),
2117            usize::cast_from(DESCRIPTIONS),
2118            "one batch per description, got {:?}",
2119            emitted.iter().map(|(b, _)| b).collect::<Vec<_>>()
2120        );
2121        for (batch, _) in &emitted {
2122            assert!(
2123                batch.lower <= batch.data_max_ts && batch.data_max_ts < batch.upper,
2124                "batch {batch:?} holds data outside the description it was written for"
2125            );
2126        }
2127
2128        let total = append_and_read_back(
2129            &target,
2130            &persist_clients,
2131            append_per_description(emitted),
2132            DONE - 1,
2133        )
2134        .await;
2135        assert_eq!(
2136            total,
2137            i64::try_from(DESCRIPTIONS).expect("small"),
2138            "every update should be readable exactly once"
2139        );
2140    }
2141
2142    /// A description that covers no updates must emit no batch, rather than open a builder that
2143    /// has no data bounds to register.
2144    #[mz_ore::test(tokio::test(flavor = "multi_thread"))]
2145    #[cfg_attr(miri, ignore)] // unsupported operation: returning ready events from epoll_wait
2146    async fn write_batches_emits_nothing_for_a_description_with_no_updates() {
2147        const SPLIT: u64 = 4;
2148        const DONE: u64 = 8;
2149
2150        // Two descriptions in hand, with data only in the second.
2151        let emitted = run_write_batches(
2152            test_target(),
2153            test_persist_clients(),
2154            vec![
2155                Step::Description(0, SPLIT),
2156                Step::Description(SPLIT, DONE),
2157                Step::AdvanceTo(SPLIT),
2158                Step::Updates(SPLIT, 8),
2159                Step::AdvanceTo(DONE),
2160            ],
2161        )
2162        .await;
2163
2164        assert_eq!(
2165            emitted
2166                .iter()
2167                .map(|(b, _)| (b.lower, b.upper))
2168                .collect::<Vec<_>>(),
2169            vec![(SPLIT, DONE)],
2170            "only the description holding data should produce a batch",
2171        );
2172    }
2173
2174    /// A snapshot at time 1 pinning the frontier while replication delivers one update at each of
2175    /// times 2..=`pinned_times`+1, with the description that covers the whole snapshot arriving
2176    /// only at the end.
2177    fn pinned_frontier_script(snapshot_rows: usize, pinned_times: u64, done: u64) -> Vec<Step> {
2178        let mut script = vec![Step::Updates(1, snapshot_rows)];
2179        for t in 2..=pinned_times + 1 {
2180            script.push(Step::Updates(t, 1));
2181        }
2182        // The minter holds a capability at the shard upper for the whole snapshot, so its one
2183        // description is emitted there, and the frontier then jumps past everything staged.
2184        script.push(Step::Description(0, done));
2185        script.push(Step::AdvanceTo(done));
2186        script
2187    }
2188
2189    /// A snapshot pins the export's frontier at its as_of while concurrent replication keeps
2190    /// delivering updates at later times. Each timestamp writes a batch of its own, all finished
2191    /// under the one description that arrives when the snapshot finishes.
2192    #[mz_ore::test(tokio::test(flavor = "multi_thread"))]
2193    #[cfg_attr(miri, ignore)] // unsupported operation: returning ready events from epoll_wait
2194    async fn write_batches_writes_one_batch_per_timestamp() {
2195        const SNAPSHOT_ROWS: usize = 4;
2196        const PINNED_TIMES: u64 = 16;
2197        const DONE: u64 = PINNED_TIMES + 2;
2198
2199        let persist_clients = test_persist_clients();
2200        let target = test_target();
2201
2202        let emitted = run_write_batches(
2203            target.clone(),
2204            Arc::clone(&persist_clients),
2205            pinned_frontier_script(SNAPSHOT_ROWS, PINNED_TIMES, DONE),
2206        )
2207        .await;
2208
2209        // Every batch carries the description's bounds, since that is what they are finished
2210        // under, and holds a single timestamp's updates.
2211        let expected: Vec<_> = std::iter::once(EmittedBatch {
2212            lower: 0,
2213            upper: DONE,
2214            data_max_ts: 1,
2215            inserts: u64::cast_from(SNAPSHOT_ROWS),
2216        })
2217        .chain((2..=PINNED_TIMES + 1).map(|ts| EmittedBatch {
2218            lower: 0,
2219            upper: DONE,
2220            data_max_ts: ts,
2221            inserts: 1,
2222        }))
2223        .collect();
2224        assert_eq!(
2225            emitted.iter().map(|(b, _)| b.clone()).collect::<Vec<_>>(),
2226            expected,
2227        );
2228
2229        let total = append_and_read_back(
2230            &target,
2231            &persist_clients,
2232            one_append(emitted, 0, DONE),
2233            DONE - 1,
2234        )
2235        .await;
2236        assert_eq!(
2237            total,
2238            i64::try_from(SNAPSHOT_ROWS).expect("small")
2239                + i64::try_from(PINNED_TIMES).expect("small"),
2240            "the same updates should be readable however they were batched"
2241        );
2242    }
2243
2244    /// The writer drains descriptions ahead of data in each round, and input keeps queueing while
2245    /// it awaits persist, so rows at a timestamp can be staged after the description covering them
2246    /// was processed. They belong in the batch that timestamp already has.
2247    #[mz_ore::test(tokio::test(flavor = "multi_thread"))]
2248    #[cfg_attr(miri, ignore)] // unsupported operation: returning ready events from epoll_wait
2249    async fn write_batches_keeps_one_batch_per_timestamp_for_rows_after_the_description() {
2250        const ROWS: usize = 4;
2251        const DONE: u64 = 4;
2252
2253        let persist_clients = test_persist_clients();
2254        let target = test_target();
2255
2256        let emitted = run_write_batches(
2257            target.clone(),
2258            Arc::clone(&persist_clients),
2259            vec![
2260                Step::Updates(1, ROWS),
2261                Step::Description(0, DONE),
2262                Step::Updates(1, ROWS),
2263                Step::AdvanceTo(DONE),
2264            ],
2265        )
2266        .await;
2267
2268        assert_eq!(
2269            emitted.iter().map(|(b, _)| b.clone()).collect::<Vec<_>>(),
2270            vec![EmittedBatch {
2271                lower: 0,
2272                upper: DONE,
2273                data_max_ts: 1,
2274                inserts: u64::cast_from(ROWS * 2),
2275            }],
2276        );
2277
2278        let total = append_and_read_back(
2279            &target,
2280            &persist_clients,
2281            one_append(emitted, 0, DONE),
2282            DONE - 1,
2283        )
2284        .await;
2285        assert_eq!(total, i64::try_from(ROWS * 2).expect("small"));
2286    }
2287
2288    /// The same snapshot with a ceiling committed first, which is what the minter does behind a
2289    /// frontier that is not moving. The description itself only arrives once the frontier reaches
2290    /// the ceiling, which is what ends the script.
2291    fn committed_ceiling_script(snapshot_rows: usize, pinned_times: u64, done: u64) -> Vec<Step> {
2292        let mut script = vec![
2293            Step::Commit(0, done),
2294            Step::AdvanceTo(1),
2295            Step::Updates(1, snapshot_rows),
2296        ];
2297        for t in 2..=pinned_times + 1 {
2298            script.push(Step::Updates(t, 1));
2299        }
2300        script.push(Step::Description(0, done));
2301        script.push(Step::AdvanceTo(done));
2302        script
2303    }
2304
2305    /// A ceiling committed ahead of the frontier gives arriving updates a bound, so a pinned
2306    /// frontier writes one batch rather than one per timestamp, and the rows sit in a builder that
2307    /// fills to the blob target instead of many single-timestamp builders that each stay under it
2308    /// and hold their rows in memory.
2309    #[mz_ore::test(tokio::test(flavor = "multi_thread"))]
2310    #[cfg_attr(miri, ignore)] // unsupported operation: returning ready events from epoll_wait
2311    async fn write_batches_routes_updates_below_the_ceiling_into_one_builder() {
2312        const SNAPSHOT_ROWS: usize = 4;
2313        const PINNED_TIMES: u64 = 16;
2314        const DONE: u64 = PINNED_TIMES + 2;
2315
2316        let persist_clients = test_persist_clients();
2317        let target = test_target();
2318
2319        let emitted = run_write_batches(
2320            target.clone(),
2321            Arc::clone(&persist_clients),
2322            committed_ceiling_script(SNAPSHOT_ROWS, PINNED_TIMES, DONE),
2323        )
2324        .await;
2325
2326        assert_eq!(
2327            emitted.len(),
2328            1,
2329            "the whole snapshot should share the one open builder, got {:?}",
2330            emitted.iter().map(|(b, _)| b).collect::<Vec<_>>()
2331        );
2332        assert_eq!(
2333            emitted[0].0,
2334            EmittedBatch {
2335                lower: 0,
2336                upper: DONE,
2337                data_max_ts: PINNED_TIMES + 1,
2338                inserts: u64::cast_from(SNAPSHOT_ROWS) + PINNED_TIMES,
2339            }
2340        );
2341
2342        let total = append_and_read_back(
2343            &target,
2344            &persist_clients,
2345            one_append(emitted, 0, DONE),
2346            DONE - 1,
2347        )
2348        .await;
2349        assert_eq!(
2350            total,
2351            i64::try_from(SNAPSHOT_ROWS).expect("small")
2352                + i64::try_from(PINNED_TIMES).expect("small"),
2353            "grouping must not change what the shard ends up holding"
2354        );
2355    }
2356
2357    #[mz_ore::test]
2358    fn next_mint_paces_the_ceiling_on_the_remap_upper() {
2359        const LOOKAHEAD: u64 = 10;
2360        let lower = frontier(0);
2361        let pinned = frontier(0);
2362
2363        // A pinned frontier derives no description, so the pass commits ahead of the remap upper,
2364        // a tick raises the ceiling, and a tick that does not clear it commits nothing.
2365        assert_eq!(
2366            next_mint(&lower, &pinned, &frontier(5), None, Some(LOOKAHEAD)),
2367            Some(Mint::Ceiling(ts(15)))
2368        );
2369        assert_eq!(
2370            next_mint(&lower, &pinned, &frontier(6), Some(ts(15)), Some(LOOKAHEAD)),
2371            Some(Mint::Ceiling(ts(16)))
2372        );
2373        assert_eq!(
2374            next_mint(&lower, &pinned, &frontier(6), Some(ts(16)), Some(LOOKAHEAD)),
2375            None
2376        );
2377        // A closed remap stream has no upper to commit past.
2378        assert_eq!(
2379            next_mint(
2380                &lower,
2381                &pinned,
2382                &Antichain::new(),
2383                Some(ts(16)),
2384                Some(LOOKAHEAD)
2385            ),
2386            None
2387        );
2388        // A lower at the ceiling, as when the current upper jumps to an as_of ahead of the observed
2389        // remap upper, leaves an empty range to commit.
2390        assert_eq!(
2391            next_mint(
2392                &frontier(15),
2393                &frontier(15),
2394                &frontier(5),
2395                None,
2396                Some(LOOKAHEAD)
2397            ),
2398            None
2399        );
2400        // Once the snapshot ends the ceiling binds until the frontier reaches it.
2401        assert_eq!(
2402            next_mint(&lower, &frontier(12), &frontier(12), Some(ts(16)), None),
2403            None
2404        );
2405        assert_eq!(
2406            next_mint(&lower, &frontier(16), &frontier(16), Some(ts(16)), None),
2407            Some(Mint::Description(frontier(16)))
2408        );
2409    }
2410}