Skip to main content

mz_compute/
compute_state.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//! Worker-local state for compute timely instances.
7
8use 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/// Cheap handles on the dyncfgs that bound how many rows a peek may examine.
91///
92/// The limit is read through handles rather than captured once, because `UpdateConfiguration`
93/// applies to peeks that are already in flight.
94#[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/// Counts the rows a peek has examined on this worker and fails it once that exceeds the limit.
114///
115/// A "row" here is a record the worker had to look at, not a record it returned. Records the MFP
116/// throws away, and records that consolidate to zero, cost scan time all the same, so they count
117/// too. A literal a trace does not hold reaches no record and so costs nothing here; the fuel
118/// budget is what bounds those seeks.
119///
120/// Exactly `limit` rows are allowed. The peek only fails when it asks for the row after that.
121#[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    /// Adopts a new limit without forgetting the rows already examined.
136    ///
137    /// Rows counted while the feature was off still count, so turning it on mid-scan accounts for
138    /// the work the peek has already caused.
139    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    /// Adds rows examined by a walk that ran before this one.
148    ///
149    /// The limit bounds a peek rather than a single walk, so a walk that continues another one
150    /// starts from the count that one reached.
151    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
173/// Worker-local state that is maintained across dataflows.
174///
175/// This state is restricted to the COMPUTE state, the deterministic, idempotent work
176/// done between data ingress and egress.
177pub struct ComputeState {
178    /// State kept for each installed compute collection.
179    ///
180    /// Each collection has exactly one frontier.
181    /// How the frontier is communicated depends on the collection type:
182    ///  * Frontiers of indexes are equal to the frontier of their corresponding traces in the
183    ///    `TraceManager`.
184    ///  * Persist sinks store their current frontier in `CollectionState::sink_write_frontier`.
185    ///  * Subscribes report their frontiers through the `subscribe_response_buffer`.
186    pub collections: BTreeMap<GlobalId, CollectionState>,
187    /// The traces available for sharing across dataflows.
188    pub traces: TraceManager,
189    /// Shared buffer with SUBSCRIBE operator instances by which they can respond.
190    ///
191    /// The entries are pairs of sink identifier (to identify the subscribe instance)
192    /// and the response itself.
193    pub subscribe_response_buffer: Rc<RefCell<Vec<(GlobalId, SubscribeResponse)>>>,
194    /// Shared buffer with S3 oneshot operator instances by which they can respond.
195    ///
196    /// The entries are pairs of sink identifier (to identify the s3 oneshot instance)
197    /// and the response itself.
198    pub copy_to_response_buffer: Rc<RefCell<Vec<(GlobalId, CopyToResponse)>>>,
199    /// Index peeks awaiting their turn on the worker, in the order the sweep serves them.
200    ///
201    /// A sweep takes from the front and returns what it could not retire to the back, so a peek
202    /// it passed over is served before the peeks it served ahead of it, and a peek that arrives
203    /// later queues behind both. Keeping the order here spares the sweep a resume point.
204    pub queued_peeks: VecDeque<IndexPeek>,
205    /// Peeks a driver has taken over, awaiting the outcome that driver hands back.
206    ///
207    /// These are polled on every sweep and draw no budget, because the work they are waiting on is
208    /// not running on the worker.
209    pub pending_peeks: VecDeque<PendingPeek>,
210    /// The persist location where we can stash large peek results.
211    pub peek_stash_persist_location: Option<PersistLocation>,
212    /// The logger, from Timely's logging framework, if logs are enabled.
213    pub compute_logger: Option<logging::compute::Logger>,
214    /// A process-global cache of (blob_uri, consensus_uri) -> PersistClient.
215    /// This is intentionally shared between workers.
216    pub persist_clients: Arc<PersistClientCache>,
217    /// Context necessary for rendering txn-wal operators.
218    pub txns_ctx: TxnsContext,
219    /// History of commands received by this workers and all its peers.
220    pub command_history: ComputeCommandHistory<UIntGauge>,
221    /// Max size in bytes of any result.
222    max_result_size: u64,
223    /// Specification for rendering linear joins.
224    pub linear_join_spec: LinearJoinSpec,
225    /// Metrics for this worker.
226    pub metrics: WorkerMetrics,
227    /// A process-global handle to tracing configuration.
228    tracing_handle: Arc<TracingHandle>,
229    /// Other configuration for compute
230    pub context: ComputeInstanceContext,
231    /// Per-worker dynamic configuration.
232    ///
233    /// This is separate from the process-global `ConfigSet` and contains config options that need
234    /// to be applied consistently with compute command order.
235    ///
236    /// For example, for options that influence dataflow rendering it is important that all workers
237    /// render the same dataflow with the same options. If these options were stored in a global
238    /// `ConfigSet`, we couldn't guarantee that all workers observe changes to them at the same
239    /// point in the stream of compute commands. Storing per-worker configuration ensures that
240    /// because each worker's configuration is only updated once that worker observes the
241    /// respective `UpdateConfiguration` command.
242    ///
243    /// Reference-counted to avoid cloning for `Context`.
244    pub worker_config: Rc<ConfigSet>,
245
246    /// The process-global metrics registry.
247    pub metrics_registry: MetricsRegistry,
248
249    /// The number of timely workers per process.
250    pub workers_per_process: usize,
251
252    /// Bounds how many offloaded peek walks run at once, shared with the other workers of the same
253    /// `serve` call. A process running two compute runtime roles has one of these per role.
254    pub peek_permits: Arc<PeekPermits>,
255
256    /// The metrics an index peek walk reports, whichever driver runs it.
257    ///
258    /// Held here rather than assembled per peek, because an offload clones it into the task and
259    /// the inline driver reads it on every activation of every pending peek.
260    peek_walk_metrics: PeekWalkMetrics,
261
262    /// What this activation may spend walking index peeks on the worker, and what is left of it.
263    ///
264    /// Begun at the top of every sweep, but armed by the first peek that asks for a slice, which
265    /// may be one arriving between two sweeps. Such a peek draws from what the last sweep left, so
266    /// a batch of arrivals costs at most one more aggregate rather than one per peek.
267    peek_budget: InlineBudget,
268
269    /// Whether the last sweep passed a peek over for want of budget.
270    ///
271    /// Says only that such a peek exists, because [`ComputeState::queued_peeks`] already says
272    /// which one it is: the sweep left it at the front of the queue.
273    peek_passed_over: bool,
274
275    /// Collections awaiting schedule instruction by the controller.
276    ///
277    /// Each entry stores a reference to a token that can be dropped to unsuspend the collection's
278    /// dataflow. Multiple collections can reference the same token if they are exported by the
279    /// same dataflow.
280    suspended_collections: BTreeMap<GlobalId, Rc<dyn Any>>,
281
282    /// Interval at which to perform server maintenance tasks. Set to a zero interval to
283    /// perform maintenance with every `step_or_park` invocation.
284    pub server_maintenance_interval: Duration,
285
286    /// The [`mz_ore::now::SYSTEM_TIME`] at which the replica was started.
287    ///
288    /// Used to compute `replica_expiration`.
289    pub init_system_time: EpochMillis,
290
291    /// The maximum time for which the replica is expected to live. If not empty, dataflows in the
292    /// replica can drop diffs associated with timestamps beyond the replica expiration.
293    /// The replica will panic if such dataflows are not dropped before the replica has expired.
294    pub replica_expiration: Antichain<Timestamp>,
295}
296
297impl ComputeState {
298    /// Whether a peek is waiting on nothing but its next turn in the sweep.
299    ///
300    /// Read by the worker loop before it parks. The sweep that serves such a peek runs on the same
301    /// thread, so the loop steps without parking rather than waking itself.
302    pub(crate) fn peeks_awaiting_turn(&self) -> bool {
303        self.peek_passed_over
304    }
305
306    /// Construct a new `ComputeState`.
307    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    /// Return a mutable reference to the identified collection.
355    ///
356    /// Panics if the collection doesn't exist.
357    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    /// Construct a new frontier probe for the given input and add it to the state of the given
364    /// collections.
365    ///
366    /// The caller is responsible for attaching the returned probe handle to the respective
367    /// dataflow input stream.
368    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    /// Apply the current `worker_config` to the compute state.
383    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        // Pager backend selection follows scratch-directory availability:
428        // a scratch dir means the file backend; no scratch dir means swap.
429        // `set_scratch_dir` and `set_backend` are both idempotent, so calling
430        // on every `apply_worker_config` tick is safe. The pager module is
431        // only compiled on Unix targets (`mz_ore::pager` is `cfg(unix)`).
432        #[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        // NB: arrangement dictionary compression is deliberately NOT applied here. Unlike the
448        // settings above, it is captured once at replica creation (see `handle_create_instance`
449        // and `InstanceConfig::arrangement_dictionary_compression`) and held fixed, so that
450        // flipping the flag does not retroactively change arrangements on existing replicas.
451
452        // Apply column-pager configuration. The arrange batchers spill
453        // through the buffer pool below, so the consumers of this budget are
454        // the MV sink's correction buffer and storage's paged upsert stash
455        // flavor, which share one policy and one underlying `mz_ore::pager`.
456        // Routes through `apply_tiered_config`, which reuses a process-wide
457        // `TieredPolicy` singleton — operator-driven tunes mutate the
458        // existing atomics rather than installing a fresh policy with a
459        // fresh budget atomic that would orphan in-flight resident tickets.
460        //
461        // Backend selection mirrors the lower-level `mz_ore::pager`
462        // already configured above: file when a scratch directory is
463        // available, swap otherwise.
464        {
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            // Budget derivation: fraction × announced memory limit, with a
473            // 128 MiB floor so the no-pressure case doesn't page per chunk.
474            // Falls back to a 4 GiB assumption if no limit was announced
475            // (e.g. dev environments).
476            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        // Install and retune the process-wide buffer pool that backs chunk
502        // spilling. Installation is the gate. The pool is constructed, and its
503        // MAP_NORESERVE address space reserved and spill threads spawned, only
504        // when a config apply runs with a spill gate on, so a process that
505        // never enables spilling never mmaps the pool. Config application
506        // reruns on every UpdateConfiguration, so flipping a gate on installs
507        // the pool on the next tick. The pool is a process singleton with no
508        // teardown: once installed it stays active for the life of the process.
509        // Turning every gate back off makes this block do nothing, so the pool
510        // keeps its last-applied budget rather than being uninstalled. Later
511        // ticks with a gate on retune the one instance in place.
512        //
513        // Storage's stash shares the singleton and gates only participation,
514        // so its spill gate installs the pool too. The worker config set is
515        // the full dyncfg aggregate, which is what makes the storage flag
516        // readable here.
517        {
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            // Set compute's leg of the process-wide chunk spill gate. The
524            // gate ORs this leg with storage's, so chunks spill while either
525            // subsystem's flag is set. Storage's config application writes
526            // only its own leg, keeping the two flags from clobbering each
527            // other. The correction buffer has a gate of its own.
528            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                // Budget derivation: fraction of physical RAM, with a 128 MiB
537                // floor so the no-pressure case doesn't page per chunk.
538                // Resident budgets derive from RAM, never from the announced
539                // memory limit, which on swap-provisioned nodes deliberately
540                // includes swap for the memory limiter's purposes. Falls back
541                // to a 4 GiB assumption if detection fails.
542                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                // No ordering is enforced between the target and the budget. A
550                // target at or below budget + warm cap leaves no compressed-tier
551                // headroom, which legally collapses the tier. Every backing
552                // write then pages out immediately, the pre-tier behavior.
553                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            // The generational depth floor below which spilled bodies store
579            // uncompressed. Subsystem-independent, so applied here alongside
580            // the rest of the process-wide chunk configuration.
581            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        // Remember the maintenance interval locally to avoid reading it from the config set on
587        // every server iteration.
588        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    /// Apply the provided replica expiration `offset` by converting it to a frontier relative to
603    /// the replica's initialization system time.
604    ///
605    /// Only expected to be called once when creating the instance. Guards against calling it
606    /// multiple times by checking if the local expiration time is set.
607    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            // Record the replica expiration in the metrics.
625            self.metrics
626                .replica_expiration_timestamp_seconds
627                .set(replica_expiration.into());
628        }
629    }
630
631    /// Returns the cc or non-cc version of "dataflow_max_inflight_bytes", as
632    /// appropriate to this replica.
633    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
646/// A wrapper around [ComputeState] with a live timely worker and response channel.
647pub(crate) struct ActiveComputeState<'a> {
648    /// The underlying Timely worker.
649    pub timely_worker: &'a mut TimelyWorker,
650    /// The compute state itself.
651    pub compute_state: &'a mut ComputeState,
652    /// The channel over which frontier information is reported.
653    pub response_tx: &'a mut ResponseSender,
654}
655
656/// A token that keeps a sink alive.
657pub struct SinkToken(#[allow(dead_code)] Box<dyn Any>);
658
659impl SinkToken {
660    /// Create a new `SinkToken`.
661    pub fn new(t: Box<dyn Any>) -> Self {
662        Self(t)
663    }
664}
665
666impl<'a> ActiveComputeState<'a> {
667    /// Entrypoint for applying a compute command.
668    #[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        // Record the command duration, per worker and command kind.
675        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        // Seed the worker configuration with the controller's snapshot before applying it, so
705        // create-time setup observes controller-synced values rather than dyncfg defaults. The
706        // same values arrive again in the following `UpdateConfiguration`, which applies globally
707        // and keeps the configuration current. An empty snapshot leaves the defaults in place.
708        config
709            .initial_config
710            .apply(&self.compute_state.worker_config);
711
712        // Ensure the state is consistent with the config before we initialize anything.
713        self.compute_state.apply_worker_config();
714
715        // Apply dictionary compression exactly once, here at instance creation, from the value the
716        // controller captured when the replica was created. We deliberately do NOT re-apply it on
717        // `handle_update_configuration`, so flipping the flag does not retroactively change this
718        // replica's arrangements. `DICTIONARY_COMPRESSION` is process-global and a replica process
719        // hosts a single instance, so this single store covers all of the replica's arrangements.
720        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        // Note: We're only updating mz_metrics from the compute state here, but not from the
761        // equivalent storage state. This is because they're running on the same process and
762        // share the metrics.
763        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        // Add the dataflow expiration to `until`.
784        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        // `StartSignal` is attached only to imported sources and imported indexes, and
824        // `import_ids` is exactly those two sets, so a dataflow with no imports has nothing
825        // suspended and begins computing as soon as it is rendered. Such a dataflow can reach
826        // hydration before its `Schedule` arrives, and the controller sends one anyway to keep
827        // protocol communication predictable, so a `started_at` stamped only from
828        // `handle_schedule` would land after `hydrated_at`. Stamping it here also keeps the row
829        // truthful from the moment it appears, rather than reporting the object as queued while
830        // nothing is queueing it.
831        let starts_immediately = dataflow.import_ids().next().is_none();
832
833        // Initialize compute and logging state for each object.
834        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        // A `Schedule` command instructs us to begin dataflow computation for a collection, so
889        // we should unsuspend it by dropping the corresponding suspension token. Note that a
890        // dataflow can export multiple collections and they all share one suspension token, so the
891        // computation of a dataflow will only start once all its exported collections have been
892        // scheduled.
893        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            // Indicates that we may drop `id`, as there are no more valid times to read.
906            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                // Acquire a copy of the trace suitable for fulfilling the peek.
919                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        // Log the receipt of the peek.
936        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        // Enable persist compaction on any allow-writes command. We
964        // assume persist only compacts after making durable changes,
965        // such as appending a batch or advancing the upper.
966        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    /// Drop the given collection.
976    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        // If this collection is an index, remove its trace.
984        self.compute_state.traces.remove(&id);
985        // If the collection is unscheduled, remove it from the list of waiting collections.
986        self.compute_state.suspended_collections.remove(&id);
987
988        // Drop the dataflow, if all its exports have been dropped.
989        if let Ok(index) = Rc::try_unwrap(collection.dataflow_index) {
990            self.timely_worker.drop_dataflow(index);
991        }
992
993        // The compute protocol requires us to send a `Frontiers` response with empty frontiers
994        // when a collection was dropped, unless:
995        //  * The frontier was already reported as empty previously, or
996        //  * The collection is a subscribe or copy-to.
997        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    /// Initializes timely dataflow logging and publishes as a view.
1015    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            // Install trace as maintained index.
1037            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            // Initialize compute and logging state for the logging index.
1043            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            // Log collections are never suspended and the controller marks them scheduled
1056            // implicitly, so no `Schedule` command ever arrives for them. Record their hydration
1057            // start here, or they would sit permanently in the illegal state of being hydrated
1058            // without having started.
1059            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        // Sanity check.
1072        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    /// Send progress information to the controller.
1081    pub fn report_frontiers(&mut self) {
1082        let mut responses = Vec::new();
1083
1084        // Maintain a single allocation for `new_frontier` to avoid allocating on every iteration.
1085        let mut new_frontier = Antichain::new();
1086
1087        for (&id, collection) in self.compute_state.collections.iter_mut() {
1088            // The compute protocol does not allow `Frontiers` responses for subscribe and copy-to
1089            // collections (database-issues#4701).
1090            if collection.is_subscribe_or_copy {
1091                continue;
1092            }
1093
1094            let reported = collection.reported_frontiers();
1095
1096            // Collect the write frontier and check for progress.
1097            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            // Collect the output frontier and check for progress.
1116            //
1117            // By default, the output frontier equals the write frontier (which is still stored in
1118            // `new_frontier`). If the collection provides a compute frontier, we construct the
1119            // output frontier by taking the meet of write and compute frontier, to avoid:
1120            //  * reporting progress through times we have not yet written
1121            //  * reporting progress through times we have not yet fully processed, for
1122            //    collections that jump their write frontiers into the future
1123            //
1124            // As a special case, in read-only mode we don't take the write frontier into account.
1125            // The dataflow doesn't have the ability to push it forward, so it can't be used as a
1126            // measure of dataflow progress.
1127            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            // Collect the input frontier and check for progress.
1139            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    /// Report per-worker metrics.
1177    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    /// Gives `peek` a turn on the worker if this activation's budget has one left, and queues it
1190    /// for a later activation otherwise.
1191    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                // A scan is opened by the slice that walks it, so passing a peek over costs it an
1196                // activation and nothing else.
1197                self.compute_state.peek_passed_over = true;
1198                self.compute_state.queued_peeks.push_back(peek);
1199            }
1200        }
1201    }
1202
1203    /// Walks `peek` for up to `fuel` cursor positions and either answers it, hands it to a driver
1204    /// that finishes it, or returns it to the queue for another turn.
1205    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        // Whether a diverted peek has somewhere to write its rows. A flag here and the location
1221        // itself only where an offload needs it, because this runs for every peek the sweep gives
1222        // a turn, including the point lookups that answer inline and never reach the stash.
1223        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        // Charged with what the slice walked rather than with what it was granted, so a peek that
1259        // answers in three positions leaves the activation's budget to the peeks behind it.
1260        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                // Read off the scan rather than decided again, so a walk that can offer a batch
1284                // always has somewhere to write it.
1285                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    /// Asks the driver that has taken `pending` over for its outcome, and sends the response when
1317    /// one is ready.
1318    fn poll_pending_peek(&mut self, mut pending: PendingPeek) {
1319        let response = match &mut pending {
1320            // An index peek reaches a driver only by leaving the queue, and the driver that takes
1321            // it over replaces it with a variant of its own.
1322            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                    // Covers the writing to the peek stash too, because the walk that produces the
1339                    // rows is the one that writes them.
1340                    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                // The task drops its sender without sending only on a cancellation, which removes
1350                // this entry, so an entry still here to be polled means the task died. Answering
1351                // keeps the peek from waiting forever on a walk nothing is running.
1352                //
1353                // NOTE: a walk dropped by a shutting-down tokio runtime arrives here the same way.
1354                // The worker is going away too, so the log line is noise rather than a lost signal.
1355                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    /// Scan the peeks a driver is finishing and the peeks awaiting a turn, and attempt to retire
1377    /// each.
1378    ///
1379    /// The queue of peeks awaiting a turn is served from the front and each peek it cannot retire
1380    /// is returned to the back, so the next sweep resumes where this one ran out of budget without
1381    /// either of them recording where that was.
1382    pub fn process_peeks(&mut self) {
1383        // Above the early return because this is the only place an activation begins. A replica
1384        // whose peeks all answer inline leaves none pending, and beginning an activation only
1385        // where there is work would let the aggregate drain across its arrivals.
1386        self.compute_state.peek_budget.start_activation();
1387
1388        // Says what this sweep found, so it is cleared before the sweep rather than carried in
1389        // from the last one. A peek cancelled or dropped between two sweeps would otherwise leave
1390        // the worker spinning on a turn nothing is waiting for.
1391        self.compute_state.peek_passed_over = false;
1392
1393        // Runs on every iteration of the worker loop, and almost every one finds no peek at all.
1394        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        // Both queues are taken out of the state for the sweep, because serving a peek borrows
1402        // `self` mutably. A peek the sweep returns for another turn lands in the emptied queue.
1403        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        // The aggregate does not refill within an activation, so the first peek the budget cannot
1409        // serve is also the last: every peek behind it would be passed over for the same reason.
1410        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        // A peek the sweep never reached keeps its place ahead of the peeks it served, so the
1420        // queue rotates instead of starving its tail.
1421        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    /// Sends a response for this peek's resolution to the coordinator.
1427    ///
1428    /// Note that this function takes ownership of the `PendingPeek`, which is
1429    /// meant to prevent multiple responses to the same peek.
1430    #[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        // Respond with the response.
1434        self.send_compute_response(ComputeResponse::PeekResponse(
1435            peek.peek().uuid,
1436            response,
1437            OpenTelemetryContext::obtain(),
1438        ));
1439
1440        // Log responding to the peek request.
1441        if let Some(logger) = self.compute_state.compute_logger.as_mut() {
1442            logger.log(&log_event);
1443        }
1444    }
1445
1446    /// Scan the shared subscribe response buffer, and forward results along.
1447    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            // Update frontier logging for this subscribe.
1451            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                // Presumably tracking state for this subscribe was already dropped by
1478                // `drop_collection`. There is nothing left to do for logging.
1479            }
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    /// Scan the shared copy to response buffer, and forward results along.
1488    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    /// Send a response to the coordinator.
1496    fn send_compute_response(&self, response: ComputeResponse) {
1497        // Ignore send errors because the coordinator is free to ignore our
1498        // responses. This happens during shutdown.
1499        let _ = self.response_tx.send(response);
1500    }
1501
1502    /// Checks for dataflow expiration. Panics if we're past the replica expiration time.
1503    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            // We error and assert separately to produce structured logs in anything that depends
1515            // on tracing.
1516            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            // Repeat condition for better error message.
1525            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    /// Returns the dataflow expiration, i.e, the timestamp beyond which diffs can be
1534    /// dropped.
1535    ///
1536    /// Returns an empty timestamp if `replica_expiration` is unset or matches conditions under
1537    /// which dataflow expiration should be disabled.
1538    pub fn determine_dataflow_expiration(
1539        &self,
1540        time_dependence: &TimeDependence,
1541        until: &Antichain<Timestamp>,
1542    ) -> Antichain<Timestamp> {
1543        // Evaluate time dependence with respect to the expiration time.
1544        // * Step time forward to ensure the expiration time is different to the moment a dataflow
1545        //   can legitimately jump to.
1546        // * We cannot expire dataflow with an until that is less or equal to the expiration time.
1547        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
1558/// A peek against either an index or a Persist collection.
1559///
1560/// Note that `PendingPeek` intentionally does not implement or derive `Clone`,
1561/// as each `PendingPeek` is meant to be dropped after it's responded to.
1562pub enum PendingPeek {
1563    /// A peek against an index. (Possibly a temporary index created for the purpose.)
1564    Index(IndexPeek),
1565    /// A peek against a Persist-backed collection.
1566    Persist(PersistPeek),
1567    /// A peek against an index whose walk was offloaded from the worker and is running as an async
1568    /// task.
1569    Offloaded(OffloadedPeek),
1570}
1571
1572impl PendingPeek {
1573    /// Produces a corresponding log event.
1574    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            // Choose the worker that does the actual peek arbitrarily but consistently.
1623            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        // Persist peeks can include at most one literal constraint.
1641        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
1705/// An in-progress Persist peek.
1706///
1707/// Note that `PendingPeek` intentionally does not implement or derive `Clone`,
1708/// as each `PendingPeek` is meant to be dropped after it's responded to.
1709pub struct PersistPeek {
1710    pub(crate) peek: Peek,
1711    /// A background task that's responsible for producing the peek results.
1712    /// If we're no longer interested in the results, we abort the task.
1713    _abort_handle: AbortOnDropHandle<()>,
1714    /// The result of the background task, eventually.
1715    result: oneshot::Receiver<(PeekResponse, Duration)>,
1716    /// The `tracing::Span` tracking this peek's operation
1717    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        // If we are using txn-wal for this collection, then the upper might
1748        // be advanced lazily and we have to go through txn-wal for reads.
1749        //
1750        // TODO: If/when we have a process-wide TxnsRead worker for clusterd,
1751        // use in here (instead of opening a new TxnsCache) to save a persist
1752        // reader registration and some txns shard read traffic.
1753        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        // Re-used state for processing and building rows.
1777        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                // Count before literal and MFP filtering because the Persist row
1795                // has already been read and must still be examined.
1796                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
1846/// An in-progress index-backed peek, and data to eventually fulfill it.
1847pub struct IndexPeek {
1848    peek: Peek,
1849    /// The data from which the trace derives.
1850    trace_bundle: TraceBundle,
1851    /// The `tracing::Span` tracking this peek's operation
1852    span: tracing::Span,
1853}
1854
1855impl IndexPeek {
1856    /// Attempts to fulfill the peek and reports success.
1857    ///
1858    /// To produce output at `peek.timestamp`, we must be certain that
1859    /// it is no longer changing. A trace guarantees that all future
1860    /// changes will be greater than or equal to an element of `upper`.
1861    ///
1862    /// If an element of `upper` is less or equal to `peek.timestamp`,
1863    /// then there can be further updates that would change the output.
1864    /// If no element of `upper` is less or equal to `peek.timestamp`,
1865    /// then for any time `t` less or equal to `peek.timestamp` it is
1866    /// not the case that `upper` is less or equal to that timestamp,
1867    /// and so the result cannot further evolve.
1868    ///
1869    /// `fuel` bounds how far the walk may go on this worker, in cursor positions, and is charged
1870    /// for the positions it visits. A walk that exhausts it with work left is offloaded rather than
1871    /// continued here, and a peek whose frontiers do not admit the read yet spends none of it.
1872    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    /// Answers the peek by scanning the traces that fulfil it, for as long as `fuel` allows.
1917    ///
1918    /// One call opens one scan and either answers from it or hands it on, so nothing survives the
1919    /// call. A scan that runs out of fuel with work left leaves with the [`PeekStatus::Offload`]
1920    /// that reports it, so the positions it walked are not walked again.
1921    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            // Both answers end the walk on this worker, so this driver accounts for it either
1938            // way.
1939            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                    // The ok phase goes unreported, because an error can come from either walk
1951                    // and its numbers describe a finished ok walk only when rows came out of it.
1952                    Err(error) => PeekResponse::Error(error),
1953                })
1954            }
1955            // The one outcome that leaves the walk unfinished, and so the one this driver
1956            // reports nothing for.
1957            //
1958            // A scan suspends out of fuel or holding a full batch, and this driver can carry on
1959            // with neither: it walks under a budget the slice has spent, and it writes no rows, so
1960            // a batch handed to it here would have to be dropped. Every position the scan walked
1961            // travels with it, and so does their cost, which is what makes offload cost one
1962            // hand-off rather than a second walk.
1963            ScanOutcome::Suspended => PeekStatus::Offload(scan),
1964        }
1965    }
1966}
1967
1968/// For keeping track of the state of pending or ready peeks, and managing
1969/// control flow.
1970enum PeekStatus {
1971    /// The frontiers of objects are not yet advanced enough, peek is still
1972    /// pending.
1973    NotReady,
1974    /// The walk stopped with work left, so it is finished away from the worker. Carries the scan,
1975    /// which resumes from the cursor positions it stopped on.
1976    ///
1977    /// A walk stops either because it spent the fuel this activation granted it or because its
1978    /// accumulated rows grew into a batch bound for the peek stash. Both leave here, because the
1979    /// driver that finishes a walk is also the one that writes to the stash.
1980    Offload(IndexPeekScan),
1981    /// The peek result is ready.
1982    Ready(PeekResponse),
1983}
1984
1985/// The frontiers we have reported to the controller for a collection.
1986#[derive(Debug)]
1987struct ReportedFrontiers {
1988    /// The reported write frontier.
1989    write_frontier: ReportedFrontier,
1990    /// The reported input frontier.
1991    input_frontier: ReportedFrontier,
1992    /// The reported output frontier.
1993    output_frontier: ReportedFrontier,
1994}
1995
1996impl ReportedFrontiers {
1997    /// Creates a new `ReportedFrontiers` instance.
1998    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/// A frontier we have reported to the controller, or the least frontier we are allowed to report.
2008#[derive(Clone, Debug)]
2009pub enum ReportedFrontier {
2010    /// A frontier has been previously reported.
2011    Reported(Antichain<Timestamp>),
2012    /// No frontier has been reported yet.
2013    NotReported {
2014        /// A lower bound for frontiers that may be reported in the future.
2015        lower: Antichain<Timestamp>,
2016    },
2017}
2018
2019impl ReportedFrontier {
2020    /// Create a new `ReportedFrontier` enforcing the minimum lower bound.
2021    pub fn new() -> Self {
2022        let lower = Antichain::from_elem(timely::progress::Timestamp::minimum());
2023        Self::NotReported { lower }
2024    }
2025
2026    /// Whether the reported frontier is the empty frontier.
2027    pub fn is_empty(&self) -> bool {
2028        match self {
2029            Self::Reported(frontier) => frontier.is_empty(),
2030            Self::NotReported { .. } => false,
2031        }
2032    }
2033
2034    /// Whether this `ReportedFrontier` allows reporting the given frontier.
2035    ///
2036    /// A `ReportedFrontier` allows reporting of another frontier if:
2037    ///  * The other frontier is greater than the reported frontier.
2038    ///  * The other frontier is greater than or equal to the lower bound.
2039    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
2047/// State maintained for a compute collection.
2048pub struct CollectionState {
2049    /// Tracks the frontiers that have been reported to the controller.
2050    reported_frontiers: ReportedFrontiers,
2051    /// The index of the dataflow computing this collection.
2052    ///
2053    /// Used for dropping the dataflow when the collection is dropped.
2054    /// The Dataflow index is wrapped in an `Rc`s and can be shared between collections, to reflect
2055    /// the possibility that a single dataflow can export multiple collections.
2056    dataflow_index: Rc<usize>,
2057    /// Whether this collection is a subscribe or copy-to.
2058    ///
2059    /// The compute protocol does not allow `Frontiers` responses for subscribe and copy-to
2060    /// collections, so we need to be able to recognize them. This is something we would like to
2061    /// change in the future (database-issues#4701).
2062    pub is_subscribe_or_copy: bool,
2063    /// The collection's initial as-of frontier.
2064    ///
2065    /// Used to determine hydration status.
2066    as_of: Antichain<Timestamp>,
2067
2068    /// A token that should be dropped when this collection is dropped to clean up associated
2069    /// sink state.
2070    ///
2071    /// Only `Some` if the collection is a sink.
2072    pub sink_token: Option<SinkToken>,
2073    /// Frontier of sink writes.
2074    ///
2075    /// Only `Some` if the collection is a sink and *not* a subscribe.
2076    pub sink_write_frontier: Option<Rc<RefCell<Antichain<Timestamp>>>>,
2077    /// Frontier probes for every input to the collection.
2078    pub input_probes: BTreeMap<GlobalId, probe::Handle<Timestamp>>,
2079    /// A probe reporting the frontier of times through which all collection outputs have been
2080    /// computed (but not necessarily written).
2081    ///
2082    /// `None` for collections with compute frontiers equal to their write frontiers.
2083    pub compute_probe: Option<probe::Handle<Timestamp>>,
2084    /// Logging state maintained for this collection.
2085    logging: Option<CollectionLogging>,
2086    /// Metrics tracked for this collection.
2087    metrics: CollectionMetrics,
2088    /// Send-side to transition a dataflow from read-only mode to read-write mode.
2089    ///
2090    /// All dataflows start in read-only mode. Only after receiving a
2091    /// `AllowWrites` command from the controller will they transition to
2092    /// read-write mode.
2093    ///
2094    /// A dataflow in read-only mode must not affect any external state.
2095    ///
2096    /// NOTE: In the future, we might want a more complicated flag, for example
2097    /// something that tells us after which timestamp we are allowed to write.
2098    /// In this first version we are keeping things as simple as possible!
2099    read_only_tx: watch::Sender<bool>,
2100    /// Receive-side to observe whether a dataflow is in read-only mode.
2101    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        // We always initialize as read_only=true. Only when we're explicitly
2112        // allowed to we switch to read-write.
2113        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    /// Return the frontiers that have been reported to the controller.
2132    fn reported_frontiers(&self) -> &ReportedFrontiers {
2133        &self.reported_frontiers
2134    }
2135
2136    /// Reset all reported frontiers to the given value.
2137    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    /// Set the write frontier that has been reported to the controller.
2144    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    /// Set the input frontier that has been reported to the controller.
2157    fn set_reported_input_frontier(&mut self, frontier: ReportedFrontier) {
2158        // Use this opportunity to update our input frontier logging.
2159        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    /// Set the output frontier that has been reported to the controller.
2170    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    /// Return whether this collection is hydrated.
2184    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    /// Allow writes for this collection.
2192    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/// Tests of the inline index-peek driver, and the fixtures [`peek_scan`]'s tests share with it.
2206#[cfg(test)]
2207pub(crate) mod index_peek_tests;
2208
2209/// Tests of the sweep that drives the pending index peeks.
2210#[cfg(test)]
2211mod peek_sweep_tests;