Skip to main content

mz_storage/source/
source_reader_pipeline.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//! Types related to the creation of dataflow raw sources.
11//!
12//! Raw sources are differential dataflow  collections of data directly produced by the
13//! upstream service. The main export of this module is [`create_raw_source`],
14//! which turns [`RawSourceCreationConfig`]s into the aforementioned streams.
15//!
16//! The full source, which is the _differential_ stream that represents the actual object
17//! created by a `CREATE SOURCE` statement, is created by composing
18//! [`create_raw_source`] with
19//! decoding, `SourceEnvelope` rendering, and more.
20//!
21
22// https://github.com/tokio-rs/prost/issues/237
23#![allow(missing_docs)]
24#![allow(clippy::needless_borrow)]
25
26use std::cell::RefCell;
27use std::collections::{BTreeMap, VecDeque};
28use std::hash::{Hash, Hasher};
29use std::rc::Rc;
30use std::sync::Arc;
31use std::time::Duration;
32
33use differential_dataflow::lattice::Lattice;
34use differential_dataflow::{AsCollection, Hashable, VecCollection};
35use futures::stream::StreamExt;
36use mz_ore::cast::CastFrom;
37use mz_ore::collections::CollectionExt;
38use mz_ore::now::NowFn;
39use mz_persist_client::cache::PersistClientCache;
40use mz_repr::{Diff, GlobalId, RelationDesc, Row};
41use mz_storage_types::configuration::StorageConfiguration;
42use mz_storage_types::controller::CollectionMetadata;
43use mz_storage_types::errors::DataflowError;
44use mz_storage_types::sources::{SourceConnection, SourceExport, SourceTimestamp};
45use mz_timely_util::antichain::AntichainExt;
46use mz_timely_util::builder_async::{OperatorBuilder as AsyncOperatorBuilder, PressOnDropButton};
47use mz_timely_util::capture::PusherCapture;
48use mz_timely_util::operator::ConcatenateFlatten;
49use mz_timely_util::reclock::reclock;
50use timely::PartialOrder;
51use timely::container::CapacityContainerBuilder;
52use timely::dataflow::channels::pact::Pipeline;
53use timely::dataflow::operators::capture::capture::Capture;
54use timely::dataflow::operators::core::Map as _;
55use timely::dataflow::operators::generic::OutputBuilder;
56use timely::dataflow::operators::generic::builder_rc::OperatorBuilder as OperatorBuilderRc;
57use timely::dataflow::operators::vec::Broadcast;
58use timely::dataflow::operators::{CapabilitySet, InspectCore, Leave};
59use timely::dataflow::{Scope, StreamVec};
60use timely::order::TotalOrder;
61use timely::progress::frontier::MutableAntichain;
62use timely::progress::{Antichain, Timestamp};
63use tokio::sync::{Semaphore, watch};
64use tokio_stream::wrappers::WatchStream;
65use tracing::trace;
66
67use crate::healthcheck::{HealthStatusMessage, HealthStatusUpdate};
68use crate::metrics::StorageMetrics;
69use crate::metrics::source::SourceMetrics;
70use crate::source::reclock::ReclockOperator;
71use crate::source::types::{Probe, SourceMessage, SourceOutput, SourceRender, StackedCollection};
72use crate::statistics::SourceStatistics;
73
74/// Shared configuration information for all source types. This is used in the
75/// `create_raw_source` functions, which produce raw sources.
76#[derive(Clone)]
77pub struct RawSourceCreationConfig {
78    /// The name to attach to the underlying timely operator.
79    pub name: String,
80    /// The ID of this instantiation of this source.
81    pub id: GlobalId,
82    /// The details of the outputs from this ingestion.
83    pub source_exports: BTreeMap<GlobalId, SourceExport<CollectionMetadata>>,
84    /// The ID of the worker on which this operator is executing
85    pub worker_id: usize,
86    /// The total count of workers
87    pub worker_count: usize,
88    /// Granularity with which timestamps should be closed (and capabilities
89    /// downgraded).
90    pub timestamp_interval: Duration,
91    /// The function to return a now time.
92    pub now_fn: NowFn,
93    /// The metrics & registry that each source instantiates.
94    pub metrics: StorageMetrics,
95    /// The upper frontier this source should resume ingestion at
96    pub as_of: Antichain<mz_repr::Timestamp>,
97    /// For each source export, the upper frontier this source should resume ingestion at in the
98    /// system time domain.
99    pub resume_uppers: BTreeMap<GlobalId, Antichain<mz_repr::Timestamp>>,
100    /// For each source export, the upper frontier this source should resume ingestion at in the
101    /// source time domain.
102    ///
103    /// Since every source has a different timestamp type we carry the timestamps of this frontier
104    /// in an encoded `Vec<Row>` form which will get decoded once we reach the connection
105    /// specialized functions.
106    pub source_resume_uppers: BTreeMap<GlobalId, Vec<Row>>,
107    /// A handle to the persist client cache
108    pub persist_clients: Arc<PersistClientCache>,
109    /// Collection of `SourceStatistics` for source and exports to share updates.
110    pub statistics: BTreeMap<GlobalId, SourceStatistics>,
111    /// Enables reporting the remap operator's write frontier.
112    pub shared_remap_upper: Rc<RefCell<Antichain<mz_repr::Timestamp>>>,
113    /// Configuration parameters, possibly from LaunchDarkly
114    pub config: StorageConfiguration,
115    /// The ID of this source remap/progress collection.
116    pub remap_collection_id: GlobalId,
117    /// The storage metadata for the remap/progress collection
118    pub remap_metadata: CollectionMetadata,
119    // A semaphore that should be acquired by async operators in order to signal that upstream
120    // operators should slow down.
121    pub busy_signal: Arc<Semaphore>,
122}
123
124/// Reduced version of [`RawSourceCreationConfig`] that is used when rendering
125/// each export.
126#[derive(Clone)]
127pub struct SourceExportCreationConfig {
128    /// The ID of this instantiation of this source.
129    pub id: GlobalId,
130    /// The ID of the worker on which this operator is executing
131    pub worker_id: usize,
132    /// The metrics & registry that each source instantiates.
133    pub metrics: StorageMetrics,
134    /// Place to share statistics updates with storage state.
135    pub source_statistics: SourceStatistics,
136}
137
138impl RawSourceCreationConfig {
139    /// Returns the worker id responsible for handling the given partition.
140    pub fn responsible_worker<P: Hash>(&self, partition: P) -> usize {
141        let mut h = std::hash::DefaultHasher::default();
142        (self.id, partition).hash(&mut h);
143        let key = usize::cast_from(h.finish());
144        key % self.worker_count
145    }
146
147    /// Returns true if this worker is responsible for handling the given partition.
148    pub fn responsible_for<P: Hash>(&self, partition: P) -> bool {
149        self.responsible_worker(partition) == self.worker_id
150    }
151}
152
153/// Creates a source dataflow operator graph from a source connection. The type of SourceConnection
154/// determines the type of connection that _should_ be created.
155///
156/// This is also the place where _reclocking_
157/// (<https://github.com/MaterializeInc/materialize/blob/main/doc/developer/design/20210714_reclocking.md>)
158/// happens.
159///
160/// See the [`source` module docs](crate::source) for more details about how raw
161/// sources are used.
162///
163/// The `resume_stream` parameter will contain frontier updates whenever times are durably
164/// recorded which allows the ingestion to release upstream resources.
165///
166/// Alongside the reclocked exports this returns a no-data stream whose frontier is the remap upper.
167pub fn create_raw_source<'scope, 'root, C>(
168    scope: Scope<'scope, mz_repr::Timestamp>,
169    root_scope: Scope<'root, ()>,
170    storage_state: &crate::storage_state::StorageState,
171    committed_upper: StreamVec<'scope, mz_repr::Timestamp, ()>,
172    config: &RawSourceCreationConfig,
173    source_connection: C,
174    start_signal: impl std::future::Future<Output = ()> + 'static,
175) -> (
176    BTreeMap<
177        GlobalId,
178        VecCollection<
179            'scope,
180            mz_repr::Timestamp,
181            Result<SourceOutput<C::Time>, DataflowError>,
182            Diff,
183        >,
184    >,
185    StreamVec<'root, (), HealthStatusMessage>,
186    StreamVec<'scope, mz_repr::Timestamp, ()>,
187    Vec<PressOnDropButton>,
188)
189where
190    C: SourceConnection + SourceRender + Clone + 'static,
191{
192    let worker_id = config.worker_id;
193    let id = config.id;
194
195    let mut tokens = vec![];
196
197    let (probed_upper_tx, probed_upper_rx) = watch::channel(None);
198
199    let source_metrics = Arc::new(config.metrics.get_source_metrics(id, worker_id));
200
201    let timestamp_desc = source_connection.timestamp_desc();
202
203    let (remap_collection, remap_token) = remap_operator(
204        scope,
205        storage_state,
206        config.clone(),
207        probed_upper_rx,
208        timestamp_desc,
209    );
210    // Need to broadcast the remap changes to all workers.
211    let remap_collection = remap_collection.inner.broadcast().as_collection();
212    tokens.push(remap_token);
213
214    // Drops the bidings, as this stream is only used to track the remap upper, which drives
215    // ceiling calculation in the persist sink during snapshots.
216    let remap_upper = remap_collection
217        .inner
218        .clone()
219        .flat_map::<Vec<()>, _, _>(|_| None::<()>);
220
221    let committed_upper = reclock_committed_upper(
222        remap_collection.clone(),
223        config.as_of.clone(),
224        committed_upper,
225        id,
226        Arc::clone(&source_metrics),
227    );
228
229    let mut reclocked_exports = BTreeMap::new();
230
231    let reclocked_exports2 = &mut reclocked_exports;
232    let (health, source_tokens) = root_scope.scoped("SourceTimeDomain", move |scope| {
233        let (exports, health_stream, source_tokens) = source_render_operator(
234            scope,
235            config,
236            source_connection,
237            probed_upper_tx,
238            committed_upper,
239            start_signal,
240        );
241
242        for (id, export) in exports {
243            let (reclock_pusher, reclocked) =
244                reclock(remap_collection.clone(), config.as_of.clone());
245            export
246                .inner
247                .map(move |(result, from_time, diff)| {
248                    let result = match result {
249                        Ok(msg) => Ok(SourceOutput {
250                            key: msg.key,
251                            value: msg.value,
252                            metadata: msg.metadata,
253                            from_time: from_time.clone(),
254                        }),
255                        Err(err) => Err(err),
256                    };
257                    (result, from_time, diff)
258                })
259                .capture_into(PusherCapture(reclock_pusher));
260            reclocked_exports2.insert(id, reclocked);
261        }
262
263        (health_stream.leave(root_scope), source_tokens)
264    });
265
266    tokens.extend(source_tokens);
267
268    (reclocked_exports, health, remap_upper, tokens)
269}
270
271/// Renders the source dataflow fragment from the given [SourceConnection]. This returns a
272/// collection timestamped with the source specific timestamp type.
273fn source_render_operator<'scope, C>(
274    scope: Scope<'scope, C::Time>,
275    config: &RawSourceCreationConfig,
276    source_connection: C,
277    probed_upper_tx: watch::Sender<Option<Probe<C::Time>>>,
278    resume_uppers: impl futures::Stream<Item = Antichain<C::Time>> + 'static,
279    start_signal: impl std::future::Future<Output = ()> + 'static,
280) -> (
281    BTreeMap<GlobalId, StackedCollection<'scope, C::Time, Result<SourceMessage, DataflowError>>>,
282    StreamVec<'scope, C::Time, HealthStatusMessage>,
283    Vec<PressOnDropButton>,
284)
285where
286    C: SourceRender + 'static,
287{
288    let source_id = config.id;
289    let worker_id = config.worker_id;
290
291    let resume_uppers = resume_uppers.inspect(move |upper| {
292        let upper = upper.pretty();
293        trace!(%upper, "timely-{worker_id} source({source_id}) received resume upper");
294    });
295
296    let (exports, health, probe_stream, tokens) =
297        source_connection.render(scope, config, resume_uppers, start_signal);
298
299    let mut export_collections = BTreeMap::new();
300
301    let source_metrics = config.metrics.get_source_metrics(config.id, worker_id);
302
303    // Compute the overall resume upper to report for the ingestion
304    let resume_upper = Antichain::from_iter(
305        config
306            .resume_uppers
307            .values()
308            .flat_map(|f| f.iter().cloned()),
309    );
310    source_metrics
311        .resume_upper
312        .set(mz_persist_client::metrics::encode_ts_metric(&resume_upper));
313
314    let mut health_streams = vec![];
315
316    for (id, export) in exports {
317        let name = format!("SourceGenericStats({})", id);
318        let mut builder = OperatorBuilderRc::new(name, scope.clone());
319
320        let (health_output, derived_health) = builder.new_output();
321        let mut health_output =
322            OutputBuilder::<_, CapacityContainerBuilder<_>>::from(health_output);
323        health_streams.push(derived_health);
324
325        let (output, new_export) = builder.new_output();
326        let mut output = OutputBuilder::<_, CapacityContainerBuilder<_>>::from(output);
327
328        let mut input = builder.new_input(export.inner, Pipeline);
329        export_collections.insert(id, new_export.as_collection());
330
331        let bytes_read_counter = config.metrics.source_defs.bytes_read.clone();
332        let source_statistics = config
333            .statistics
334            .get(&id)
335            .expect("statistics initialized")
336            .clone();
337
338        builder.build(move |mut caps| {
339            let mut health_cap = Some(caps.remove(0));
340
341            move |frontiers| {
342                let mut last_status = None;
343                let mut health_output = health_output.activate();
344
345                if frontiers[0].is_empty() {
346                    health_cap = None;
347                    return;
348                }
349                let health_cap = health_cap.as_mut().unwrap();
350
351                input.for_each(|cap, data| {
352                    for (message, _, _) in data.iter() {
353                        match message {
354                            Ok(message) => {
355                                source_statistics.inc_messages_received_by(1);
356                                let key_len = u64::cast_from(message.key.byte_len());
357                                let value_len = u64::cast_from(message.value.byte_len());
358                                bytes_read_counter.inc_by(key_len + value_len);
359                                source_statistics.inc_bytes_received_by(key_len + value_len);
360                            }
361                            Err(error) => {
362                                // All errors coming into the data stream are definite.
363                                // Downstream consumers of this data will preserve this
364                                // status.
365                                let hint = match error {
366                                    DataflowError::SourceError(e) if e.hint.is_some() => {
367                                        e.hint.as_deref().map(str::to_string)
368                                    }
369                                    _ => Some(
370                                        "retracting the errored value may resume the source"
371                                            .to_string(),
372                                    ),
373                                };
374                                let update = HealthStatusUpdate::stalled(error.to_string(), hint);
375                                let status = HealthStatusMessage {
376                                    id: Some(id),
377                                    namespace: C::STATUS_NAMESPACE.clone(),
378                                    update,
379                                };
380                                if last_status.as_ref() != Some(&status) {
381                                    last_status = Some(status.clone());
382                                    health_output.session(&health_cap).give(status);
383                                }
384                            }
385                        }
386                    }
387                    let mut output = output.activate();
388                    output.session(&cap).give_container(data);
389                });
390            }
391        });
392    }
393
394    // Broadcasting does more work than necessary, which would be to exchange the probes to the
395    // worker that will be the one minting the bindings but we'd have to thread this information
396    // through and couple the two functions enough that it's not worth the optimization (I think).
397    // Use `InspectCore::inspect_container` instead of `Inspect::inspect`.
398    // `Inspect` carries a `where for<'a> &'a C: IntoIterator` bound, and on
399    // macOS the solver can satisfy that bound by chasing objc2's
400    // `&Retained<T>: IntoIterator` blanket impl into an endless
401    // `Retained<Retained<…>>` chain, overflowing the recursion limit.
402    // `InspectCore` has no such bound, so the cascade never starts. We
403    // iterate the container by hand to recover the per-item callback.
404    probe_stream.broadcast().inspect_container(move |event| {
405        if let Ok((_, data)) = event {
406            for probe in data {
407                // We don't care if the receiver is gone
408                let _ = probed_upper_tx.send(Some(probe.clone()));
409            }
410        }
411    });
412
413    (
414        export_collections,
415        health.concatenate_flatten::<_, CapacityContainerBuilder<_>>(health_streams),
416        tokens,
417    )
418}
419
420/// Mints new contents for the remap shard based on summaries about the source
421/// upper it receives from the raw reader operators.
422///
423/// Only one worker will be active and write to the remap shard. All source
424/// upper summaries will be exchanged to it.
425fn remap_operator<'scope, FromTime>(
426    scope: Scope<'scope, mz_repr::Timestamp>,
427    storage_state: &crate::storage_state::StorageState,
428    config: RawSourceCreationConfig,
429    mut probed_upper: watch::Receiver<Option<Probe<FromTime>>>,
430    remap_relation_desc: RelationDesc,
431) -> (
432    VecCollection<'scope, mz_repr::Timestamp, FromTime, Diff>,
433    PressOnDropButton,
434)
435where
436    FromTime: SourceTimestamp,
437{
438    let RawSourceCreationConfig {
439        name,
440        id,
441        source_exports: _,
442        worker_id,
443        worker_count,
444        timestamp_interval: _,
445        remap_metadata,
446        as_of,
447        resume_uppers: _,
448        source_resume_uppers: _,
449        metrics: _,
450        now_fn,
451        persist_clients,
452        statistics: _,
453        shared_remap_upper,
454        config: _,
455        remap_collection_id,
456        busy_signal: _,
457    } = config;
458
459    let read_only_rx = storage_state.read_only_rx.clone();
460    let error_handler = storage_state.error_handler("remap_operator", id);
461
462    let chosen_worker = usize::cast_from(id.hashed() % u64::cast_from(worker_count));
463    let active_worker = chosen_worker == worker_id;
464
465    let operator_name = format!("remap({})", id);
466    let mut remap_op = AsyncOperatorBuilder::new(operator_name, scope.clone());
467    let (remap_output, remap_stream) = remap_op.new_output::<CapacityContainerBuilder<_>>();
468
469    let button = remap_op.build(move |capabilities| async move {
470        if !active_worker {
471            // This worker is not writing, so make sure it's "taken out" of the
472            // calculation by advancing to the empty frontier.
473            shared_remap_upper.borrow_mut().clear();
474            return;
475        }
476
477        let mut cap_set = CapabilitySet::from_elem(capabilities.into_element());
478
479        let remap_handle = crate::source::reclock::compat::PersistHandle::<FromTime, _>::new(
480            Arc::clone(&persist_clients),
481            read_only_rx,
482            remap_metadata.clone(),
483            as_of.clone(),
484            shared_remap_upper,
485            id,
486            "remap",
487            worker_id,
488            worker_count,
489            remap_relation_desc,
490            remap_collection_id,
491        )
492        .await;
493
494        let remap_handle = match remap_handle {
495            Ok(handle) => handle,
496            Err(e) => {
497                error_handler
498                    .report_and_stop(
499                        e.context(format!("Failed to create remap handle for source {name}")),
500                    )
501                    .await
502            }
503        };
504
505        let (mut timestamper, mut initial_batch) = ReclockOperator::new(remap_handle).await;
506
507        // Emit initial snapshot of the remap_shard, bootstrapping
508        // downstream reclock operators.
509        trace!(
510            "timely-{worker_id} remap({id}) emitting remap snapshot: trace_updates={:?}",
511            &initial_batch.updates
512        );
513
514        let cap = cap_set.delayed(cap_set.first().unwrap());
515        remap_output.give_container(&cap, &mut initial_batch.updates);
516        drop(cap);
517        cap_set.downgrade(initial_batch.upper);
518
519        let mut prev_probe_ts: Option<mz_repr::Timestamp> = None;
520
521        while !cap_set.is_empty() {
522            // We only mint bindings after a successful probe.
523            let new_probe = probed_upper
524                .wait_for(|new_probe| match (prev_probe_ts, new_probe) {
525                    (None, Some(_)) => true,
526                    (Some(prev_ts), Some(new)) => prev_ts < new.probe_ts,
527                    _ => false,
528                })
529                .await
530                .map(|probe| (*probe).clone())
531                .unwrap_or_else(|_| {
532                    Some(Probe {
533                        probe_ts: now_fn().into(),
534                        upstream_frontier: Antichain::new(),
535                    })
536                });
537
538            let probe = new_probe.expect("known to be Some");
539            prev_probe_ts = Some(probe.probe_ts);
540
541            let binding_ts = probe.probe_ts;
542            let cur_source_upper = probe.upstream_frontier;
543
544            let new_into_upper = Antichain::from_elem(binding_ts.step_forward());
545
546            let mut remap_trace_batch = timestamper
547                .mint(binding_ts, new_into_upper, cur_source_upper.borrow())
548                .await;
549
550            trace!(
551                "timely-{worker_id} remap({id}) minted new bindings: \
552                updates={:?} \
553                source_upper={} \
554                trace_upper={}",
555                &remap_trace_batch.updates,
556                cur_source_upper.pretty(),
557                remap_trace_batch.upper.pretty()
558            );
559
560            let cap = cap_set.delayed(cap_set.first().unwrap());
561            remap_output.give_container(&cap, &mut remap_trace_batch.updates);
562            cap_set.downgrade(remap_trace_batch.upper);
563        }
564    });
565
566    (remap_stream.as_collection(), button.press_on_drop())
567}
568
569/// Reclocks an `IntoTime` frontier stream into a `FromTime` frontier stream. This is used for the
570/// virtual (through persist) feedback edge so that we convert the `IntoTime` resumption frontier
571/// into the `FromTime` frontier that is used with the source's `OffsetCommiter`.
572fn reclock_committed_upper<'scope, T, FromTime>(
573    bindings: VecCollection<'scope, T, FromTime, Diff>,
574    as_of: Antichain<T>,
575    committed_upper: StreamVec<'scope, T, ()>,
576    id: GlobalId,
577    metrics: Arc<SourceMetrics>,
578) -> impl futures::stream::Stream<Item = Antichain<FromTime>> + 'static
579where
580    T: Timestamp + Lattice + TotalOrder,
581    FromTime: SourceTimestamp,
582{
583    let (tx, rx) = watch::channel(Antichain::from_elem(FromTime::minimum()));
584    let scope = bindings.scope().clone();
585
586    let name = format!("ReclockCommitUpper({id})");
587    let mut builder = OperatorBuilderRc::new(name, scope);
588
589    let mut bindings = builder.new_input(bindings.inner.clone(), Pipeline);
590    let _ = builder.new_input(committed_upper.clone(), Pipeline);
591
592    builder.build(move |_| {
593        // Remap bindings beyond the upper
594        use timely::progress::ChangeBatch;
595        let mut accepted_times: ChangeBatch<(T, FromTime)> = ChangeBatch::new();
596        // The upper frontier of the bindings
597        let mut upper = Antichain::from_elem(Timestamp::minimum());
598        // Remap bindings not beyond upper
599        let mut ready_times = VecDeque::new();
600        let mut source_upper = MutableAntichain::new();
601
602        move |frontiers| {
603            // Accept new bindings
604            bindings.for_each(|_, data| {
605                accepted_times.extend(data.drain(..).map(|(from, mut into, diff)| {
606                    into.advance_by(as_of.borrow());
607                    ((into, from), diff.into_inner())
608                }));
609            });
610            // Extract ready bindings
611            let new_upper = frontiers[0].frontier();
612            if PartialOrder::less_than(&upper.borrow(), &new_upper) {
613                upper = new_upper.to_owned();
614                // Drain consolidated accepted times not greater or equal to `upper` into `ready_times`.
615                // Retain accepted times greater or equal to `upper` in
616                let mut pending_times = std::mem::take(&mut accepted_times).into_inner();
617                // These should already be sorted, as part of `.into_inner()`, but sort defensively in case.
618                pending_times.sort_unstable_by(|a, b| a.0.cmp(&b.0));
619                for ((into, from), diff) in pending_times.drain(..) {
620                    if !upper.less_equal(&into) {
621                        ready_times.push_back((from, into, diff));
622                    } else {
623                        accepted_times.update((into, from), diff);
624                    }
625                }
626            }
627
628            // The received times only accumulate correctly for times beyond the as_of.
629            if as_of.iter().all(|t| !upper.less_equal(t)) {
630                let committed_upper = frontiers[1].frontier();
631                if as_of.iter().all(|t| !committed_upper.less_equal(t)) {
632                    // We have committed this source up until `committed_upper`. Because we have
633                    // required that IntoTime is a total order this will be either a singleton set
634                    // or the empty set.
635                    //
636                    // * Case 1: committed_upper is the empty set {}
637                    //
638                    // There won't be any future IntoTime timestamps that we will produce so we can
639                    // provide feedback to the source that it can forget about everything.
640                    //
641                    // * Case 2: committed_upper is a singleton set {t_next}
642                    //
643                    // We know that t_next cannot be the minimum timestamp because we have required
644                    // that all times of the as_of frontier are not beyond some time of
645                    // committed_upper. Therefore t_next has a predecessor timestamp t_prev.
646                    //
647                    // We don't know what remap[t_next] is yet, but we do know that we will have to
648                    // emit all source updates `u: remap[t_prev] <= time(u) <= remap[t_next]`.
649                    // Since `t_next` is the minimum undetermined timestamp and we know that t1 <=
650                    // t2 => remap[t1] <= remap[t2] we know that we will never need any source
651                    // updates `u: !(remap[t_prev] <= time(u))`.
652                    //
653                    // Therefore we can provide feedback to the source that it can forget about any
654                    // updates that are not beyond remap[t_prev].
655                    //
656                    // Important: We are *NOT* saying that the source can *compact* its data using
657                    // remap[t_prev] as the compaction frontier. If the source were to compact its
658                    // collection to remap[t_prev] we would lose the distinction between updates
659                    // that happened *at* t_prev versus updates that happened ealier and were
660                    // advanced to t_prev. If the source needs to communicate a compaction frontier
661                    // upstream then the specific source implementation needs to further adjust the
662                    // reclocked committed_upper and calculate a suitable compaction frontier in
663                    // the same way we adjust uppers of collections in the controller with the
664                    // LagWriteFrontier read policy.
665                    //
666                    // == What about IntoTime times that are general lattices?
667                    //
668                    // Reversing the upper for a general lattice is much more involved but it boils
669                    // down to computing the meet of all the times in `committed_upper` and then
670                    // treating that as `t_next` (I think). Until we need to deal with that though
671                    // we can just assume TotalOrder.
672                    let reclocked_upper = match committed_upper.as_option() {
673                        Some(t_next) => {
674                            let idx = ready_times.partition_point(|(_, t, _)| t < t_next);
675                            let updates = ready_times
676                                .drain(0..idx)
677                                .map(|(from_time, _, diff)| (from_time, diff));
678                            source_upper.update_iter(updates);
679                            // At this point source_upper contains all updates that are less than
680                            // t_next, which is equal to remap[t_prev]
681                            source_upper.frontier().to_owned()
682                        }
683                        None => Antichain::new(),
684                    };
685                    tx.send_replace(reclocked_upper);
686                }
687            }
688
689            metrics
690                .commit_upper_accepted_times
691                .set(u64::cast_from(accepted_times.len()));
692            metrics
693                .commit_upper_ready_times
694                .set(u64::cast_from(ready_times.len()));
695        }
696    });
697
698    WatchStream::from_changes(rx)
699}