Skip to main content

mz_storage/
storage_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 storage timely instances.
7//!
8//! One instance of a [`Worker`], along with its contained [`StorageState`], is
9//! part of an ensemble of storage workers that all run inside the same timely
10//! cluster. We call this worker a _storage worker_ to disambiguate it from
11//! other kinds of workers, potentially other components that might be sharing
12//! the same timely cluster.
13//!
14//! ## Controller and internal communication
15//!
16//! A worker receives _external_ [`StorageCommands`](StorageCommand) from the
17//! storage controller, via a channel. Storage workers also share an _internal_
18//! control/command fabric ([`internal_control`]). Internal commands go through
19//! a sequencer dataflow that ensures that all workers receive all commands in
20//! the same consistent order.
21//!
22//! We need to make sure that commands that cause dataflows to be rendered are
23//! processed in the same consistent order across all workers because timely
24//! requires this. To achieve this, we make sure that only internal commands can
25//! cause dataflows to be rendered. External commands (from the controller)
26//! cause internal commands to be broadcast (by only one worker), to get
27//! dataflows rendered.
28//!
29//! The internal command fabric is also used to broadcast messages from a local
30//! operator/worker to all workers. For example, when we need to tear down and
31//! restart a dataflow on all workers when an error is encountered.
32//!
33//! ## Async Storage Worker
34//!
35//! The storage worker has a companion [`AsyncStorageWorker`] that must be used
36//! when running code that requires `async`. This is needed because a timely
37//! main loop cannot run `async` code.
38//!
39//! ## Example flow of commands for `RunIngestion`
40//!
41//! With external commands, internal commands, and the async worker,
42//! understanding where and how commands from the controller are realized can
43//! get complicated. We will follow the complete flow for `RunIngestion`, as an
44//! example:
45//!
46//! 1. Worker receives a [`StorageCommand::RunIngestion`] command from the
47//!    controller.
48//! 2. This command is processed in [`StorageState::handle_storage_command`].
49//!    This step cannot render dataflows, because it does not have access to the
50//!    timely worker. It will only set up state that stays over the whole
51//!    lifetime of the source, such as the `reported_frontier`. Putting in place
52//!    this reported frontier will enable frontier reporting for that source. We
53//!    will not start reporting when we only see an internal command for
54//!    rendering a dataflow, which can "overtake" the external `RunIngestion`
55//!    command.
56//! 3. During processing of that command, we call
57//!    [`AsyncStorageWorker::update_ingestion_frontiers`], which causes a command to
58//!    be sent to the async worker.
59//! 4. We eventually get a response from the async worker:
60//!    [`AsyncStorageWorkerResponse::IngestionFrontiersUpdated`].
61//! 5. This response is handled in [`Worker::handle_async_worker_response`].
62//! 6. Handling that response causes a
63//!    [`InternalStorageCommand::CreateIngestionDataflow`] to be broadcast to
64//!    all workers via the internal command fabric.
65//! 7. This message will be processed (on each worker) in
66//!    [`Worker::handle_internal_storage_command`]. This is what will cause the
67//!    required dataflow to be rendered on all workers.
68//!
69//! The process described above assumes that the `RunIngestion` is _not_ an
70//! update, i.e. it is in response to a `CREATE SOURCE`-like statement.
71//!
72//! The primary distinction when handling a `RunIngestion` that represents an
73//! update, is that it might fill out new internal state in the mid-level
74//! clients on the way toward being run.
75
76use std::cell::RefCell;
77use std::collections::{BTreeMap, BTreeSet, VecDeque};
78use std::path::PathBuf;
79use std::rc::Rc;
80use std::sync::Arc;
81use std::thread;
82use std::time::Duration;
83
84use fail::fail_point;
85use mz_ore::now::NowFn;
86use mz_ore::soft_assert_or_log;
87use mz_ore::tracing::TracingHandle;
88use mz_persist_client::batch::ProtoBatch;
89use mz_persist_client::cache::PersistClientCache;
90use mz_persist_client::operators::shard_source::ErrorHandler;
91use mz_repr::{GlobalId, Timestamp};
92use mz_rocksdb::config::SharedWriteBufferManager;
93use mz_storage_client::client::{
94    RunIngestionCommand, StatusUpdate, StorageCommand, StorageResponse,
95};
96use mz_storage_types::AlterCompatible;
97use mz_storage_types::configuration::StorageConfiguration;
98use mz_storage_types::connections::ConnectionContext;
99use mz_storage_types::controller::CollectionMetadata;
100use mz_storage_types::dyncfgs::STORAGE_SERVER_MAINTENANCE_INTERVAL;
101use mz_storage_types::oneshot_sources::OneshotIngestionDescription;
102use mz_storage_types::sinks::StorageSinkDesc;
103use mz_storage_types::sources::IngestionDescription;
104use mz_timely_util::builder_async::PressOnDropButton;
105use mz_txn_wal::operator::TxnsContext;
106use timely::order::PartialOrder;
107use timely::progress::Timestamp as _;
108use timely::progress::frontier::Antichain;
109use timely::worker::Worker as TimelyWorker;
110use tokio::sync::mpsc::error::TryRecvError;
111use tokio::sync::{mpsc, watch};
112use tokio::time::Instant;
113use tracing::{debug, info, warn};
114use uuid::Uuid;
115
116use crate::internal_control::{
117    self, DataflowParameters, InternalCommandReceiver, InternalCommandSender,
118    InternalStorageCommand,
119};
120use crate::metrics::StorageMetrics;
121use crate::statistics::{AggregatedStatistics, SinkStatistics, SourceStatistics};
122use crate::storage_state::async_storage_worker::{AsyncStorageWorker, AsyncStorageWorkerResponse};
123
124pub mod async_storage_worker;
125
126type CommandReceiver = mpsc::UnboundedReceiver<StorageCommand>;
127type ResponseSender = mpsc::UnboundedSender<StorageResponse>;
128
129/// State maintained for each worker thread.
130///
131/// Much of this state can be viewed as local variables for the worker thread,
132/// holding state that persists across function calls.
133pub struct Worker<'w> {
134    /// The underlying Timely worker.
135    ///
136    /// NOTE: This is `pub` for testing.
137    pub timely_worker: &'w mut TimelyWorker,
138    /// The channel over which communication handles for newly connected clients
139    /// are delivered.
140    pub client_rx: mpsc::UnboundedReceiver<(Uuid, CommandReceiver, ResponseSender)>,
141    /// The state associated with collection ingress and egress.
142    pub storage_state: StorageState,
143}
144
145impl<'w> Worker<'w> {
146    /// Creates new `Worker` state from the given components.
147    pub fn new(
148        timely_worker: &'w mut TimelyWorker,
149        client_rx: mpsc::UnboundedReceiver<(Uuid, CommandReceiver, ResponseSender)>,
150        metrics: StorageMetrics,
151        now: NowFn,
152        connection_context: ConnectionContext,
153        instance_context: StorageInstanceContext,
154        persist_clients: Arc<PersistClientCache>,
155        txns_ctx: TxnsContext,
156        tracing_handle: Arc<TracingHandle>,
157        shared_rocksdb_write_buffer_manager: SharedWriteBufferManager,
158    ) -> Self {
159        // It is very important that we only create the internal control
160        // flow/command sequencer once because a) the worker state is re-used
161        // when a new client connects and b) dataflows that have already been
162        // rendered into the timely worker are reused as well.
163        //
164        // If we created a new sequencer every time we get a new client (likely
165        // because the controller re-started and re-connected), dataflows that
166        // were rendered before would still hold a handle to the old sequencer
167        // but we would not read their commands anymore.
168        let (internal_cmd_tx, internal_cmd_rx) =
169            internal_control::setup_command_sequencer(timely_worker);
170
171        let storage_state = StorageState::new_guest(
172            timely_worker.index(),
173            timely_worker.peers(),
174            internal_cmd_tx,
175            Some(internal_cmd_rx),
176            metrics,
177            now,
178            connection_context,
179            instance_context,
180            persist_clients,
181            txns_ctx,
182            tracing_handle,
183            shared_rocksdb_write_buffer_manager,
184        );
185
186        // TODO(aljoscha): We might want `async_worker` and `internal_cmd_tx` to
187        // be fields of `Worker` instead of `StorageState`, but at least for the
188        // command flow sources and sinks need access to that. We can refactor
189        // this once we have a clearer boundary between what sources/sinks need
190        // and the full "power" of the internal command flow, which should stay
191        // internal to the worker/not be exposed to source/sink implementations.
192        Self {
193            timely_worker,
194            client_rx,
195            storage_state,
196        }
197    }
198}
199
200impl StorageState {
201    /// Creates per-worker storage state, for hosting on any Timely worker.
202    ///
203    /// The caller provides the internal command channel endpoints: the native storage worker wires
204    /// them to the sequencer dataflow, a foreign host wires the sender to its own sequencing
205    /// channel and passes no receiver, since it dispatches internal commands itself.
206    /// Must be called on the hosting worker's thread, because the async worker unparks the
207    /// creating thread.
208    pub fn new_guest(
209        timely_worker_index: usize,
210        timely_worker_peers: usize,
211        internal_cmd_tx: InternalCommandSender,
212        internal_cmd_rx: Option<InternalCommandReceiver>,
213        metrics: StorageMetrics,
214        now: NowFn,
215        connection_context: ConnectionContext,
216        instance_context: StorageInstanceContext,
217        persist_clients: Arc<PersistClientCache>,
218        txns_ctx: TxnsContext,
219        tracing_handle: Arc<TracingHandle>,
220        shared_rocksdb_write_buffer_manager: SharedWriteBufferManager,
221    ) -> Self {
222        let storage_configuration =
223            StorageConfiguration::new(connection_context, mz_dyncfgs::all_dyncfgs());
224
225        // We always initialize as read_only=true. Only when we're explicitly
226        // allowed do we switch to doing writes.
227        let (read_only_tx, read_only_rx) = watch::channel(true);
228
229        // Similar to the internal command sequencer, it is very important that
230        // we only create the async worker once because a) the worker state is
231        // re-used when a new client connects and b) commands that have already
232        // been sent and might yield a response will be lost if a new iteration
233        // of `run_client` creates a new async worker.
234        //
235        // If we created a new async worker every time we get a new client
236        // (likely because the controller re-started and re-connected), we can
237        // get into an inconsistent state where we think that a dataflow has
238        // been rendered, for example because there is an entry in
239        // `StorageState::ingestions`, while there is not yet a dataflow. This
240        // happens because the dataflow only gets rendered once we get a
241        // response from the async worker and send off an internal command.
242        //
243        // The core idea is that both the sequencer and the async worker are
244        // part of the per-worker state, and must be treated as such, meaning
245        // they must survive between invocations of `run_client`.
246
247        // TODO(aljoscha): This thread unparking business seems brittle, but that's
248        // also how the command channel works currently. We can wrap it inside a
249        // struct that holds both a channel and a `Thread`, but I don't
250        // think that would help too much.
251        let async_worker = async_storage_worker::AsyncStorageWorker::new(
252            thread::current(),
253            Arc::clone(&persist_clients),
254        );
255        let cluster_memory_limit = instance_context.cluster_memory_limit;
256
257        let storage_state = StorageState {
258            source_uppers: BTreeMap::new(),
259            source_tokens: BTreeMap::new(),
260            metrics,
261            reported_frontiers: BTreeMap::new(),
262            ingestions: BTreeMap::new(),
263            exports: BTreeMap::new(),
264            oneshot_ingestions: BTreeMap::new(),
265            now,
266            timely_worker_index,
267            timely_worker_peers,
268            instance_context,
269            persist_clients,
270            txns_ctx,
271            sink_tokens: BTreeMap::new(),
272            sink_write_frontiers: BTreeMap::new(),
273            dropped_ids: Vec::new(),
274            aggregated_statistics: AggregatedStatistics::new(
275                timely_worker_index,
276                timely_worker_peers,
277            ),
278            shared_status_updates: Default::default(),
279            latest_status_updates: Default::default(),
280            initial_status_reported: Default::default(),
281            internal_cmd_tx,
282            internal_cmd_rx,
283            read_only_tx,
284            read_only_rx,
285            async_worker,
286            storage_configuration,
287            dataflow_parameters: DataflowParameters::new(
288                shared_rocksdb_write_buffer_manager,
289                cluster_memory_limit,
290            ),
291            tracing_handle,
292            server_maintenance_interval: Duration::ZERO,
293        };
294
295        storage_state
296    }
297}
298
299/// Worker-local state related to the ingress or egress of collections of data.
300pub struct StorageState {
301    /// The highest observed upper frontier for collection.
302    ///
303    /// This is shared among all source instances, so that they can jointly advance the
304    /// frontier even as other instances are created and dropped. Ideally, the Storage
305    /// module would eventually provide one source of truth on this rather than multiple,
306    /// and we should aim for that but are not there yet.
307    pub source_uppers: BTreeMap<GlobalId, Rc<RefCell<Antichain<mz_repr::Timestamp>>>>,
308    /// Handles to created sources, keyed by ID
309    /// NB: The type of the tokens must not be changed to something other than `PressOnDropButton`
310    /// to prevent usage of custom shutdown tokens that are tricky to get right.
311    pub source_tokens: BTreeMap<GlobalId, Vec<PressOnDropButton>>,
312    /// Metrics for storage objects.
313    pub metrics: StorageMetrics,
314    /// Tracks the conditional write frontiers we have reported.
315    pub reported_frontiers: BTreeMap<GlobalId, Antichain<Timestamp>>,
316    /// Descriptions of each installed ingestion.
317    pub ingestions: BTreeMap<GlobalId, IngestionDescription<CollectionMetadata>>,
318    /// Descriptions of each installed export.
319    pub exports: BTreeMap<GlobalId, StorageSinkDesc<CollectionMetadata, mz_repr::Timestamp>>,
320    /// Descriptions of oneshot ingestions that are currently running.
321    pub oneshot_ingestions: BTreeMap<uuid::Uuid, OneshotIngestionDescription<ProtoBatch>>,
322    /// Undocumented
323    pub now: NowFn,
324    /// Index of the associated timely dataflow worker.
325    pub timely_worker_index: usize,
326    /// Peers in the associated timely dataflow worker.
327    pub timely_worker_peers: usize,
328    /// Other configuration for sources and sinks.
329    pub instance_context: StorageInstanceContext,
330    /// A process-global cache of (blob_uri, consensus_uri) -> PersistClient.
331    /// This is intentionally shared between workers
332    pub persist_clients: Arc<PersistClientCache>,
333    /// Context necessary for rendering txn-wal operators.
334    pub txns_ctx: TxnsContext,
335    /// Tokens that should be dropped when a dataflow is dropped to clean up
336    /// associated state.
337    /// NB: The type of the tokens must not be changed to something other than `PressOnDropButton`
338    /// to prevent usage of custom shutdown tokens that are tricky to get right.
339    pub sink_tokens: BTreeMap<GlobalId, Vec<PressOnDropButton>>,
340    /// Frontier of sink writes (all subsequent writes will be at times at or
341    /// equal to this frontier)
342    pub sink_write_frontiers: BTreeMap<GlobalId, Rc<RefCell<Antichain<Timestamp>>>>,
343    /// Collection ids that have been dropped but not yet reported as dropped
344    pub dropped_ids: Vec<GlobalId>,
345
346    /// Statistics for sources and sinks.
347    pub aggregated_statistics: AggregatedStatistics,
348
349    /// A place shared with running dataflows, so that health operators, can
350    /// report status updates back to us.
351    ///
352    /// **NOTE**: Operators that append to this collection should take care to only add new
353    /// status updates if the status of the ingestion/export in question has _changed_.
354    pub shared_status_updates: Rc<RefCell<Vec<StatusUpdate>>>,
355
356    /// The latest status update for each object.
357    pub latest_status_updates: BTreeMap<GlobalId, StatusUpdate>,
358
359    /// Whether we have reported the initial status after connecting to a new client.
360    /// This is reset to false when a new client connects.
361    pub initial_status_reported: bool,
362
363    /// Sender for cluster-internal storage commands. These can be sent from
364    /// within workers/operators and will be distributed to all workers. For
365    /// example, for shutting down an entire dataflow from within a
366    /// operator/worker.
367    pub internal_cmd_tx: InternalCommandSender,
368    /// Receiver for cluster-internal storage commands. `None` when the state is hosted outside
369    /// the storage server, whose host dispatches internal commands itself.
370    pub internal_cmd_rx: Option<InternalCommandReceiver>,
371
372    /// When this replica/cluster is in read-only mode it must not affect any
373    /// changes to external state. This flag can only be changed by a
374    /// [StorageCommand::AllowWrites].
375    ///
376    /// Everything running on this replica/cluster must obey this flag. At the
377    /// time of writing, nothing currently looks at this flag.
378    /// TODO(benesch): fix this.
379    ///
380    /// NOTE: In the future, we might want a more complicated flag, for example
381    /// something that tells us after which timestamp we are allowed to write.
382    /// In this first version we are keeping things as simple as possible!
383    pub read_only_rx: watch::Receiver<bool>,
384
385    /// Send-side for read-only state.
386    pub read_only_tx: watch::Sender<bool>,
387
388    /// Async worker companion, used for running code that requires async, which
389    /// the timely main loop cannot do.
390    pub async_worker: AsyncStorageWorker<mz_repr::Timestamp>,
391
392    /// Configuration for source and sink connections.
393    pub storage_configuration: StorageConfiguration,
394    /// Dynamically configurable parameters that control how dataflows are rendered.
395    /// NOTE(guswynn): we should consider moving these into `storage_configuration`.
396    pub dataflow_parameters: DataflowParameters,
397
398    /// A process-global handle to tracing configuration.
399    pub tracing_handle: Arc<TracingHandle>,
400
401    /// Interval at which to perform server maintenance tasks. Set to a zero interval to
402    /// perform maintenance with every `step_or_park` invocation.
403    pub server_maintenance_interval: Duration,
404}
405
406impl StorageState {
407    /// Return an error handler that triggers a suspend and restart of the corresponding storage
408    /// dataflow.
409    pub fn error_handler(&self, context: &'static str, id: GlobalId) -> ErrorHandler {
410        let tx = self.internal_cmd_tx.clone();
411        ErrorHandler::signal(move |e| {
412            tx.send(InternalStorageCommand::SuspendAndRestart {
413                id,
414                reason: format!("{context}: {e:#}"),
415            })
416        })
417    }
418}
419
420/// Extra context for a storage instance.
421/// This is extra information that is used when rendering source
422/// and sinks that is not tied to the source/connection configuration itself.
423#[derive(Clone)]
424pub struct StorageInstanceContext {
425    /// A directory that can be used for scratch work.
426    pub scratch_directory: Option<PathBuf>,
427    /// The memory limit of the materialize cluster replica. This will
428    /// be used to calculate and configure the maximum inflight bytes for backpressure
429    pub cluster_memory_limit: Option<usize>,
430}
431
432impl StorageInstanceContext {
433    /// Build a new `StorageInstanceContext`.
434    pub fn new(scratch_directory: Option<PathBuf>, cluster_memory_limit: Option<usize>) -> Self {
435        Self {
436            scratch_directory,
437            cluster_memory_limit,
438        }
439    }
440
441    /// Returns a `rocksdb::Env` for a new RocksDB instance.
442    ///
443    /// With a scratch directory this is the default `Env`, which stores data
444    /// on the host filesystem. Without one, RocksDB runs in memory, and every
445    /// call returns a fresh in-memory `Env`. State written through an `Env`
446    /// is only reachable through that same `Env`, so a per-instance `Env`
447    /// isolates instances from each other and from previous incarnations of
448    /// themselves. Background threads are process-wide either way, both
449    /// variants delegate them to the default `Env`.
450    pub fn rocksdb_env(&self) -> Result<rocksdb::Env, rocksdb::Error> {
451        if self.scratch_directory.is_some() {
452            rocksdb::Env::new()
453        } else {
454            rocksdb::Env::mem_env()
455        }
456    }
457}
458
459impl<'w> Worker<'w> {
460    /// Waits for client connections and runs them to completion.
461    pub fn run(&mut self) {
462        while let Some((_nonce, rx, tx)) = self.client_rx.blocking_recv() {
463            self.run_client(rx, tx);
464        }
465    }
466
467    /// Runs this (timely) storage worker until the given `command_rx` is
468    /// disconnected.
469    ///
470    /// See the [module documentation](crate::storage_state) for this
471    /// workers responsibilities, how it communicates with the other workers and
472    /// how commands flow from the controller and through the workers.
473    fn run_client(&mut self, mut command_rx: CommandReceiver, response_tx: ResponseSender) {
474        // At this point, all workers are still reading from the command flow.
475        if self.reconcile(&mut command_rx).is_err() {
476            return;
477        }
478
479        // The last time we reported statistics.
480        let mut last_stats_time = Instant::now();
481
482        // The last time we did periodic maintenance.
483        let mut last_maintenance = std::time::Instant::now();
484
485        let mut disconnected = false;
486        while !disconnected {
487            let config = &self.storage_state.storage_configuration;
488            let stats_interval = config.parameters.statistics_collection_interval;
489
490            let maintenance_interval = self.storage_state.server_maintenance_interval;
491
492            let now = std::time::Instant::now();
493            // Determine if we need to perform maintenance, which is true if `maintenance_interval`
494            // time has passed since the last maintenance.
495            let sleep_duration;
496            if now >= last_maintenance + maintenance_interval {
497                last_maintenance = now;
498                sleep_duration = None;
499
500                self.report_frontier_progress(&response_tx);
501            } else {
502                // We didn't perform maintenance, sleep until the next maintenance interval.
503                let next_maintenance = last_maintenance + maintenance_interval;
504                sleep_duration = Some(next_maintenance.saturating_duration_since(now))
505            }
506
507            // Ask Timely to execute a unit of work.
508            //
509            // If there are no pending commands or responses from the async
510            // worker, we ask Timely to park the thread if there's nothing to
511            // do. We rely on another thread unparking us when there's new work
512            // to be done, e.g., when sending a command or when new Kafka
513            // messages have arrived.
514            //
515            // It is critical that we allow Timely to park iff there are no
516            // pending commands or responses. The command may have already been
517            // consumed by the call to `client_rx.recv`. See:
518            // https://github.com/MaterializeInc/materialize/pull/13973#issuecomment-1200312212
519            if command_rx.is_empty() && self.storage_state.async_worker.is_empty() {
520                // Make sure we wake up again to report any pending statistics updates.
521                let mut park_duration = stats_interval.saturating_sub(last_stats_time.elapsed());
522                if let Some(sleep_duration) = sleep_duration {
523                    park_duration = std::cmp::min(sleep_duration, park_duration);
524                }
525                self.timely_worker.step_or_park(Some(park_duration));
526            } else {
527                self.timely_worker.step();
528            }
529
530            // Rerport any dropped ids
531            for id in std::mem::take(&mut self.storage_state.dropped_ids) {
532                self.send_storage_response(&response_tx, StorageResponse::DroppedId(id));
533            }
534
535            self.process_oneshot_ingestions(&response_tx);
536
537            self.report_status_updates(&response_tx);
538
539            if last_stats_time.elapsed() >= stats_interval {
540                self.report_storage_statistics(&response_tx);
541                last_stats_time = Instant::now();
542            }
543
544            // Handle any received commands.
545            loop {
546                match command_rx.try_recv() {
547                    Ok(cmd) => self.storage_state.handle_storage_command(cmd),
548                    Err(TryRecvError::Empty) => break,
549                    Err(TryRecvError::Disconnected) => {
550                        disconnected = true;
551                        break;
552                    }
553                }
554            }
555
556            // Handle responses from the async worker.
557            while let Ok(response) = self.storage_state.async_worker.try_recv() {
558                self.handle_async_worker_response(response);
559            }
560
561            // Handle any received commands.
562            while let Some(command) = self
563                .storage_state
564                .internal_cmd_rx
565                .as_ref()
566                .expect("storage server always wires a receiver")
567                .try_recv()
568            {
569                self.handle_internal_storage_command(command);
570            }
571        }
572    }
573
574    /// Entry point for applying a response from the async storage worker.
575    pub fn handle_async_worker_response(
576        &self,
577        async_response: AsyncStorageWorkerResponse<mz_repr::Timestamp>,
578    ) {
579        // NOTE: If we want to share the load of async processing we
580        // have to change `handle_storage_command` and change this
581        // assert.
582        assert_eq!(
583            self.timely_worker.index(),
584            0,
585            "only worker #0 is doing async processing"
586        );
587        match async_response {
588            AsyncStorageWorkerResponse::IngestionFrontiersUpdated {
589                id,
590                ingestion_description,
591                as_of,
592                resume_uppers,
593                source_resume_uppers,
594            } => {
595                self.storage_state.internal_cmd_tx.send(
596                    InternalStorageCommand::CreateIngestionDataflow {
597                        id,
598                        ingestion_description,
599                        as_of,
600                        resume_uppers,
601                        source_resume_uppers,
602                    },
603                );
604            }
605            AsyncStorageWorkerResponse::ExportFrontiersUpdated { id, description } => {
606                self.storage_state
607                    .internal_cmd_tx
608                    .send(InternalStorageCommand::RunSinkDataflow(id, description));
609            }
610            AsyncStorageWorkerResponse::DropDataflow(id) => {
611                self.storage_state
612                    .internal_cmd_tx
613                    .send(InternalStorageCommand::DropDataflow(vec![id]));
614            }
615        }
616    }
617
618    /// Entry point for applying an internal storage command.
619    pub fn handle_internal_storage_command(&mut self, internal_cmd: InternalStorageCommand) {
620        match internal_cmd {
621            InternalStorageCommand::SuspendAndRestart { id, reason } => {
622                info!(
623                    "worker {}/{} initiating suspend-and-restart for {id} because of: {reason}",
624                    self.timely_worker.index(),
625                    self.timely_worker.peers(),
626                );
627
628                let maybe_ingestion = self.storage_state.ingestions.get(&id).cloned();
629                if let Some(ingestion_description) = maybe_ingestion {
630                    // Yank the token of the previously existing source dataflow.Note that this
631                    // token also includes any source exports/subsources.
632                    let maybe_token = self.storage_state.source_tokens.remove(&id);
633                    if maybe_token.is_none() {
634                        // Something has dropped the source. Make sure we don't
635                        // accidentally re-create it.
636                        return;
637                    }
638
639                    // This needs to be done by one worker, which will
640                    // broadcasts a `CreateIngestionDataflow` command to all
641                    // workers based on the response that contains the
642                    // resumption upper.
643                    //
644                    // Doing this separately on each worker could lead to
645                    // differing resume_uppers which might lead to all kinds of
646                    // mayhem.
647                    //
648                    // TODO(aljoscha): If we ever become worried that this is
649                    // putting undue pressure on worker 0 we can pick the
650                    // designated worker for a source/sink based on `id.hash()`.
651                    if self.timely_worker.index() == 0 {
652                        for (id, _) in ingestion_description.source_exports.iter() {
653                            self.storage_state
654                                .aggregated_statistics
655                                .advance_global_epoch(*id);
656                        }
657                        self.storage_state
658                            .async_worker
659                            .update_ingestion_frontiers(id, ingestion_description);
660                    }
661
662                    // Continue with other commands.
663                    return;
664                }
665
666                let maybe_sink = self.storage_state.exports.get(&id).cloned();
667                if let Some(sink_description) = maybe_sink {
668                    // Yank the token of the previously existing sink
669                    // dataflow.
670                    let maybe_token = self.storage_state.sink_tokens.remove(&id);
671
672                    if maybe_token.is_none() {
673                        // Something has dropped the sink. Make sure we don't
674                        // accidentally re-create it.
675                        return;
676                    }
677
678                    // This needs to be broadcast by one worker and go through
679                    // the internal command fabric, to ensure consistent
680                    // ordering of dataflow rendering across all workers.
681                    if self.timely_worker.index() == 0 {
682                        self.storage_state
683                            .aggregated_statistics
684                            .advance_global_epoch(id);
685                        self.storage_state
686                            .async_worker
687                            .update_sink_frontiers(id, sink_description);
688                    }
689
690                    // Continue with other commands.
691                    return;
692                }
693
694                if !self
695                    .storage_state
696                    .ingestions
697                    .values()
698                    .any(|v| v.source_exports.contains_key(&id))
699                {
700                    // Our current approach to dropping a source results in a race between shard
701                    // finalization (which happens in the controller) and dataflow shutdown (which
702                    // happens in clusterd). If a source is created and dropped fast enough -or the
703                    // two commands get sufficiently delayed- then it's possible to receive a
704                    // SuspendAndRestart command for an unknown source. We cannot assert that this
705                    // never happens but we log an error here to track how often this happens.
706                    warn!(
707                        "got InternalStorageCommand::SuspendAndRestart for something that is not a source or sink: {id}"
708                    );
709                }
710            }
711            InternalStorageCommand::CreateIngestionDataflow {
712                id: ingestion_id,
713                mut ingestion_description,
714                as_of,
715                mut resume_uppers,
716                mut source_resume_uppers,
717            } => {
718                info!(
719                    ?as_of,
720                    ?resume_uppers,
721                    "worker {}/{} trying to (re-)start ingestion {ingestion_id}",
722                    self.timely_worker.index(),
723                    self.timely_worker.peers(),
724                );
725
726                // We initialize statistics before we prune finished exports. We
727                // still want to export statistics for these, plus the rendering
728                // machinery will get confused if there are not at least
729                // statistics for the "main" source.
730                for (export_id, export) in ingestion_description.source_exports.iter() {
731                    let resume_upper = resume_uppers[export_id].clone();
732                    self.storage_state.aggregated_statistics.initialize_source(
733                        *export_id,
734                        ingestion_id,
735                        resume_upper.clone(),
736                        || {
737                            SourceStatistics::new(
738                                *export_id,
739                                self.storage_state.timely_worker_index,
740                                &self.storage_state.metrics.source_statistics,
741                                ingestion_id,
742                                &export.storage_metadata.data_shard,
743                                export.data_config.envelope.clone(),
744                                resume_upper,
745                            )
746                        },
747                    );
748                }
749
750                let finished_exports: BTreeSet<GlobalId> = resume_uppers
751                    .iter()
752                    .filter(|(_, frontier)| frontier.is_empty())
753                    .map(|(id, _)| *id)
754                    .collect();
755
756                resume_uppers.retain(|id, _| !finished_exports.contains(id));
757                source_resume_uppers.retain(|id, _| !finished_exports.contains(id));
758                ingestion_description
759                    .source_exports
760                    .retain(|id, _| !finished_exports.contains(id));
761
762                for id in ingestion_description.collection_ids() {
763                    // If there is already a shared upper, we re-use it, to make
764                    // sure that parties that are already using the shared upper
765                    // can continue doing so.
766                    let source_upper = self
767                        .storage_state
768                        .source_uppers
769                        .entry(id.clone())
770                        .or_insert_with(|| {
771                            Rc::new(RefCell::new(Antichain::from_elem(Timestamp::minimum())))
772                        });
773
774                    let mut source_upper = source_upper.borrow_mut();
775                    if !source_upper.is_empty() {
776                        source_upper.clear();
777                        source_upper.insert(mz_repr::Timestamp::minimum());
778                    }
779                }
780
781                // If all subsources of the source are finished, we can skip rendering entirely.
782                // Also, if `as_of` is empty, the dataflow has been finalized, so we can skip it as
783                // well.
784                //
785                // TODO(guswynn|petrosagg): this is a bit hacky, and is a consequence of storage state
786                // management being a bit of a mess. we should clean this up and remove weird if
787                // statements like this.
788                if resume_uppers.values().all(|frontier| frontier.is_empty()) || as_of.is_empty() {
789                    info!(
790                        ?resume_uppers,
791                        ?as_of,
792                        "worker {}/{} skipping building ingestion dataflow \
793                        for {ingestion_id} because the ingestion is finished",
794                        self.timely_worker.index(),
795                        self.timely_worker.peers(),
796                    );
797                    return;
798                }
799
800                crate::render::build_ingestion_dataflow(
801                    self.timely_worker,
802                    &mut self.storage_state,
803                    ingestion_id,
804                    ingestion_description,
805                    as_of,
806                    resume_uppers,
807                    source_resume_uppers,
808                );
809            }
810            InternalStorageCommand::RunOneshotIngestion {
811                ingestion_id,
812                collection_id,
813                collection_meta,
814                request,
815            } => {
816                crate::render::build_oneshot_ingestion_dataflow(
817                    self.timely_worker,
818                    &mut self.storage_state,
819                    ingestion_id,
820                    collection_id,
821                    collection_meta,
822                    request,
823                );
824            }
825            InternalStorageCommand::RunSinkDataflow(sink_id, sink_description) => {
826                info!(
827                    "worker {}/{} trying to (re-)start sink {sink_id}",
828                    self.timely_worker.index(),
829                    self.timely_worker.peers(),
830                );
831
832                {
833                    // If there is already a shared write frontier, we re-use it, to
834                    // make sure that parties that are already using the shared
835                    // frontier can continue doing so.
836                    let sink_write_frontier = self
837                        .storage_state
838                        .sink_write_frontiers
839                        .entry(sink_id.clone())
840                        .or_insert_with(|| Rc::new(RefCell::new(Antichain::new())));
841
842                    let mut sink_write_frontier = sink_write_frontier.borrow_mut();
843                    sink_write_frontier.clear();
844                    sink_write_frontier.insert(mz_repr::Timestamp::minimum());
845                }
846                self.storage_state
847                    .aggregated_statistics
848                    .initialize_sink(sink_id, || {
849                        SinkStatistics::new(
850                            sink_id,
851                            self.storage_state.timely_worker_index,
852                            &self.storage_state.metrics.sink_statistics,
853                        )
854                    });
855
856                crate::render::build_export_dataflow(
857                    self.timely_worker,
858                    &mut self.storage_state,
859                    sink_id,
860                    sink_description,
861                );
862            }
863            InternalStorageCommand::DropDataflow(ids) => {
864                for id in &ids {
865                    // Clean up per-source / per-sink state.
866                    self.storage_state.source_uppers.remove(id);
867                    self.storage_state.source_tokens.remove(id);
868
869                    self.storage_state.sink_tokens.remove(id);
870                    self.storage_state.sink_write_frontiers.remove(id);
871
872                    self.storage_state.aggregated_statistics.deinitialize(*id);
873                }
874            }
875            InternalStorageCommand::UpdateConfiguration { storage_parameters } => {
876                self.storage_state
877                    .dataflow_parameters
878                    .update(storage_parameters.clone());
879                self.storage_state
880                    .storage_configuration
881                    .update(storage_parameters);
882
883                // Clear out the updates as we no longer forward them to anyone else to process.
884                // We clone `StorageState::storage_configuration` many times during rendering
885                // and want to avoid cloning these unused updates.
886                self.storage_state
887                    .storage_configuration
888                    .parameters
889                    .dyncfg_updates = Default::default();
890
891                // Remember the maintenance interval locally to avoid reading it from the config set on
892                // every server iteration.
893                self.storage_state.server_maintenance_interval =
894                    STORAGE_SERVER_MAINTENANCE_INTERVAL
895                        .get(self.storage_state.storage_configuration.config_set());
896
897                // Apply storage's upsert spill flag to both stash flavors'
898                // mechanisms: the storage leg of the process-wide chunk
899                // spill gate (chunked flavor) and the storage-owned column
900                // pager (paged flavor). The buffer pool, the pager pool, and
901                // their budgets are the shared ones configured by compute's
902                // `apply_worker_config` (compute and storage run in the same
903                // process). The chunk gate ORs storage's leg with compute's,
904                // so chunks spill while either subsystem's flag is set.
905                //
906                // The flag is replica-scoped: the storage controller merges
907                // per-replica overrides into the `UpdateConfiguration`
908                // commands it sends, so reading this worker's `ConfigSet`
909                // here observes them.
910                {
911                    use mz_storage_types::dyncfgs::ENABLE_UPSERT_PAGED_SPILL;
912
913                    let enabled = ENABLE_UPSERT_PAGED_SPILL
914                        .get(self.storage_state.storage_configuration.config_set());
915                    debug!(
916                        worker = self.timely_worker.index(),
917                        enabled, "upsert stash spill: applying gate",
918                    );
919                    crate::upsert::upsert_stash_spill::set_enabled(enabled);
920                    crate::upsert::upsert_stash_pager::set_enabled(enabled);
921                }
922            }
923            InternalStorageCommand::StatisticsUpdate { sources, sinks } => self
924                .storage_state
925                .aggregated_statistics
926                .ingest(sources, sinks),
927        }
928    }
929
930    /// Emit information about write frontier progress, along with information that should
931    /// be made durable for this to be the case.
932    ///
933    /// The write frontier progress is "conditional" in that it is not until the information is made
934    /// durable that the data are emitted to downstream workers, and indeed they should not rely on
935    /// the completeness of what they hear until the information is made durable.
936    ///
937    /// Specifically, this sends information about new timestamp bindings created by dataflow workers,
938    /// with the understanding if that if made durable (and ack'd back to the workers) the source will
939    /// in fact progress with this write frontier.
940    pub fn report_frontier_progress(&mut self, response_tx: &ResponseSender) {
941        let mut new_uppers = Vec::new();
942
943        // Check if any observed frontier should advance the reported frontiers.
944        for (id, frontier) in self
945            .storage_state
946            .source_uppers
947            .iter()
948            .chain(self.storage_state.sink_write_frontiers.iter())
949        {
950            let Some(reported_frontier) = self.storage_state.reported_frontiers.get_mut(id) else {
951                // Frontier reporting has not yet been started for this object.
952                // Potentially because this timely worker has not yet seen the
953                // `CreateSources` command.
954                continue;
955            };
956
957            let observed_frontier = frontier.borrow();
958
959            // Only do a thing if it *advances* the frontier, not just *changes* the frontier.
960            // This is protection against `frontier` lagging behind what we have conditionally reported.
961            if PartialOrder::less_than(reported_frontier, &observed_frontier) {
962                new_uppers.push((*id, observed_frontier.clone()));
963                reported_frontier.clone_from(&observed_frontier);
964            }
965        }
966
967        for (id, upper) in new_uppers {
968            self.send_storage_response(response_tx, StorageResponse::FrontierUpper(id, upper));
969        }
970    }
971
972    /// Pumps latest status updates from the buffer shared with operators and
973    /// reports any updates that need reporting.
974    pub fn report_status_updates(&mut self, response_tx: &ResponseSender) {
975        // If we haven't done the initial status report, report all current statuses
976        if !self.storage_state.initial_status_reported {
977            // We pull initially reported status updates to "now", so that they
978            // sort as the latest update in internal status collections. This
979            // makes it so that a newly bootstrapped envd can append status
980            // updates to internal status collections that report an accurate
981            // view as of the time when they came up.
982            let now_ts = mz_ore::now::to_datetime((self.storage_state.now)());
983            let status_updates = self
984                .storage_state
985                .latest_status_updates
986                .values()
987                .cloned()
988                .map(|mut update| {
989                    update.timestamp = now_ts.clone();
990                    update
991                });
992            for update in status_updates {
993                self.send_storage_response(response_tx, StorageResponse::StatusUpdate(update));
994            }
995            self.storage_state.initial_status_reported = true;
996        }
997
998        // Pump updates into our state and stage them for reporting.
999        for shared_update in self.storage_state.shared_status_updates.take() {
1000            self.send_storage_response(
1001                response_tx,
1002                StorageResponse::StatusUpdate(shared_update.clone()),
1003            );
1004
1005            self.storage_state
1006                .latest_status_updates
1007                .insert(shared_update.id, shared_update);
1008        }
1009    }
1010
1011    /// Report source statistics back to the controller.
1012    pub fn report_storage_statistics(&mut self, response_tx: &ResponseSender) {
1013        let (sources, sinks) = self.storage_state.aggregated_statistics.emit_local();
1014        if !sources.is_empty() || !sinks.is_empty() {
1015            self.storage_state
1016                .internal_cmd_tx
1017                .send(InternalStorageCommand::StatisticsUpdate { sources, sinks })
1018        }
1019
1020        let (sources, sinks) = self.storage_state.aggregated_statistics.snapshot();
1021        if !sources.is_empty() || !sinks.is_empty() {
1022            self.send_storage_response(
1023                response_tx,
1024                StorageResponse::StatisticsUpdates(sources, sinks),
1025            );
1026        }
1027    }
1028
1029    /// Send a response to the coordinator.
1030    pub fn send_storage_response(&self, response_tx: &ResponseSender, response: StorageResponse) {
1031        // Ignore send errors because the coordinator is free to ignore our
1032        // responses. This happens during shutdown.
1033        let _ = response_tx.send(response);
1034    }
1035
1036    /// Forward completed oneshot ingestion results to the coordinator.
1037    pub fn process_oneshot_ingestions(&mut self, response_tx: &ResponseSender) {
1038        for (ingestion_id, ingestion_state) in &mut self.storage_state.oneshot_ingestions {
1039            loop {
1040                match ingestion_state.results.try_recv() {
1041                    Ok(result) => {
1042                        let response = match result {
1043                            Ok(maybe_batch) => maybe_batch.into_iter().map(Result::Ok).collect(),
1044                            Err(err) => vec![Err(err)],
1045                        };
1046                        let staged_batches = BTreeMap::from([(*ingestion_id, response)]);
1047                        let _ = response_tx.send(StorageResponse::StagedBatches(staged_batches));
1048                    }
1049                    Err(TryRecvError::Empty) => {
1050                        break;
1051                    }
1052                    Err(TryRecvError::Disconnected) => {
1053                        break;
1054                    }
1055                }
1056            }
1057        }
1058    }
1059
1060    /// Extract commands until `InitializationComplete`, and make the worker
1061    /// reflect those commands. If the worker can not be made to reflect the
1062    /// commands, return an error.
1063    fn reconcile(&mut self, command_rx: &mut CommandReceiver) -> Result<(), ()> {
1064        // To initialize the connection, we want to drain all commands until we
1065        // receive a `StorageCommand::InitializationComplete` command to form a
1066        // target command state.
1067        let mut commands = vec![];
1068        loop {
1069            match command_rx.blocking_recv().ok_or(())? {
1070                StorageCommand::InitializationComplete => break,
1071                command => commands.push(command),
1072            }
1073        }
1074
1075        self.reconcile_commands(commands);
1076        Ok(())
1077    }
1078
1079    /// Reconciles the worker state with the given target command state,
1080    /// which the caller has drained from a new client connection up to (exclusive) the
1081    /// `InitializationComplete` marker.
1082    pub fn reconcile_commands(&mut self, mut commands: Vec<StorageCommand>) {
1083        let worker_id = self.timely_worker.index();
1084
1085        // Track which frontiers this envd expects; we will also set their
1086        // initial timestamp to the minimum timestamp to reset them as we don't
1087        // know what frontiers the new envd expects.
1088        let mut expected_objects = BTreeSet::new();
1089
1090        let mut drop_commands = BTreeSet::new();
1091        let mut running_ingestion_descriptions = self.storage_state.ingestions.clone();
1092        let mut running_exports_descriptions = self.storage_state.exports.clone();
1093
1094        let mut create_oneshot_ingestions: BTreeSet<Uuid> = BTreeSet::new();
1095        let mut cancel_oneshot_ingestions: BTreeSet<Uuid> = BTreeSet::new();
1096
1097        for command in &mut commands {
1098            match command {
1099                StorageCommand::Hello { .. } => {
1100                    panic!("Hello must be captured before")
1101                }
1102                StorageCommand::AllowCompaction(id, since) => {
1103                    info!(%worker_id, ?id, ?since, "reconcile: received AllowCompaction command");
1104
1105                    // collect all "drop commands". These are `AllowCompaction`
1106                    // commands that compact to the empty since. Then, later, we make sure
1107                    // we retain only those `Create*` commands that are not dropped. We
1108                    // assume that the `AllowCompaction` command is ordered after the
1109                    // `Create*` commands but don't assert that.
1110                    // WIP: Should we assert?
1111                    if since.is_empty() {
1112                        drop_commands.insert(*id);
1113                    }
1114                }
1115                StorageCommand::RunIngestion(ingestion) => {
1116                    info!(%worker_id, ?ingestion, "reconcile: received RunIngestion command");
1117
1118                    // Ensure that ingestions are forward-rolling alter compatible.
1119                    let prev = running_ingestion_descriptions
1120                        .insert(ingestion.id, ingestion.description.clone());
1121
1122                    if let Some(prev_ingest) = prev {
1123                        // If the new ingestion is not exactly equal to the currently running
1124                        // ingestion, we must either track that we need to synthesize an update
1125                        // command to change the ingestion, or panic.
1126                        prev_ingest
1127                            .alter_compatible(ingestion.id, &ingestion.description)
1128                            .expect("only alter compatible ingestions permitted");
1129                    }
1130                }
1131                StorageCommand::RunSink(export) => {
1132                    info!(%worker_id, ?export, "reconcile: received RunSink command");
1133
1134                    // Ensure that exports are forward-rolling alter compatible.
1135                    let prev =
1136                        running_exports_descriptions.insert(export.id, export.description.clone());
1137
1138                    if let Some(prev_export) = prev {
1139                        prev_export
1140                            .alter_compatible(export.id, &export.description)
1141                            .expect("only alter compatible exports permitted");
1142                    }
1143                }
1144                StorageCommand::RunOneshotIngestion(ingestion) => {
1145                    info!(
1146                        %worker_id,
1147                        ingestion_id = %ingestion.ingestion_id,
1148                        collection_id = %ingestion.collection_id,
1149                        "reconcile: received RunOneshotIngestion command",
1150                    );
1151                    create_oneshot_ingestions.insert(ingestion.ingestion_id);
1152                }
1153                StorageCommand::CancelOneshotIngestion(uuid) => {
1154                    info!(%worker_id, %uuid, "reconcile: received CancelOneshotIngestion command");
1155                    cancel_oneshot_ingestions.insert(*uuid);
1156                }
1157                StorageCommand::InitializationComplete
1158                | StorageCommand::AllowWrites
1159                | StorageCommand::UpdateConfiguration(_) => (),
1160            }
1161        }
1162
1163        let mut seen_most_recent_definition = BTreeSet::new();
1164
1165        // We iterate over this backward to ensure that we keep only the most recent ingestion
1166        // description.
1167        let mut filtered_commands = VecDeque::new();
1168        for mut command in commands.into_iter().rev() {
1169            let mut should_keep = true;
1170            match &mut command {
1171                StorageCommand::Hello { .. } => {
1172                    panic!("Hello must be captured before")
1173                }
1174                StorageCommand::RunIngestion(ingestion) => {
1175                    // Subsources can be dropped independently of their
1176                    // primary source, so we evaluate them in a separate
1177                    // loop.
1178                    for export_id in ingestion
1179                        .description
1180                        .source_exports
1181                        .keys()
1182                        .filter(|export_id| **export_id != ingestion.id)
1183                    {
1184                        if drop_commands.remove(export_id) {
1185                            info!(%worker_id, %export_id, "reconcile: dropping subsource");
1186                            self.storage_state.dropped_ids.push(*export_id);
1187                        }
1188                    }
1189
1190                    if drop_commands.remove(&ingestion.id)
1191                        || self.storage_state.dropped_ids.contains(&ingestion.id)
1192                    {
1193                        info!(%worker_id, %ingestion.id, "reconcile: dropping ingestion");
1194
1195                        // If an ingestion is dropped, so too must all of
1196                        // its subsources (i.e. ingestion exports, as well
1197                        // as its progress subsource).
1198                        for id in ingestion.description.collection_ids() {
1199                            drop_commands.remove(&id);
1200                            self.storage_state.dropped_ids.push(id);
1201                        }
1202                        should_keep = false;
1203                    } else {
1204                        let most_recent_defintion =
1205                            seen_most_recent_definition.insert(ingestion.id);
1206
1207                        if most_recent_defintion {
1208                            // If this is the most recent definition, this
1209                            // is what we will be running when
1210                            // reconciliation completes. This definition
1211                            // must not include any dropped subsources.
1212                            ingestion.description.source_exports.retain(|export_id, _| {
1213                                !self.storage_state.dropped_ids.contains(export_id)
1214                            });
1215
1216                            // After clearing any dropped subsources, we can
1217                            // state that we expect all of these to exist.
1218                            expected_objects.extend(ingestion.description.collection_ids());
1219                        }
1220
1221                        let running_ingestion = self.storage_state.ingestions.get(&ingestion.id);
1222
1223                        // We keep only:
1224                        // - The most recent version of the ingestion, which
1225                        //   is why these commands are run in reverse.
1226                        // - Ingestions whose descriptions are not exactly
1227                        //   those that are currently running.
1228                        should_keep = most_recent_defintion
1229                            && running_ingestion != Some(&ingestion.description)
1230                    }
1231                }
1232                StorageCommand::RunSink(export) => {
1233                    if drop_commands.remove(&export.id)
1234                        // If there were multiple `RunSink` in the command
1235                        // stream, we want to ensure none of them are
1236                        // retained.
1237                        || self.storage_state.dropped_ids.contains(&export.id)
1238                    {
1239                        info!(%worker_id, %export.id, "reconcile: dropping sink");
1240
1241                        // Make sure that we report back that the ID was
1242                        // dropped.
1243                        self.storage_state.dropped_ids.push(export.id);
1244
1245                        should_keep = false
1246                    } else {
1247                        expected_objects.insert(export.id);
1248
1249                        let running_sink = self.storage_state.exports.get(&export.id);
1250
1251                        // We keep only:
1252                        // - The most recent version of the sink, which
1253                        //   is why these commands are run in reverse.
1254                        // - Sinks whose descriptions are not exactly
1255                        //   those that are currently running.
1256                        should_keep = seen_most_recent_definition.insert(export.id)
1257                            && running_sink != Some(&export.description);
1258                    }
1259                }
1260                StorageCommand::RunOneshotIngestion(ingestion) => {
1261                    let already_running = self
1262                        .storage_state
1263                        .oneshot_ingestions
1264                        .contains_key(&ingestion.ingestion_id);
1265                    let was_canceled = cancel_oneshot_ingestions.contains(&ingestion.ingestion_id);
1266
1267                    should_keep = !already_running && !was_canceled;
1268                }
1269                StorageCommand::CancelOneshotIngestion(ingestion_id) => {
1270                    let already_running = self
1271                        .storage_state
1272                        .oneshot_ingestions
1273                        .contains_key(ingestion_id);
1274                    should_keep = already_running;
1275                }
1276                StorageCommand::InitializationComplete
1277                | StorageCommand::AllowWrites
1278                | StorageCommand::UpdateConfiguration(_)
1279                | StorageCommand::AllowCompaction(_, _) => (),
1280            }
1281            if should_keep {
1282                filtered_commands.push_front(command);
1283            }
1284        }
1285        let commands = filtered_commands;
1286
1287        // Make sure all the "drop commands" matched up with a source or sink.
1288        // This is also what the regular handler logic for `AllowCompaction`
1289        // would do.
1290        soft_assert_or_log!(
1291            drop_commands.is_empty(),
1292            "AllowCompaction commands for non-existent IDs {:?}",
1293            drop_commands
1294        );
1295
1296        // Determine the ID of all objects we did _not_ see; these are
1297        // considered stale.
1298        let stale_objects = self
1299            .storage_state
1300            .ingestions
1301            .values()
1302            .map(|i| i.collection_ids())
1303            .flatten()
1304            .chain(self.storage_state.exports.keys().copied())
1305            // Objects are considered stale if we did not see them re-created.
1306            .filter(|id| !expected_objects.contains(id))
1307            .collect::<Vec<_>>();
1308        let stale_oneshot_ingestions = self
1309            .storage_state
1310            .oneshot_ingestions
1311            .keys()
1312            .filter(|ingestion_id| {
1313                let to_create = create_oneshot_ingestions.contains(ingestion_id);
1314                let to_drop = cancel_oneshot_ingestions.contains(ingestion_id);
1315                mz_ore::soft_assert_or_log!(
1316                    !(!to_create && to_drop),
1317                    "attempting to drop oneshot source {ingestion_id} that is not expected to be created during reconciliation"
1318                );
1319                !to_create && !to_drop
1320            })
1321            .copied()
1322            .collect::<Vec<_>>();
1323
1324        info!(
1325            %worker_id, ?expected_objects, ?stale_objects, ?stale_oneshot_ingestions,
1326            "reconcile: modifing storage state to match expected objects",
1327        );
1328
1329        for id in stale_objects {
1330            self.storage_state.drop_collection(id);
1331        }
1332        for id in stale_oneshot_ingestions {
1333            self.storage_state.drop_oneshot_ingestion(id);
1334        }
1335
1336        // Do not report dropping any objects that do not belong to expected
1337        // objects.
1338        self.storage_state
1339            .dropped_ids
1340            .retain(|id| expected_objects.contains(id));
1341
1342        // Do not report any frontiers that do not belong to expected objects.
1343        // Note that this set of objects can differ from the set of sources and
1344        // sinks.
1345        self.storage_state
1346            .reported_frontiers
1347            .retain(|id, _| expected_objects.contains(id));
1348
1349        // Reset the reported frontiers for the remaining objects.
1350        for (_, frontier) in &mut self.storage_state.reported_frontiers {
1351            *frontier = Antichain::from_elem(<_>::minimum());
1352        }
1353
1354        // Reset the initial status reported flag when a new client connects
1355        self.storage_state.initial_status_reported = false;
1356
1357        // Execute the modified commands.
1358        for command in commands {
1359            self.storage_state.handle_storage_command(command);
1360        }
1361    }
1362}
1363
1364impl StorageState {
1365    /// Entry point for applying a storage command.
1366    ///
1367    /// NOTE: This does not have access to the timely worker and therefore
1368    /// cannot render dataflows. For dataflow rendering, this needs to either
1369    /// send asynchronous command to the `async_worker` or internal
1370    /// commands to the `internal_cmd_tx`.
1371    pub fn handle_storage_command(&mut self, cmd: StorageCommand) {
1372        match cmd {
1373            StorageCommand::Hello { .. } => panic!("Hello must be captured before"),
1374            StorageCommand::InitializationComplete => (),
1375            StorageCommand::AllowWrites => {
1376                self.read_only_tx
1377                    .send(false)
1378                    .expect("we're holding one other end");
1379                self.persist_clients.cfg().enable_compaction();
1380            }
1381            StorageCommand::UpdateConfiguration(params) => {
1382                // These can be done from all workers safely.
1383                debug!("Applying configuration update: {params:?}");
1384
1385                // We serialize the dyncfg updates in StorageParameters, but configure
1386                // persist separately.
1387                self.persist_clients
1388                    .cfg()
1389                    .apply_from(&params.dyncfg_updates);
1390
1391                params.tracing.apply(self.tracing_handle.as_ref());
1392
1393                if let Some(log_filter) = &params.tracing.log_filter {
1394                    self.storage_configuration
1395                        .connection_context
1396                        .librdkafka_log_level =
1397                        mz_ore::tracing::crate_level(&log_filter.clone().into(), "librdkafka");
1398                }
1399
1400                // This needs to be broadcast by one worker and go through
1401                // the internal command fabric, to ensure consistent
1402                // ordering of dataflow rendering across all workers.
1403                if self.timely_worker_index == 0 {
1404                    self.internal_cmd_tx
1405                        .send(InternalStorageCommand::UpdateConfiguration {
1406                            storage_parameters: *params,
1407                        })
1408                }
1409            }
1410            StorageCommand::RunIngestion(ingestion) => {
1411                let RunIngestionCommand { id, description } = *ingestion;
1412
1413                // Remember the ingestion description to facilitate possible
1414                // reconciliation later.
1415                self.ingestions.insert(id, description.clone());
1416
1417                // Initialize shared frontier reporting.
1418                for id in description.collection_ids() {
1419                    self.reported_frontiers
1420                        .entry(id)
1421                        .or_insert_with(|| Antichain::from_elem(mz_repr::Timestamp::minimum()));
1422                }
1423
1424                // This needs to be done by one worker, which will broadcasts a
1425                // `CreateIngestionDataflow` command to all workers based on the response that
1426                // contains the resumption upper.
1427                //
1428                // Doing this separately on each worker could lead to differing resume_uppers
1429                // which might lead to all kinds of mayhem.
1430                //
1431                // n.b. the ingestion on each worker uses the description from worker 0––not the
1432                // ingestion in the local storage state. This is something we might have
1433                // interest in fixing in the future, e.g. materialize#19907
1434                if self.timely_worker_index == 0 {
1435                    self.async_worker
1436                        .update_ingestion_frontiers(id, description);
1437                }
1438            }
1439            StorageCommand::RunOneshotIngestion(oneshot) => {
1440                if self.timely_worker_index == 0 {
1441                    self.internal_cmd_tx
1442                        .send(InternalStorageCommand::RunOneshotIngestion {
1443                            ingestion_id: oneshot.ingestion_id,
1444                            collection_id: oneshot.collection_id,
1445                            collection_meta: oneshot.collection_meta,
1446                            request: oneshot.request,
1447                        });
1448                }
1449            }
1450            StorageCommand::CancelOneshotIngestion(id) => {
1451                self.drop_oneshot_ingestion(id);
1452            }
1453            StorageCommand::RunSink(export) => {
1454                // Remember the sink description to facilitate possible
1455                // reconciliation later.
1456                let prev = self.exports.insert(export.id, export.description.clone());
1457
1458                // New sink, add state.
1459                if prev.is_none() {
1460                    self.reported_frontiers.insert(
1461                        export.id,
1462                        Antichain::from_elem(mz_repr::Timestamp::minimum()),
1463                    );
1464                }
1465
1466                // This needs to be broadcast by one worker and go through the internal command
1467                // fabric, to ensure consistent ordering of dataflow rendering across all
1468                // workers.
1469                if self.timely_worker_index == 0 {
1470                    self.internal_cmd_tx
1471                        .send(InternalStorageCommand::RunSinkDataflow(
1472                            export.id,
1473                            export.description,
1474                        ));
1475                }
1476            }
1477            StorageCommand::AllowCompaction(id, frontier) => {
1478                soft_assert_or_log!(
1479                    self.exports.contains_key(&id) || self.reported_frontiers.contains_key(&id),
1480                    "AllowCompaction command for non-existent {id}"
1481                );
1482
1483                if frontier.is_empty() {
1484                    // Indicates that we may drop `id`, as there are no more valid times to read.
1485                    self.drop_collection(id);
1486                }
1487            }
1488        }
1489    }
1490
1491    /// Drop the identified storage collection from the storage state.
1492    fn drop_collection(&mut self, id: GlobalId) {
1493        fail_point!("crash_on_drop");
1494
1495        self.ingestions.remove(&id);
1496        self.exports.remove(&id);
1497
1498        let _ = self.latest_status_updates.remove(&id);
1499
1500        // This will stop reporting of frontiers.
1501        //
1502        // If this object still has its frontiers reported, we will notify the
1503        // client envd of the drop.
1504        if self.reported_frontiers.remove(&id).is_some() {
1505            // The only actions left are internal cleanup, so we can commit to
1506            // the client that these objects have been dropped.
1507            //
1508            // This must be done now rather than in response to `DropDataflow`,
1509            // otherwise we introduce the possibility of a timing issue where:
1510            // - We remove all tracking state from the storage state and send
1511            //   `DropDataflow` (i.e. this block).
1512            // - While waiting to process that command, we reconcile with a new
1513            //   envd. That envd has already committed to its catalog that this
1514            //   object no longer exists.
1515            // - We process the `DropDataflow` command, and identify that this
1516            //   object has been dropped.
1517            // - The next time `dropped_ids` is processed, we send a response
1518            //   that this ID has been dropped, but the upstream state has no
1519            //   record of that object having ever existed.
1520            self.dropped_ids.push(id);
1521        }
1522
1523        // Send through async worker for correct ordering with RunIngestion, and
1524        // dropping the dataflow is done on async worker response.
1525        if self.timely_worker_index == 0 {
1526            self.async_worker.drop_dataflow(id);
1527        }
1528    }
1529
1530    /// Drop the identified oneshot ingestion from the storage state.
1531    fn drop_oneshot_ingestion(&mut self, ingestion_id: uuid::Uuid) {
1532        let prev = self.oneshot_ingestions.remove(&ingestion_id);
1533        info!(%ingestion_id, existed = %prev.is_some(), "dropping oneshot ingestion");
1534    }
1535}