1use std::any::Any;
9use std::cell::RefCell;
10use std::cmp::Ordering;
11use std::collections::{BTreeMap, BTreeSet, VecDeque};
12use std::num::NonZeroUsize;
13use std::rc::Rc;
14use std::sync::Arc;
15use std::time::{Duration, Instant};
16
17use differential_dataflow::Hashable;
18use differential_dataflow::lattice::Lattice;
19use differential_dataflow::trace::TraceReader;
20use mz_compute_client::logging::LoggingConfig;
21use mz_compute_client::protocol::command::{
22 ComputeCommand, ComputeParameters, InstanceConfig, Peek, PeekTarget,
23};
24use mz_compute_client::protocol::history::ComputeCommandHistory;
25use mz_compute_client::protocol::response::{
26 ComputeResponse, CopyToResponse, FrontiersResponse, PeekError, PeekResponse, SubscribeResponse,
27};
28use mz_compute_types::dataflows::DataflowDescription;
29use mz_compute_types::dyncfgs::{
30 ENABLE_PEEK_RESPONSE_STASH, ENABLE_PEEK_ROW_ITERATION_LIMIT, PEEK_RESPONSE_STASH_BATCH_BYTES,
31 PEEK_RESPONSE_STASH_THRESHOLD_BYTES, PEEK_ROW_ITERATION_LIMIT,
32};
33use mz_compute_types::plan::render_plan::RenderPlan;
34use mz_dyncfg::{ConfigSet, ConfigValHandle};
35use mz_expr::SafeMfpPlan;
36use mz_expr::row::RowCollection;
37use mz_ore::cast::{CastFrom, CastLossy};
38use mz_ore::collections::CollectionExt;
39use mz_ore::metrics::{MetricsRegistry, UIntGauge};
40use mz_ore::now::EpochMillis;
41use mz_ore::soft_panic_or_log;
42use mz_ore::task::AbortOnDropHandle;
43use mz_ore::tracing::{OpenTelemetryContext, TracingHandle};
44use mz_persist_client::Diagnostics;
45use mz_persist_client::cache::PersistClientCache;
46use mz_persist_client::cfg::USE_CRITICAL_SINCE_SNAPSHOT;
47use mz_persist_client::read::ReadHandle;
48use mz_persist_types::PersistLocation;
49use mz_persist_types::codec_impls::UnitSchema;
50use mz_repr::{DatumVec, GlobalId, Row, RowArena, Timestamp};
51use mz_storage_operators::stats::StatsCursor;
52use mz_storage_types::StorageDiff;
53use mz_storage_types::controller::CollectionMetadata;
54use mz_storage_types::dyncfgs::ORE_OVERFLOWING_BEHAVIOR;
55use mz_storage_types::sources::SourceData;
56use mz_storage_types::time_dependence::TimeDependence;
57use mz_txn_wal::operator::TxnsContext;
58use mz_txn_wal::txn_cache::TxnsCache;
59use timely::dataflow::operators::probe;
60use timely::order::PartialOrder;
61use timely::progress::frontier::Antichain;
62use timely::worker::Worker as TimelyWorker;
63use tokio::sync::{oneshot, watch};
64use tracing::{Level, debug, error, info, span, trace, warn};
65use uuid::Uuid;
66
67use crate::arrangement::manager::{TraceBundle, TraceManager};
68use crate::compute_state::peek_budget::InlineBudget;
69use crate::compute_state::peek_metrics::{IndexPeekMetrics, PeekWalkMetrics};
70pub(crate) use crate::compute_state::peek_offload::PeekPermits;
71use crate::compute_state::peek_offload::{OffloadConfig, OffloadedPeek};
72use crate::compute_state::peek_scan::{
73 IndexPeekScan, PeekScan, ScanOutcome, StashBounds, entry_byte_len, rows_response,
74};
75use crate::logging;
76use crate::logging::compute::{CollectionLogging, ComputeEvent, PeekEvent};
77use crate::logging::initialize::LoggingTraces;
78use crate::metrics::{CollectionMetrics, WorkerMetrics};
79use crate::render::{LinearJoinSpec, StartSignal};
80use crate::server::{ComputeInstanceContext, ResponseSender};
81
82mod error_scan;
83mod peek_budget;
84mod peek_metrics;
85mod peek_offload;
86mod peek_result_iterator;
87mod peek_scan;
88mod peek_stash;
89
90#[derive(Clone, Debug)]
95struct PeekRowIterationConfig {
96 enabled: ConfigValHandle<bool>,
97 limit: ConfigValHandle<usize>,
98}
99
100impl PeekRowIterationConfig {
101 fn new(config: &ConfigSet) -> Self {
102 Self {
103 enabled: ENABLE_PEEK_ROW_ITERATION_LIMIT.handle(config),
104 limit: PEEK_ROW_ITERATION_LIMIT.handle(config),
105 }
106 }
107
108 fn current_limit(&self) -> Option<usize> {
109 self.enabled.get().then(|| self.limit.get())
110 }
111}
112
113#[derive(Debug)]
122pub(crate) struct PeekRowIterationTracker {
123 limit: Option<usize>,
124 rows_iterated: usize,
125}
126
127impl PeekRowIterationTracker {
128 fn new(limit: Option<usize>, rows_iterated: usize) -> Self {
129 Self {
130 limit,
131 rows_iterated,
132 }
133 }
134
135 fn set_limit(&mut self, limit: Option<usize>) {
140 self.limit = limit;
141 }
142
143 fn rows_iterated(&self) -> usize {
144 self.rows_iterated
145 }
146
147 fn add_rows_iterated(&mut self, rows_iterated: usize) {
152 self.rows_iterated = self.rows_iterated.saturating_add(rows_iterated);
153 }
154
155 fn track_next(&mut self) -> Result<(), PeekError> {
156 if let Some(limit) = self.limit
157 && self.rows_iterated >= limit
158 {
159 return Err(PeekError::RowIterationLimitExceeded { limit });
160 }
161
162 self.rows_iterated = self.rows_iterated.saturating_add(1);
163 Ok(())
164 }
165}
166
167fn peek_row_iteration_limit(config: &ConfigSet) -> Option<usize> {
168 ENABLE_PEEK_ROW_ITERATION_LIMIT
169 .get(config)
170 .then(|| PEEK_ROW_ITERATION_LIMIT.get(config))
171}
172
173pub struct ComputeState {
178 pub collections: BTreeMap<GlobalId, CollectionState>,
187 pub traces: TraceManager,
189 pub subscribe_response_buffer: Rc<RefCell<Vec<(GlobalId, SubscribeResponse)>>>,
194 pub copy_to_response_buffer: Rc<RefCell<Vec<(GlobalId, CopyToResponse)>>>,
199 pub queued_peeks: VecDeque<IndexPeek>,
205 pub pending_peeks: VecDeque<PendingPeek>,
210 pub peek_stash_persist_location: Option<PersistLocation>,
212 pub compute_logger: Option<logging::compute::Logger>,
214 pub persist_clients: Arc<PersistClientCache>,
217 pub txns_ctx: TxnsContext,
219 pub command_history: ComputeCommandHistory<UIntGauge>,
221 max_result_size: u64,
223 pub linear_join_spec: LinearJoinSpec,
225 pub metrics: WorkerMetrics,
227 tracing_handle: Arc<TracingHandle>,
229 pub context: ComputeInstanceContext,
231 pub worker_config: Rc<ConfigSet>,
245
246 pub metrics_registry: MetricsRegistry,
248
249 pub workers_per_process: usize,
251
252 pub peek_permits: Arc<PeekPermits>,
255
256 peek_walk_metrics: PeekWalkMetrics,
261
262 peek_budget: InlineBudget,
268
269 peek_passed_over: bool,
274
275 suspended_collections: BTreeMap<GlobalId, Rc<dyn Any>>,
281
282 pub server_maintenance_interval: Duration,
285
286 pub init_system_time: EpochMillis,
290
291 pub replica_expiration: Antichain<Timestamp>,
295}
296
297impl ComputeState {
298 pub(crate) fn peeks_awaiting_turn(&self) -> bool {
303 self.peek_passed_over
304 }
305
306 pub fn new(
308 persist_clients: Arc<PersistClientCache>,
309 txns_ctx: TxnsContext,
310 metrics: WorkerMetrics,
311 tracing_handle: Arc<TracingHandle>,
312 context: ComputeInstanceContext,
313 metrics_registry: MetricsRegistry,
314 workers_per_process: usize,
315 peek_permits: Arc<PeekPermits>,
316 ) -> Self {
317 let worker_config: Rc<ConfigSet> = mz_dyncfgs::all_dyncfgs().into();
318 let traces = TraceManager::new(metrics.clone());
319 let command_history = ComputeCommandHistory::new(metrics.for_history());
320 let peek_walk_metrics = PeekWalkMetrics::new(&metrics);
321 let peek_budget = InlineBudget::new(&worker_config);
322
323 Self {
324 collections: Default::default(),
325 traces,
326 subscribe_response_buffer: Default::default(),
327 copy_to_response_buffer: Default::default(),
328 queued_peeks: Default::default(),
329 pending_peeks: Default::default(),
330 peek_stash_persist_location: None,
331 compute_logger: None,
332 persist_clients,
333 txns_ctx,
334 command_history,
335 max_result_size: u64::MAX,
336 linear_join_spec: Default::default(),
337 metrics,
338 tracing_handle,
339 context,
340 worker_config,
341 metrics_registry,
342 workers_per_process,
343 peek_permits,
344 peek_walk_metrics,
345 peek_budget,
346 peek_passed_over: false,
347 suspended_collections: Default::default(),
348 server_maintenance_interval: Duration::ZERO,
349 init_system_time: mz_ore::now::SYSTEM_TIME(),
350 replica_expiration: Antichain::default(),
351 }
352 }
353
354 pub fn expect_collection_mut(&mut self, id: GlobalId) -> &mut CollectionState {
358 self.collections
359 .get_mut(&id)
360 .expect("collection must exist")
361 }
362
363 pub fn input_probe_for(
369 &mut self,
370 input_id: GlobalId,
371 collection_ids: impl Iterator<Item = GlobalId>,
372 ) -> probe::Handle<Timestamp> {
373 let probe = probe::Handle::default();
374 for id in collection_ids {
375 if let Some(collection) = self.collections.get_mut(&id) {
376 collection.input_probes.insert(input_id, probe.clone());
377 }
378 }
379 probe
380 }
381
382 fn apply_worker_config(&mut self) {
384 use mz_compute_types::dyncfgs::*;
385
386 let config = &self.worker_config;
387
388 self.linear_join_spec = LinearJoinSpec::from_config(config);
389
390 if ENABLE_LGALLOC.get(config) {
391 if let Some(path) = &self.context.scratch_directory {
392 let clear_bytes = LGALLOC_SLOW_CLEAR_BYTES.get(config);
393 let eager_return = ENABLE_LGALLOC_EAGER_RECLAMATION.get(config);
394 let file_growth_dampener = LGALLOC_FILE_GROWTH_DAMPENER.get(config);
395 let interval = LGALLOC_BACKGROUND_INTERVAL.get(config);
396 let local_buffer_bytes = LGALLOC_LOCAL_BUFFER_BYTES.get(config);
397 info!(
398 ?path,
399 backgrund_interval=?interval,
400 clear_bytes,
401 eager_return,
402 file_growth_dampener,
403 local_buffer_bytes,
404 "enabling lgalloc"
405 );
406 let background_worker_config = lgalloc::BackgroundWorkerConfig {
407 interval,
408 clear_bytes,
409 };
410 lgalloc::lgalloc_set_config(
411 lgalloc::LgAlloc::new()
412 .enable()
413 .with_path(path.clone())
414 .with_background_config(background_worker_config)
415 .eager_return(eager_return)
416 .file_growth_dampener(file_growth_dampener)
417 .local_buffer_bytes(local_buffer_bytes),
418 );
419 } else {
420 debug!("not enabling lgalloc, scratch directory not specified");
421 }
422 } else {
423 info!("disabling lgalloc");
424 lgalloc::lgalloc_set_config(lgalloc::LgAlloc::new().disable());
425 }
426
427 #[cfg(unix)]
433 if let Some(path) = &self.context.scratch_directory {
434 mz_ore::pager::set_scratch_dir(path.clone());
435 mz_ore::pager::set_backend(mz_ore::pager::Backend::File);
436 } else {
437 mz_ore::pager::set_backend(mz_ore::pager::Backend::Swap);
438 }
439
440 crate::memory_limiter::apply_limiter_config(config);
441
442 mz_ore::region::ENABLE_LGALLOC_REGION.store(
443 ENABLE_COLUMNATION_LGALLOC.get(config),
444 std::sync::atomic::Ordering::Relaxed,
445 );
446
447 {
465 use mz_ore::pager::Backend;
466 use mz_timely_util::column_pager::{Codec, apply_tiered_config};
467
468 let enabled = ENABLE_COLUMN_PAGED_BATCHER_SPILL.get(config);
469 let codec = COLUMN_PAGED_BATCHER_LZ4.get(config).then_some(Codec::Lz4);
470 let swap_pageout = COLUMN_PAGED_BATCHER_SWAP_PAGEOUT.get(config);
471
472 const MIB: usize = 1024 * 1024;
477 const DEFAULT_MEM_LIMIT: usize = 4 * 1024 * MIB;
478 let mem_limit = crate::memory_limiter::get_memory_limit().unwrap_or(DEFAULT_MEM_LIMIT);
479 let fraction = COLUMN_PAGED_BATCHER_BUDGET_FRACTION.get(config).max(0.0);
480 let total = usize::cast_lossy(f64::cast_lossy(mem_limit) * fraction).max(128 * MIB);
481
482 let backend = if self.context.scratch_directory.is_some() {
483 Backend::File
484 } else {
485 Backend::Swap
486 };
487
488 debug!(
489 enabled,
490 ?backend,
491 ?codec,
492 swap_pageout,
493 fraction,
494 mem_limit,
495 budget_bytes = total,
496 "column-paged batcher: applying tiered config",
497 );
498 apply_tiered_config(enabled, total, backend, codec, swap_pageout);
499 }
500
501 {
518 use mz_timely_util::pool_config::{PoolPagerConfig, apply_pool_config};
519
520 let compute_spill = ENABLE_COLUMN_PAGED_BATCHER_SPILL.get(config);
521 let storage_spill = mz_storage_types::dyncfgs::ENABLE_UPSERT_PAGED_SPILL.get(config);
522 let sink_spill = ENABLE_CORRECTION_V2_SPILL.get(config);
523 mz_timely_util::columnar::chunk::set_compute_spill_enabled(compute_spill);
529 mz_timely_util::columnar::chunk::set_sink_spill_enabled(sink_spill);
530 if !(compute_spill || storage_spill || sink_spill) {
531 debug!("chunk spill: gates off, leaving the buffer pool uninstalled");
532 } else {
533 let spill_threads = COLUMN_PAGED_BATCHER_SPILL_WORKER_COUNT.get(config);
534 let eager_backing = COLUMN_PAGED_BATCHER_EAGER_BACKING.get(config);
535
536 const MIB: usize = 1024 * 1024;
543 const DEFAULT_RAM: usize = 4 * 1024 * MIB;
544 let ram = mz_ore::memory::physical_memory_bytes().unwrap_or(DEFAULT_RAM);
545 let of_ram =
546 |fraction: f64| usize::cast_lossy(f64::cast_lossy(ram) * fraction.max(0.0));
547 let fraction = COLUMN_PAGED_BATCHER_BUDGET_FRACTION.get(config);
548 let total = of_ram(fraction).max(128 * MIB);
549 let rss_target = of_ram(COLUMN_PAGED_BATCHER_POOL_RSS_TARGET_FRACTION.get(config));
554
555 let applied = apply_pool_config(PoolPagerConfig {
556 budget_bytes: total,
557 spill_threads,
558 eager_backing,
559 rss_target_bytes: rss_target,
560 });
561 if applied {
562 info!(
563 compute_spill,
564 storage_spill,
565 fraction,
566 ram,
567 budget_bytes = total,
568 spill_threads,
569 eager_backing,
570 rss_target_bytes = rss_target,
571 "chunk spill: applying buffer-pool config",
572 );
573 } else {
574 warn!("chunk spill: buffer pool unavailable; chunks stay resident");
575 }
576 }
577
578 let compress_min_depth =
582 u8::try_from(COLUMN_CHUNK_COMPRESS_MIN_DEPTH.get(config)).unwrap_or(u8::MAX);
583 mz_timely_util::columnar::chunk::set_compress_min_depth(compress_min_depth);
584 }
585
586 self.server_maintenance_interval = COMPUTE_SERVER_MAINTENANCE_INTERVAL.get(config);
589
590 let overflowing_behavior = ORE_OVERFLOWING_BEHAVIOR.get(config);
591 match overflowing_behavior.parse() {
592 Ok(behavior) => mz_ore::overflowing::set_behavior(behavior),
593 Err(err) => {
594 error!(
595 err,
596 overflowing_behavior, "Invalid value for ore_overflowing_behavior"
597 );
598 }
599 }
600 }
601
602 pub fn apply_expiration_offset(&mut self, offset: Duration) {
608 if self.replica_expiration.is_empty() {
609 let offset: EpochMillis = offset
610 .as_millis()
611 .try_into()
612 .expect("duration must fit within u64");
613 let replica_expiration_millis = self.init_system_time + offset;
614 let replica_expiration = Timestamp::from(replica_expiration_millis);
615
616 info!(
617 offset = %offset,
618 replica_expiration_millis = %replica_expiration_millis,
619 replica_expiration_utc = %mz_ore::now::to_datetime(replica_expiration_millis),
620 "setting replica expiration",
621 );
622 self.replica_expiration = Antichain::from_elem(replica_expiration);
623
624 self.metrics
626 .replica_expiration_timestamp_seconds
627 .set(replica_expiration.into());
628 }
629 }
630
631 pub fn dataflow_max_inflight_bytes(&self) -> Option<usize> {
634 use mz_compute_types::dyncfgs::{
635 DATAFLOW_MAX_INFLIGHT_BYTES, DATAFLOW_MAX_INFLIGHT_BYTES_CC,
636 };
637
638 if self.persist_clients.cfg.is_cc_active {
639 DATAFLOW_MAX_INFLIGHT_BYTES_CC.get(&self.worker_config)
640 } else {
641 DATAFLOW_MAX_INFLIGHT_BYTES.get(&self.worker_config)
642 }
643 }
644}
645
646pub(crate) struct ActiveComputeState<'a> {
648 pub timely_worker: &'a mut TimelyWorker,
650 pub compute_state: &'a mut ComputeState,
652 pub response_tx: &'a mut ResponseSender,
654}
655
656pub struct SinkToken(#[allow(dead_code)] Box<dyn Any>);
658
659impl SinkToken {
660 pub fn new(t: Box<dyn Any>) -> Self {
662 Self(t)
663 }
664}
665
666impl<'a> ActiveComputeState<'a> {
667 #[mz_ore::instrument(level = "debug")]
669 pub fn handle_compute_command(&mut self, cmd: ComputeCommand) {
670 use ComputeCommand::*;
671
672 self.compute_state.command_history.push(cmd.clone());
673
674 let timer = self
676 .compute_state
677 .metrics
678 .handle_command_duration_seconds
679 .for_command(&cmd)
680 .start_timer();
681
682 match cmd {
683 Hello { .. } => panic!("Hello must be captured before"),
684 CreateInstance(instance_config) => self.handle_create_instance(*instance_config),
685 InitializationComplete => (),
686 UpdateConfiguration(params) => self.handle_update_configuration(*params),
687 CreateDataflow(dataflow) => self.handle_create_dataflow(*dataflow),
688 Schedule(id) => self.handle_schedule(id),
689 AllowCompaction { id, frontier } => self.handle_allow_compaction(id, frontier),
690 Peek(peek) => {
691 peek.otel_ctx.attach_as_parent();
692 self.handle_peek(*peek)
693 }
694 CancelPeek { uuid } => self.handle_cancel_peek(uuid),
695 AllowWrites(id) => {
696 self.handle_allow_writes(id);
697 }
698 }
699
700 timer.observe_duration();
701 }
702
703 fn handle_create_instance(&mut self, config: InstanceConfig) {
704 config
709 .initial_config
710 .apply(&self.compute_state.worker_config);
711
712 self.compute_state.apply_worker_config();
714
715 mz_row_spine::DICTIONARY_COMPRESSION.store(
721 config.arrangement_dictionary_compression,
722 std::sync::atomic::Ordering::Relaxed,
723 );
724
725 if let Some(offset) = config.expiration_offset {
726 self.compute_state.apply_expiration_offset(offset);
727 }
728
729 self.initialize_logging(config.logging);
730
731 self.compute_state.peek_stash_persist_location = Some(config.peek_stash_persist_location);
732 }
733
734 fn handle_update_configuration(&mut self, params: ComputeParameters) {
735 debug!("Applying configuration update: {params:?}");
736
737 let ComputeParameters {
738 workload_class,
739 max_result_size,
740 tracing,
741 grpc_client: _grpc_client,
742 dyncfg_updates,
743 } = params;
744
745 if let Some(v) = workload_class {
746 self.compute_state.metrics.set_workload_class(v);
747 }
748 if let Some(v) = max_result_size {
749 self.compute_state.max_result_size = v;
750 }
751
752 tracing.apply(self.compute_state.tracing_handle.as_ref());
753
754 dyncfg_updates.apply(&self.compute_state.worker_config);
755 self.compute_state
756 .persist_clients
757 .cfg()
758 .apply_from(&dyncfg_updates);
759
760 mz_metrics::update_dyncfg(&dyncfg_updates);
764
765 self.compute_state.apply_worker_config();
766 }
767
768 fn handle_create_dataflow(
769 &mut self,
770 dataflow: DataflowDescription<RenderPlan, CollectionMetadata>,
771 ) {
772 let dataflow_index = Rc::new(self.timely_worker.next_dataflow_index());
773 let as_of = dataflow.as_of.clone().unwrap();
774
775 let dataflow_expiration = dataflow
776 .time_dependence
777 .as_ref()
778 .map(|time_dependence| {
779 self.determine_dataflow_expiration(time_dependence, &dataflow.until)
780 })
781 .unwrap_or_default();
782
783 let until = dataflow.until.meet(&dataflow_expiration);
785
786 if dataflow.is_transient() {
787 debug!(
788 name = %dataflow.debug_name,
789 import_ids = %dataflow.display_import_ids(),
790 export_ids = %dataflow.display_export_ids(),
791 as_of = ?as_of.elements(),
792 time_dependence = ?dataflow.time_dependence,
793 expiration = ?dataflow_expiration.elements(),
794 expiration_datetime = ?dataflow_expiration
795 .as_option()
796 .map(|t| mz_ore::now::to_datetime(t.into())),
797 plan_until = ?dataflow.until.elements(),
798 until = ?until.elements(),
799 "creating dataflow",
800 );
801 } else {
802 info!(
803 name = %dataflow.debug_name,
804 import_ids = %dataflow.display_import_ids(),
805 export_ids = %dataflow.display_export_ids(),
806 as_of = ?as_of.elements(),
807 time_dependence = ?dataflow.time_dependence,
808 expiration = ?dataflow_expiration.elements(),
809 expiration_datetime = ?dataflow_expiration
810 .as_option()
811 .map(|t| mz_ore::now::to_datetime(t.into())),
812 plan_until = ?dataflow.until.elements(),
813 until = ?until.elements(),
814 "creating dataflow",
815 );
816 };
817
818 let subscribe_copy_ids: BTreeSet<_> = dataflow
819 .subscribe_ids()
820 .chain(dataflow.copy_to_ids())
821 .collect();
822
823 let starts_immediately = dataflow.import_ids().next().is_none();
832
833 for object_id in dataflow.export_ids() {
835 let is_subscribe_or_copy = subscribe_copy_ids.contains(&object_id);
836 let metrics = self.compute_state.metrics.for_collection(object_id);
837 let mut collection = CollectionState::new(
838 Rc::clone(&dataflow_index),
839 is_subscribe_or_copy,
840 as_of.clone(),
841 metrics,
842 );
843
844 if let Some(logger) = self.compute_state.compute_logger.clone() {
845 let logging = CollectionLogging::new(
846 object_id,
847 logger,
848 *dataflow_index,
849 dataflow.import_ids(),
850 );
851 if starts_immediately {
852 logging.set_hydration_start();
853 }
854 collection.logging = Some(logging);
855 }
856
857 collection.reset_reported_frontiers(ReportedFrontier::NotReported {
858 lower: as_of.clone(),
859 });
860
861 let existing = self.compute_state.collections.insert(object_id, collection);
862 if existing.is_some() {
863 error!(
864 id = ?object_id,
865 "existing collection for newly created dataflow",
866 );
867 }
868 }
869
870 let (start_signal, suspension_token) = StartSignal::new();
871 for id in dataflow.export_ids() {
872 self.compute_state
873 .suspended_collections
874 .insert(id, Rc::clone(&suspension_token));
875 }
876
877 crate::render::build_compute_dataflow(
878 self.timely_worker,
879 self.compute_state,
880 dataflow,
881 start_signal,
882 until,
883 dataflow_expiration,
884 );
885 }
886
887 fn handle_schedule(&mut self, id: GlobalId) {
888 let suspension_token = self.compute_state.suspended_collections.remove(&id);
894 drop(suspension_token);
895
896 if let Some(collection) = self.compute_state.collections.get(&id) {
897 if let Some(logging) = &collection.logging {
898 logging.set_hydration_start();
899 }
900 }
901 }
902
903 fn handle_allow_compaction(&mut self, id: GlobalId, frontier: Antichain<Timestamp>) {
904 if frontier.is_empty() {
905 self.drop_collection(id);
907 } else {
908 self.compute_state
909 .traces
910 .allow_compaction(id, frontier.borrow());
911 }
912 }
913
914 #[mz_ore::instrument(level = "debug")]
915 fn handle_peek(&mut self, peek: Peek) {
916 let pending = match &peek.target {
917 PeekTarget::Index { id } => {
918 let trace_bundle = self.compute_state.traces.get(id).unwrap().clone();
920 PendingPeek::index(peek, trace_bundle)
921 }
922 PeekTarget::Persist { metadata, .. } => {
923 let metadata = metadata.clone();
924 PendingPeek::persist(
925 peek,
926 Arc::clone(&self.compute_state.persist_clients),
927 metadata,
928 usize::cast_from(self.compute_state.max_result_size),
929 self.timely_worker,
930 PeekRowIterationConfig::new(&self.compute_state.worker_config),
931 )
932 }
933 };
934
935 if let Some(logger) = self.compute_state.compute_logger.as_mut() {
937 logger.log(&pending.as_log_event(true));
938 }
939
940 match pending {
941 PendingPeek::Index(peek) => self.serve_index_peek(&mut Antichain::new(), peek),
942 pending => self.poll_pending_peek(pending),
943 }
944 }
945
946 fn handle_cancel_peek(&mut self, uuid: Uuid) {
947 let queued = &mut self.compute_state.queued_peeks;
948 if let Some(index) = queued.iter().position(|peek| peek.peek.uuid == uuid) {
949 let peek = queued.remove(index).expect("found above");
950 self.send_peek_response(PendingPeek::Index(peek), PeekResponse::Canceled);
951 return;
952 }
953
954 let pending = &mut self.compute_state.pending_peeks;
955 let Some(index) = pending.iter().position(|peek| peek.peek().uuid == uuid) else {
956 return;
957 };
958 let peek = pending.remove(index).expect("found above");
959 self.send_peek_response(peek, PeekResponse::Canceled);
960 }
961
962 fn handle_allow_writes(&mut self, id: GlobalId) {
963 self.compute_state.persist_clients.cfg().enable_compaction();
967
968 if let Some(collection) = self.compute_state.collections.get_mut(&id) {
969 collection.allow_writes();
970 } else {
971 soft_panic_or_log!("allow writes for unknown collection {id}");
972 }
973 }
974
975 fn drop_collection(&mut self, id: GlobalId) {
977 let collection = self
978 .compute_state
979 .collections
980 .remove(&id)
981 .expect("dropped untracked collection");
982
983 self.compute_state.traces.remove(&id);
985 self.compute_state.suspended_collections.remove(&id);
987
988 if let Ok(index) = Rc::try_unwrap(collection.dataflow_index) {
990 self.timely_worker.drop_dataflow(index);
991 }
992
993 if !collection.is_subscribe_or_copy {
998 let reported = collection.reported_frontiers;
999 let write_frontier = (!reported.write_frontier.is_empty()).then(Antichain::new);
1000 let input_frontier = (!reported.input_frontier.is_empty()).then(Antichain::new);
1001 let output_frontier = (!reported.output_frontier.is_empty()).then(Antichain::new);
1002
1003 let frontiers = FrontiersResponse {
1004 write_frontier,
1005 input_frontier,
1006 output_frontier,
1007 };
1008 if frontiers.has_updates() {
1009 self.send_compute_response(ComputeResponse::Frontiers(id, frontiers));
1010 }
1011 }
1012 }
1013
1014 pub fn initialize_logging(&mut self, config: LoggingConfig) {
1016 if self.compute_state.compute_logger.is_some() {
1017 panic!("dataflow server has already initialized logging");
1018 }
1019
1020 let LoggingTraces {
1021 traces,
1022 dataflow_index,
1023 compute_logger: logger,
1024 } = logging::initialize(
1025 self.timely_worker,
1026 &config,
1027 self.compute_state.metrics_registry.clone(),
1028 self.compute_state.metrics.for_logging(),
1029 Rc::clone(&self.compute_state.worker_config),
1030 self.compute_state.workers_per_process,
1031 );
1032
1033 let dataflow_index = Rc::new(dataflow_index);
1034 let mut log_index_ids = config.index_logs;
1035 for (log, trace) in traces {
1036 let id = log_index_ids
1038 .remove(&log)
1039 .expect("`logging::initialize` does not invent logs");
1040 self.compute_state.traces.set(id, trace);
1041
1042 let is_subscribe_or_copy = false;
1044 let as_of = Antichain::from_elem(Timestamp::MIN);
1045 let metrics = self.compute_state.metrics.for_collection(id);
1046 let mut collection = CollectionState::new(
1047 Rc::clone(&dataflow_index),
1048 is_subscribe_or_copy,
1049 as_of,
1050 metrics,
1051 );
1052
1053 let logging =
1054 CollectionLogging::new(id, logger.clone(), *dataflow_index, std::iter::empty());
1055 logging.set_hydration_start();
1060 collection.logging = Some(logging);
1061
1062 let existing = self.compute_state.collections.insert(id, collection);
1063 if existing.is_some() {
1064 error!(
1065 id = ?id,
1066 "existing collection for newly initialized logging export",
1067 );
1068 }
1069 }
1070
1071 assert!(
1073 log_index_ids.is_empty(),
1074 "failed to create requested logging indexes: {log_index_ids:?}",
1075 );
1076
1077 self.compute_state.compute_logger = Some(logger);
1078 }
1079
1080 pub fn report_frontiers(&mut self) {
1082 let mut responses = Vec::new();
1083
1084 let mut new_frontier = Antichain::new();
1086
1087 for (&id, collection) in self.compute_state.collections.iter_mut() {
1088 if collection.is_subscribe_or_copy {
1091 continue;
1092 }
1093
1094 let reported = collection.reported_frontiers();
1095
1096 new_frontier.clear();
1098 if let Some(traces) = self.compute_state.traces.get_mut(&id) {
1099 assert!(
1100 collection.sink_write_frontier.is_none(),
1101 "collection {id} has multiple frontiers"
1102 );
1103 traces.oks_mut().read_upper(&mut new_frontier);
1104 } else if let Some(frontier) = &collection.sink_write_frontier {
1105 new_frontier.clone_from(&frontier.borrow());
1106 } else {
1107 error!(id = ?id, "collection without write frontier");
1108 continue;
1109 }
1110 let new_write_frontier = reported
1111 .write_frontier
1112 .allows_reporting(&new_frontier)
1113 .then(|| new_frontier.clone());
1114
1115 if let Some(probe) = &collection.compute_probe {
1128 if *collection.read_only_rx.borrow() {
1129 new_frontier.clear();
1130 }
1131 probe.with_frontier(|frontier| new_frontier.extend(frontier.iter().copied()));
1132 }
1133 let new_output_frontier = reported
1134 .output_frontier
1135 .allows_reporting(&new_frontier)
1136 .then(|| new_frontier.clone());
1137
1138 new_frontier.clear();
1140 for probe in collection.input_probes.values() {
1141 probe.with_frontier(|frontier| new_frontier.extend(frontier.iter().copied()));
1142 }
1143 let new_input_frontier = reported
1144 .input_frontier
1145 .allows_reporting(&new_frontier)
1146 .then(|| new_frontier.clone());
1147
1148 if let Some(frontier) = &new_write_frontier {
1149 collection
1150 .set_reported_write_frontier(ReportedFrontier::Reported(frontier.clone()));
1151 }
1152 if let Some(frontier) = &new_input_frontier {
1153 collection
1154 .set_reported_input_frontier(ReportedFrontier::Reported(frontier.clone()));
1155 }
1156 if let Some(frontier) = &new_output_frontier {
1157 collection
1158 .set_reported_output_frontier(ReportedFrontier::Reported(frontier.clone()));
1159 }
1160
1161 let response = FrontiersResponse {
1162 write_frontier: new_write_frontier,
1163 input_frontier: new_input_frontier,
1164 output_frontier: new_output_frontier,
1165 };
1166 if response.has_updates() {
1167 responses.push((id, response));
1168 }
1169 }
1170
1171 for (id, frontiers) in responses {
1172 self.send_compute_response(ComputeResponse::Frontiers(id, frontiers));
1173 }
1174 }
1175
1176 pub(crate) fn report_metrics(&self) {
1178 if let Some(expiration) = self.compute_state.replica_expiration.as_option() {
1179 let now = Duration::from_millis(mz_ore::now::SYSTEM_TIME()).as_secs_f64();
1180 let expiration = Duration::from_millis(<u64>::from(expiration)).as_secs_f64();
1181 let remaining = expiration - now;
1182 self.compute_state
1183 .metrics
1184 .replica_expiration_remaining_seconds
1185 .set(remaining)
1186 }
1187 }
1188
1189 fn serve_index_peek(&mut self, upper: &mut Antichain<Timestamp>, peek: IndexPeek) {
1192 match self.compute_state.peek_budget.grant() {
1193 Some(fuel) => self.walk_index_peek(upper, peek, fuel),
1194 None => {
1195 self.compute_state.peek_passed_over = true;
1198 self.compute_state.queued_peeks.push_back(peek);
1199 }
1200 }
1201 }
1202
1203 fn walk_index_peek(
1206 &mut self,
1207 upper: &mut Antichain<Timestamp>,
1208 mut peek: IndexPeek,
1209 fuel: usize,
1210 ) {
1211 let start = Instant::now();
1212
1213 let row_iteration_limit = peek_row_iteration_limit(&self.compute_state.worker_config);
1214
1215 let peek_stash_eligible = peek
1216 .peek
1217 .finishing
1218 .is_streamable(peek.peek.result_desc.arity());
1219
1220 let has_stash_location = {
1224 let enabled = ENABLE_PEEK_RESPONSE_STASH.get(&self.compute_state.worker_config);
1225 let located = self.compute_state.peek_stash_persist_location.is_some();
1226 if !located && enabled {
1227 error!("missing peek_stash_persist_location but peek stash is enabled");
1228 }
1229 enabled && located
1230 };
1231
1232 let stash = StashBounds {
1233 eligible: peek_stash_eligible && has_stash_location,
1234 threshold_bytes: PEEK_RESPONSE_STASH_THRESHOLD_BYTES
1235 .get(&self.compute_state.worker_config),
1236 batch_bytes: PEEK_RESPONSE_STASH_BATCH_BYTES.get(&self.compute_state.worker_config),
1237 };
1238
1239 let metrics = IndexPeekMetrics {
1240 seek_fulfillment_seconds: &self
1241 .compute_state
1242 .metrics
1243 .index_peek_seek_fulfillment_seconds,
1244 frontier_check_seconds: &self.compute_state.metrics.index_peek_frontier_check_seconds,
1245 walk: &self.compute_state.peek_walk_metrics,
1246 };
1247
1248 let mut unspent = fuel;
1249 let status = peek.seek_fulfillment(
1250 upper,
1251 self.compute_state.max_result_size,
1252 stash,
1253 row_iteration_limit,
1254 &mut unspent,
1255 &metrics,
1256 );
1257
1258 self.compute_state
1261 .peek_budget
1262 .charge(fuel.saturating_sub(unspent));
1263
1264 self.compute_state
1265 .metrics
1266 .index_peek_total_seconds
1267 .observe(start.elapsed().as_secs_f64());
1268
1269 match status {
1270 PeekStatus::Ready(response) => {
1271 let _span =
1272 span!(parent: &peek.span, Level::DEBUG, "process_peek_response").entered();
1273 self.send_peek_response(PendingPeek::Index(peek), response);
1274 }
1275 PeekStatus::NotReady => self.compute_state.queued_peeks.push_back(peek),
1276 PeekStatus::Offload(scan) => {
1277 let _span = span!(parent: &peek.span, Level::DEBUG, "offload_index_peek").entered();
1278
1279 let permits = Arc::clone(&self.compute_state.peek_permits);
1280 let config = OffloadConfig::new(&self.compute_state.worker_config);
1281 let walk_metrics = self.compute_state.peek_walk_metrics.clone();
1282 let worker = std::thread::current();
1283 let stash = self
1286 .compute_state
1287 .peek_stash_persist_location
1288 .as_ref()
1289 .filter(|_| scan.stash_eligible())
1290 .cloned()
1291 .map(|location| {
1292 peek_stash::StashTarget::new(
1293 &peek.peek,
1294 Arc::clone(&self.compute_state.persist_clients),
1295 location,
1296 )
1297 });
1298
1299 let offloaded = OffloadedPeek::start(
1300 peek.peek,
1301 scan,
1302 stash,
1303 permits,
1304 config,
1305 walk_metrics,
1306 worker,
1307 );
1308
1309 self.compute_state
1310 .pending_peeks
1311 .push_back(PendingPeek::Offloaded(offloaded));
1312 }
1313 }
1314 }
1315
1316 fn poll_pending_peek(&mut self, mut pending: PendingPeek) {
1319 let response = match &mut pending {
1320 PendingPeek::Index(peek) => {
1323 soft_panic_or_log!(
1324 "index peek on {} polled as if a driver had taken it over",
1325 peek.peek.target.id()
1326 );
1327 None
1328 }
1329 PendingPeek::Persist(peek) => peek.result.try_recv().ok().map(|(result, duration)| {
1330 self.compute_state
1331 .metrics
1332 .persist_peek_seconds
1333 .observe(duration.as_secs_f64());
1334 result
1335 }),
1336 PendingPeek::Offloaded(offloaded) => match offloaded.result.try_recv() {
1337 Ok((response, duration)) => {
1338 self.compute_state
1341 .metrics
1342 .index_peek_offload_seconds
1343 .observe(duration.as_secs_f64());
1344
1345 trace!(?offloaded.peek, ?duration, "finished offloaded index peek walk");
1346 Some(response)
1347 }
1348 Err(oneshot::error::TryRecvError::Empty) => None,
1349 Err(oneshot::error::TryRecvError::Closed) => {
1356 soft_panic_or_log!(
1357 "offloaded walk of peek on {} ended without an outcome",
1358 offloaded.peek.target.id()
1359 );
1360 Some(PeekResponse::Error(PeekError::unstructured(
1361 "offloaded peek walk failed",
1362 )))
1363 }
1364 },
1365 };
1366
1367 if let Some(response) = response {
1368 let _span =
1369 span!(parent: pending.span(), Level::DEBUG, "process_peek_response").entered();
1370 self.send_peek_response(pending, response)
1371 } else {
1372 self.compute_state.pending_peeks.push_back(pending);
1373 }
1374 }
1375
1376 pub fn process_peeks(&mut self) {
1383 self.compute_state.peek_budget.start_activation();
1387
1388 self.compute_state.peek_passed_over = false;
1392
1393 if self.compute_state.pending_peeks.is_empty() && self.compute_state.queued_peeks.is_empty()
1395 {
1396 return;
1397 }
1398
1399 let mut upper = Antichain::new();
1400
1401 let mut pending_peeks = std::mem::take(&mut self.compute_state.pending_peeks);
1404 while let Some(peek) = pending_peeks.pop_front() {
1405 self.poll_pending_peek(peek);
1406 }
1407
1408 let mut queued_peeks = std::mem::take(&mut self.compute_state.queued_peeks);
1411 while let Some(peek) = queued_peeks.pop_front() {
1412 let Some(fuel) = self.compute_state.peek_budget.grant() else {
1413 queued_peeks.push_front(peek);
1414 break;
1415 };
1416 self.walk_index_peek(&mut upper, peek, fuel);
1417 }
1418
1419 self.compute_state.peek_passed_over = !queued_peeks.is_empty();
1422 let served = std::mem::replace(&mut self.compute_state.queued_peeks, queued_peeks);
1423 self.compute_state.queued_peeks.extend(served);
1424 }
1425
1426 #[mz_ore::instrument(level = "debug")]
1431 fn send_peek_response(&mut self, peek: PendingPeek, response: PeekResponse) {
1432 let log_event = peek.as_log_event(false);
1433 self.send_compute_response(ComputeResponse::PeekResponse(
1435 peek.peek().uuid,
1436 response,
1437 OpenTelemetryContext::obtain(),
1438 ));
1439
1440 if let Some(logger) = self.compute_state.compute_logger.as_mut() {
1442 logger.log(&log_event);
1443 }
1444 }
1445
1446 pub fn process_subscribes(&mut self) {
1448 let mut subscribe_responses = self.compute_state.subscribe_response_buffer.borrow_mut();
1449 for (sink_id, mut response) in subscribe_responses.drain(..) {
1450 if let Some(collection) = self.compute_state.collections.get_mut(&sink_id) {
1452 let new_frontier = match &response {
1453 SubscribeResponse::Batch(b) => b.upper.clone(),
1454 SubscribeResponse::DroppedAt(_) => Antichain::new(),
1455 };
1456
1457 let reported = collection.reported_frontiers();
1458 assert!(
1459 reported.write_frontier.allows_reporting(&new_frontier),
1460 "subscribe write frontier regression: {:?} -> {:?}",
1461 reported.write_frontier,
1462 new_frontier,
1463 );
1464 assert!(
1465 reported.input_frontier.allows_reporting(&new_frontier),
1466 "subscribe input frontier regression: {:?} -> {:?}",
1467 reported.input_frontier,
1468 new_frontier,
1469 );
1470
1471 collection
1472 .set_reported_write_frontier(ReportedFrontier::Reported(new_frontier.clone()));
1473 collection
1474 .set_reported_input_frontier(ReportedFrontier::Reported(new_frontier.clone()));
1475 collection.set_reported_output_frontier(ReportedFrontier::Reported(new_frontier));
1476 } else {
1477 }
1480
1481 response
1482 .to_error_if_exceeds(usize::try_from(self.compute_state.max_result_size).unwrap());
1483 self.send_compute_response(ComputeResponse::SubscribeResponse(sink_id, response));
1484 }
1485 }
1486
1487 pub fn process_copy_tos(&self) {
1489 let mut responses = self.compute_state.copy_to_response_buffer.borrow_mut();
1490 for (sink_id, response) in responses.drain(..) {
1491 self.send_compute_response(ComputeResponse::CopyToResponse(sink_id, response));
1492 }
1493 }
1494
1495 fn send_compute_response(&self, response: ComputeResponse) {
1497 let _ = self.response_tx.send(response);
1500 }
1501
1502 pub(crate) fn check_expiration(&self) {
1504 let now = mz_ore::now::SYSTEM_TIME();
1505 if self.compute_state.replica_expiration.less_than(&now.into()) {
1506 let now_datetime = mz_ore::now::to_datetime(now);
1507 let expiration_datetime = self
1508 .compute_state
1509 .replica_expiration
1510 .as_option()
1511 .map(Into::into)
1512 .map(mz_ore::now::to_datetime);
1513
1514 error!(
1517 now,
1518 now_datetime = ?now_datetime,
1519 expiration = ?self.compute_state.replica_expiration.elements(),
1520 expiration_datetime = ?expiration_datetime,
1521 "replica expired"
1522 );
1523
1524 assert!(
1526 !self.compute_state.replica_expiration.less_than(&now.into()),
1527 "replica expired. now: {now} ({now_datetime:?}), expiration: {:?} ({expiration_datetime:?})",
1528 self.compute_state.replica_expiration.elements(),
1529 );
1530 }
1531 }
1532
1533 pub fn determine_dataflow_expiration(
1539 &self,
1540 time_dependence: &TimeDependence,
1541 until: &Antichain<Timestamp>,
1542 ) -> Antichain<Timestamp> {
1543 let iter = self
1548 .compute_state
1549 .replica_expiration
1550 .iter()
1551 .filter_map(|t| time_dependence.apply(*t))
1552 .filter_map(|t| Timestamp::try_step_forward(&t))
1553 .filter(|expiration| !until.less_equal(expiration));
1554 Antichain::from_iter(iter)
1555 }
1556}
1557
1558pub enum PendingPeek {
1563 Index(IndexPeek),
1565 Persist(PersistPeek),
1567 Offloaded(OffloadedPeek),
1570}
1571
1572impl PendingPeek {
1573 pub fn as_log_event(&self, installed: bool) -> ComputeEvent {
1575 let peek = self.peek();
1576 let (id, peek_type) = match &peek.target {
1577 PeekTarget::Index { id } => (*id, logging::compute::PeekType::Index),
1578 PeekTarget::Persist { id, .. } => (*id, logging::compute::PeekType::Persist),
1579 };
1580 let uuid = peek.uuid.into_bytes();
1581 ComputeEvent::Peek(PeekEvent {
1582 id,
1583 time: peek.timestamp,
1584 uuid,
1585 peek_type,
1586 installed,
1587 })
1588 }
1589
1590 fn index(peek: Peek, mut trace_bundle: TraceBundle) -> Self {
1591 let empty_frontier = Antichain::new();
1592 let timestamp_frontier = Antichain::from_elem(peek.timestamp);
1593 trace_bundle
1594 .oks_mut()
1595 .set_logical_compaction(timestamp_frontier.borrow());
1596 trace_bundle
1597 .errs_mut()
1598 .set_logical_compaction(timestamp_frontier.borrow());
1599 trace_bundle
1600 .oks_mut()
1601 .set_physical_compaction(empty_frontier.borrow());
1602 trace_bundle
1603 .errs_mut()
1604 .set_physical_compaction(empty_frontier.borrow());
1605
1606 PendingPeek::Index(IndexPeek {
1607 peek,
1608 trace_bundle,
1609 span: tracing::Span::current(),
1610 })
1611 }
1612
1613 fn persist(
1614 peek: Peek,
1615 persist_clients: Arc<PersistClientCache>,
1616 metadata: CollectionMetadata,
1617 max_result_size: usize,
1618 timely_worker: &TimelyWorker,
1619 row_iteration_config: PeekRowIterationConfig,
1620 ) -> Self {
1621 let active_worker = {
1622 let chosen_index = usize::cast_from(peek.uuid.hashed()) % timely_worker.peers();
1624 chosen_index == timely_worker.index()
1625 };
1626 let activator = timely_worker.sync_activator_for([].into());
1627 let peek_uuid = peek.uuid;
1628
1629 let (result_tx, result_rx) = oneshot::channel();
1630 let timestamp = peek.timestamp;
1631 let mfp_plan = peek.map_filter_project.clone();
1632 let max_results_needed = peek
1633 .finishing
1634 .limit
1635 .map(|l| usize::cast_from(u64::from(l)))
1636 .unwrap_or(usize::MAX)
1637 + peek.finishing.offset;
1638 let order_by = peek.finishing.order_by.clone();
1639
1640 let literal_constraint = peek
1642 .literal_constraints
1643 .clone()
1644 .map(|rows| rows.into_element());
1645
1646 let task_handle = mz_ore::task::spawn(|| "persist::peek", async move {
1647 let start = Instant::now();
1648 let result = if active_worker {
1649 PersistPeek::do_peek(
1650 &persist_clients,
1651 metadata,
1652 timestamp,
1653 literal_constraint,
1654 mfp_plan,
1655 max_result_size,
1656 max_results_needed,
1657 row_iteration_config,
1658 )
1659 .await
1660 } else {
1661 Ok(vec![])
1662 };
1663 let result = match result {
1664 Ok(rows) => PeekResponse::Rows(vec![RowCollection::new(rows, &order_by)]),
1665 Err(error) => PeekResponse::Error(error),
1666 };
1667 match result_tx.send((result, start.elapsed())) {
1668 Ok(()) => {}
1669 Err((_result, elapsed)) => {
1670 debug!(duration =? elapsed, "dropping result for cancelled peek {peek_uuid}")
1671 }
1672 }
1673 match activator.activate() {
1674 Ok(()) => {}
1675 Err(_) => {
1676 debug!("unable to wake timely after completed peek {peek_uuid}");
1677 }
1678 }
1679 });
1680 PendingPeek::Persist(PersistPeek {
1681 peek,
1682 _abort_handle: task_handle.abort_on_drop(),
1683 result: result_rx,
1684 span: tracing::Span::current(),
1685 })
1686 }
1687
1688 fn span(&self) -> &tracing::Span {
1689 match self {
1690 PendingPeek::Index(p) => &p.span,
1691 PendingPeek::Persist(p) => &p.span,
1692 PendingPeek::Offloaded(p) => &p.span,
1693 }
1694 }
1695
1696 pub(crate) fn peek(&self) -> &Peek {
1697 match self {
1698 PendingPeek::Index(p) => &p.peek,
1699 PendingPeek::Persist(p) => &p.peek,
1700 PendingPeek::Offloaded(p) => &p.peek,
1701 }
1702 }
1703}
1704
1705pub struct PersistPeek {
1710 pub(crate) peek: Peek,
1711 _abort_handle: AbortOnDropHandle<()>,
1714 result: oneshot::Receiver<(PeekResponse, Duration)>,
1716 span: tracing::Span,
1718}
1719
1720impl PersistPeek {
1721 async fn do_peek(
1722 persist_clients: &PersistClientCache,
1723 metadata: CollectionMetadata,
1724 as_of: Timestamp,
1725 literal_constraint: Option<Row>,
1726 mfp_plan: SafeMfpPlan,
1727 max_result_size: usize,
1728 mut limit_remaining: usize,
1729 row_iteration_config: PeekRowIterationConfig,
1730 ) -> Result<Vec<(Row, NonZeroUsize)>, PeekError> {
1731 let client = persist_clients
1732 .open(metadata.persist_location)
1733 .await
1734 .map_err(|e| PeekError::unstructured(e.to_string()))?;
1735
1736 let mut reader: ReadHandle<SourceData, (), Timestamp, StorageDiff> = client
1737 .open_leased_reader(
1738 metadata.data_shard,
1739 Arc::new(metadata.relation_desc.clone()),
1740 Arc::new(UnitSchema),
1741 Diagnostics::from_purpose("persist::peek"),
1742 USE_CRITICAL_SINCE_SNAPSHOT.get(client.dyncfgs()),
1743 )
1744 .await
1745 .map_err(|e| PeekError::unstructured(e.to_string()))?;
1746
1747 let mut txns_read = if let Some(txns_id) = metadata.txns_shard {
1754 Some(TxnsCache::open(&client, txns_id, Some(metadata.data_shard)).await)
1755 } else {
1756 None
1757 };
1758
1759 let metrics = client.metrics();
1760
1761 let mut cursor = StatsCursor::new(
1762 &mut reader,
1763 txns_read.as_mut(),
1764 metrics,
1765 &mfp_plan,
1766 &metadata.relation_desc,
1767 Antichain::from_elem(as_of),
1768 )
1769 .await
1770 .map_err(|since| {
1771 PeekError::unstructured(format!(
1772 "attempted to peek at {as_of}, but the since has advanced to {since:?}"
1773 ))
1774 })?;
1775
1776 let mut result = vec![];
1778 let mut datum_vec = DatumVec::new();
1779 let mut row_builder = Row::default();
1780 let arena = RowArena::new();
1781 let mut total_size = 0usize;
1782 let mut row_iteration_tracker = PeekRowIterationTracker::new(None, 0);
1783
1784 let literal_len = match &literal_constraint {
1785 None => 0,
1786 Some(row) => row.iter().count(),
1787 };
1788
1789 'collect: while limit_remaining > 0 {
1790 let Some(batch) = cursor.next().await else {
1791 break;
1792 };
1793 for (data, _, d) in batch {
1794 row_iteration_tracker.set_limit(row_iteration_config.current_limit());
1797 row_iteration_tracker.track_next()?;
1798
1799 let row = data.map_err(PeekError::from)?;
1800
1801 if let Some(literal) = &literal_constraint {
1802 match row.iter().take(literal_len).cmp(literal.iter()) {
1803 Ordering::Less => continue,
1804 Ordering::Equal => {}
1805 Ordering::Greater => break 'collect,
1806 }
1807 }
1808
1809 let count: usize = d.try_into().map_err(|_| {
1810 error!(
1811 shard = %metadata.data_shard, diff = d, ?row,
1812 "persist peek encountered negative multiplicities",
1813 );
1814 PeekError::unstructured(format!(
1815 "Invalid data in source, \
1816 saw retractions ({}) for row that does not exist: {:?}",
1817 -d, row,
1818 ))
1819 })?;
1820 let Some(count) = NonZeroUsize::new(count) else {
1821 continue;
1822 };
1823 let mut datum_local = datum_vec.borrow_with(&row);
1824 let eval_result = mfp_plan
1825 .evaluate_into(&mut datum_local, &arena, &mut row_builder)
1826 .map(|row| row.cloned())
1827 .map_err(PeekError::from)?;
1828 if let Some(row) = eval_result {
1829 total_size = total_size.saturating_add(entry_byte_len(&row));
1830 if total_size > max_result_size {
1831 return Err(PeekError::ResultExceedsMaxSize { max_result_size });
1832 }
1833 result.push((row, count));
1834 limit_remaining = limit_remaining.saturating_sub(count.get());
1835 if limit_remaining == 0 {
1836 break;
1837 }
1838 }
1839 }
1840 }
1841
1842 Ok(result)
1843 }
1844}
1845
1846pub struct IndexPeek {
1848 peek: Peek,
1849 trace_bundle: TraceBundle,
1851 span: tracing::Span,
1853}
1854
1855impl IndexPeek {
1856 fn seek_fulfillment(
1873 &mut self,
1874 upper: &mut Antichain<Timestamp>,
1875 max_result_size: u64,
1876 stash: StashBounds,
1877 row_iteration_limit: Option<usize>,
1878 fuel: &mut usize,
1879 metrics: &IndexPeekMetrics<'_>,
1880 ) -> PeekStatus {
1881 let method_start = Instant::now();
1882
1883 self.trace_bundle.oks_mut().read_upper(upper);
1884 if upper.less_equal(&self.peek.timestamp) {
1885 return PeekStatus::NotReady;
1886 }
1887 self.trace_bundle.errs_mut().read_upper(upper);
1888 if upper.less_equal(&self.peek.timestamp) {
1889 return PeekStatus::NotReady;
1890 }
1891
1892 let read_frontier = self.trace_bundle.compaction_frontier();
1893 if !read_frontier.less_equal(&self.peek.timestamp) {
1894 let error = format!(
1895 "Arrangement compaction frontier ({:?}) is beyond the time of the attempted read ({})",
1896 read_frontier.elements(),
1897 self.peek.timestamp,
1898 );
1899 return PeekStatus::Ready(PeekResponse::Error(PeekError::unstructured(error)));
1900 }
1901
1902 metrics
1903 .frontier_check_seconds
1904 .observe(method_start.elapsed().as_secs_f64());
1905
1906 let result =
1907 self.collect_finished_data(max_result_size, stash, row_iteration_limit, fuel, metrics);
1908
1909 metrics
1910 .seek_fulfillment_seconds
1911 .observe(method_start.elapsed().as_secs_f64());
1912
1913 result
1914 }
1915
1916 fn collect_finished_data(
1922 &mut self,
1923 max_result_size: u64,
1924 stash: StashBounds,
1925 row_iteration_limit: Option<usize>,
1926 fuel: &mut usize,
1927 metrics: &IndexPeekMetrics<'_>,
1928 ) -> PeekStatus {
1929 let peek = &self.peek;
1930 let (oks, errs) = self.trace_bundle.oks_errs_mut();
1931 let mut scan = PeekScan::new(peek, errs, oks, max_result_size, stash);
1932
1933 let outcome = scan.step(row_iteration_limit, fuel);
1934
1935 let phases = scan.phases();
1936 match outcome {
1937 ScanOutcome::Finished(result) => {
1940 metrics.walk.walked_inline();
1941 metrics.walk.observe_error_phase(&phases);
1942 PeekStatus::Ready(match result {
1943 Ok(rows) => {
1944 metrics.walk.observe_ok_phase(&phases);
1945 let start = Instant::now();
1946 let response = rows_response(rows, &self.peek.finishing.order_by);
1947 metrics.walk.observe_row_collection(start.elapsed());
1948 response
1949 }
1950 Err(error) => PeekResponse::Error(error),
1953 })
1954 }
1955 ScanOutcome::Suspended => PeekStatus::Offload(scan),
1964 }
1965 }
1966}
1967
1968enum PeekStatus {
1971 NotReady,
1974 Offload(IndexPeekScan),
1981 Ready(PeekResponse),
1983}
1984
1985#[derive(Debug)]
1987struct ReportedFrontiers {
1988 write_frontier: ReportedFrontier,
1990 input_frontier: ReportedFrontier,
1992 output_frontier: ReportedFrontier,
1994}
1995
1996impl ReportedFrontiers {
1997 fn new() -> Self {
1999 Self {
2000 write_frontier: ReportedFrontier::new(),
2001 input_frontier: ReportedFrontier::new(),
2002 output_frontier: ReportedFrontier::new(),
2003 }
2004 }
2005}
2006
2007#[derive(Clone, Debug)]
2009pub enum ReportedFrontier {
2010 Reported(Antichain<Timestamp>),
2012 NotReported {
2014 lower: Antichain<Timestamp>,
2016 },
2017}
2018
2019impl ReportedFrontier {
2020 pub fn new() -> Self {
2022 let lower = Antichain::from_elem(timely::progress::Timestamp::minimum());
2023 Self::NotReported { lower }
2024 }
2025
2026 pub fn is_empty(&self) -> bool {
2028 match self {
2029 Self::Reported(frontier) => frontier.is_empty(),
2030 Self::NotReported { .. } => false,
2031 }
2032 }
2033
2034 fn allows_reporting(&self, other: &Antichain<Timestamp>) -> bool {
2040 match self {
2041 Self::Reported(frontier) => PartialOrder::less_than(frontier, other),
2042 Self::NotReported { lower } => PartialOrder::less_equal(lower, other),
2043 }
2044 }
2045}
2046
2047pub struct CollectionState {
2049 reported_frontiers: ReportedFrontiers,
2051 dataflow_index: Rc<usize>,
2057 pub is_subscribe_or_copy: bool,
2063 as_of: Antichain<Timestamp>,
2067
2068 pub sink_token: Option<SinkToken>,
2073 pub sink_write_frontier: Option<Rc<RefCell<Antichain<Timestamp>>>>,
2077 pub input_probes: BTreeMap<GlobalId, probe::Handle<Timestamp>>,
2079 pub compute_probe: Option<probe::Handle<Timestamp>>,
2084 logging: Option<CollectionLogging>,
2086 metrics: CollectionMetrics,
2088 read_only_tx: watch::Sender<bool>,
2100 pub read_only_rx: watch::Receiver<bool>,
2102}
2103
2104impl CollectionState {
2105 fn new(
2106 dataflow_index: Rc<usize>,
2107 is_subscribe_or_copy: bool,
2108 as_of: Antichain<Timestamp>,
2109 metrics: CollectionMetrics,
2110 ) -> Self {
2111 let (read_only_tx, read_only_rx) = watch::channel(true);
2114
2115 Self {
2116 reported_frontiers: ReportedFrontiers::new(),
2117 dataflow_index,
2118 is_subscribe_or_copy,
2119 as_of,
2120 sink_token: None,
2121 sink_write_frontier: None,
2122 input_probes: Default::default(),
2123 compute_probe: None,
2124 logging: None,
2125 metrics,
2126 read_only_tx,
2127 read_only_rx,
2128 }
2129 }
2130
2131 fn reported_frontiers(&self) -> &ReportedFrontiers {
2133 &self.reported_frontiers
2134 }
2135
2136 pub fn reset_reported_frontiers(&mut self, frontier: ReportedFrontier) {
2138 self.reported_frontiers.write_frontier = frontier.clone();
2139 self.reported_frontiers.input_frontier = frontier.clone();
2140 self.reported_frontiers.output_frontier = frontier;
2141 }
2142
2143 fn set_reported_write_frontier(&mut self, frontier: ReportedFrontier) {
2145 if let Some(logging) = &mut self.logging {
2146 let time = match &frontier {
2147 ReportedFrontier::Reported(frontier) => frontier.get(0).copied(),
2148 ReportedFrontier::NotReported { .. } => Some(Timestamp::MIN),
2149 };
2150 logging.set_frontier(time);
2151 }
2152
2153 self.reported_frontiers.write_frontier = frontier;
2154 }
2155
2156 fn set_reported_input_frontier(&mut self, frontier: ReportedFrontier) {
2158 if let Some(logging) = &mut self.logging {
2160 for (id, probe) in &self.input_probes {
2161 let new_time = probe.with_frontier(|frontier| frontier.as_option().copied());
2162 logging.set_import_frontier(*id, new_time);
2163 }
2164 }
2165
2166 self.reported_frontiers.input_frontier = frontier;
2167 }
2168
2169 fn set_reported_output_frontier(&mut self, frontier: ReportedFrontier) {
2171 let already_hydrated = self.hydrated();
2172
2173 self.reported_frontiers.output_frontier = frontier;
2174
2175 if !already_hydrated && self.hydrated() {
2176 if let Some(logging) = &mut self.logging {
2177 logging.set_hydrated();
2178 }
2179 self.metrics.record_collection_hydrated();
2180 }
2181 }
2182
2183 fn hydrated(&self) -> bool {
2185 match &self.reported_frontiers.output_frontier {
2186 ReportedFrontier::Reported(frontier) => PartialOrder::less_than(&self.as_of, frontier),
2187 ReportedFrontier::NotReported { .. } => false,
2188 }
2189 }
2190
2191 fn allow_writes(&self) {
2193 info!(
2194 dataflow_index = *self.dataflow_index,
2195 export = ?self.logging.as_ref().map(|l| l.export_id()),
2196 "allowing writes for dataflow",
2197 );
2198 let _ = self.read_only_tx.send(false);
2199 }
2200}
2201
2202#[cfg(test)]
2203mod tests;
2204
2205#[cfg(test)]
2207pub(crate) mod index_peek_tests;
2208
2209#[cfg(test)]
2211mod peek_sweep_tests;