1#![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#[derive(Clone)]
77pub struct RawSourceCreationConfig {
78 pub name: String,
80 pub id: GlobalId,
82 pub source_exports: BTreeMap<GlobalId, SourceExport<CollectionMetadata>>,
84 pub worker_id: usize,
86 pub worker_count: usize,
88 pub timestamp_interval: Duration,
91 pub now_fn: NowFn,
93 pub metrics: StorageMetrics,
95 pub as_of: Antichain<mz_repr::Timestamp>,
97 pub resume_uppers: BTreeMap<GlobalId, Antichain<mz_repr::Timestamp>>,
100 pub source_resume_uppers: BTreeMap<GlobalId, Vec<Row>>,
107 pub persist_clients: Arc<PersistClientCache>,
109 pub statistics: BTreeMap<GlobalId, SourceStatistics>,
111 pub shared_remap_upper: Rc<RefCell<Antichain<mz_repr::Timestamp>>>,
113 pub config: StorageConfiguration,
115 pub remap_collection_id: GlobalId,
117 pub remap_metadata: CollectionMetadata,
119 pub busy_signal: Arc<Semaphore>,
122}
123
124#[derive(Clone)]
127pub struct SourceExportCreationConfig {
128 pub id: GlobalId,
130 pub worker_id: usize,
132 pub metrics: StorageMetrics,
134 pub source_statistics: SourceStatistics,
136}
137
138impl RawSourceCreationConfig {
139 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 pub fn responsible_for<P: Hash>(&self, partition: P) -> bool {
149 self.responsible_worker(partition) == self.worker_id
150 }
151}
152
153pub 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 let remap_collection = remap_collection.inner.broadcast().as_collection();
212 tokens.push(remap_token);
213
214 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
271fn 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 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 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 probe_stream.broadcast().inspect_container(move |event| {
405 if let Ok((_, data)) = event {
406 for probe in data {
407 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
420fn 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 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 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 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
569fn 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 use timely::progress::ChangeBatch;
595 let mut accepted_times: ChangeBatch<(T, FromTime)> = ChangeBatch::new();
596 let mut upper = Antichain::from_elem(Timestamp::minimum());
598 let mut ready_times = VecDeque::new();
600 let mut source_upper = MutableAntichain::new();
601
602 move |frontiers| {
603 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 let new_upper = frontiers[0].frontier();
612 if PartialOrder::less_than(&upper.borrow(), &new_upper) {
613 upper = new_upper.to_owned();
614 let mut pending_times = std::mem::take(&mut accepted_times).into_inner();
617 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 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 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 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}