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 ¤t_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}