Skip to main content

mz_storage/
internal_control.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//! Types for cluster-internal control messages that can be broadcast to all
7//! workers from individual operators/workers.
8
9use std::cell::RefCell;
10use std::collections::{BTreeMap, BTreeSet};
11use std::rc::Rc;
12use std::sync::mpsc;
13
14use mz_repr::{GlobalId, Row};
15use mz_rocksdb::config::SharedWriteBufferManager;
16use mz_storage_types::controller::CollectionMetadata;
17use mz_storage_types::oneshot_sources::OneshotIngestionRequest;
18use mz_storage_types::parameters::StorageParameters;
19use mz_storage_types::sinks::StorageSinkDesc;
20use mz_storage_types::sources::IngestionDescription;
21use mz_timely_util::scope_label::ScopeExt;
22use serde::{Deserialize, Serialize};
23use timely::dataflow::channels::pact::{Exchange, Pipeline};
24use timely::dataflow::operators::Operator;
25use timely::dataflow::operators::generic::source;
26use timely::dataflow::operators::vec::Broadcast;
27use timely::progress::Antichain;
28use timely::scheduling::Activator;
29use timely::worker::Worker as TimelyWorker;
30
31use crate::statistics::{SinkStatisticsRecord, SourceStatisticsRecord};
32
33/// _Dynamic_ storage instance configuration parameters that are used during dataflow rendering.
34/// Changes to these parameters are applied to `StorageWorker`s in a consistent order
35/// with source and sink creation.
36#[derive(Debug)]
37pub struct DataflowParameters {
38    /// Configuration/tuning for RocksDB. This also contains
39    /// some shared objects, which is why its separate.
40    pub upsert_rocksdb_tuning_config: mz_rocksdb::RocksDBConfig,
41}
42
43impl DataflowParameters {
44    /// Creates a new instance of `DataflowParameters` with given shared rocksdb write buffer manager
45    /// and the cluster memory limit
46    pub fn new(
47        shared_rocksdb_write_buffer_manager: SharedWriteBufferManager,
48        cluster_memory_limit: Option<usize>,
49    ) -> Self {
50        Self {
51            upsert_rocksdb_tuning_config: mz_rocksdb::RocksDBConfig::new(
52                shared_rocksdb_write_buffer_manager,
53                cluster_memory_limit,
54            ),
55        }
56    }
57    /// Update the `DataflowParameters` with new configuration.
58    pub fn update(&mut self, storage_parameters: StorageParameters) {
59        self.upsert_rocksdb_tuning_config
60            .apply(storage_parameters.upsert_rocksdb_tuning_config.clone());
61    }
62}
63
64/// Internal commands that can be sent by individual operators/workers that will
65/// be broadcast to all workers. The worker main loop will receive those and act
66/// on them.
67#[derive(Clone, Debug, Serialize, Deserialize, PartialEq)]
68pub enum InternalStorageCommand {
69    /// Suspend and restart the dataflow identified by the `GlobalId`.
70    SuspendAndRestart {
71        /// The id of the dataflow that should be restarted.
72        id: GlobalId,
73        /// The reason for the restart request.
74        reason: String,
75    },
76    /// Render an ingestion dataflow at the given resumption frontier.
77    CreateIngestionDataflow {
78        /// ID of the ingestion/sourve.
79        id: GlobalId,
80        /// The description of the ingestion/source.
81        ingestion_description: IngestionDescription<CollectionMetadata>,
82        /// The frontier beyond which ingested updates should be uncompacted. Inputs to the
83        /// ingestion are guaranteed to be readable at this frontier.
84        as_of: Antichain<mz_repr::Timestamp>,
85        /// A frontier in the Materialize time domain with the property that all updates not beyond
86        /// it have already been durably ingested.
87        resume_uppers: BTreeMap<GlobalId, Antichain<mz_repr::Timestamp>>,
88        /// A frontier in the source time domain with the property that all updates not beyond it
89        /// have already been durably ingested.
90        source_resume_uppers: BTreeMap<GlobalId, Vec<Row>>,
91    },
92    /// Render a oneshot ingestion dataflow that fetches data from an external system and stages
93    /// batches in Persist, that can later be appended to the shard.
94    RunOneshotIngestion {
95        /// ID of the running dataflow that is doing the ingestion.
96        ingestion_id: uuid::Uuid,
97        /// ID of the collection we'll create batches for.
98        collection_id: GlobalId,
99        /// Metadata of the collection we'll create batches for.
100        collection_meta: CollectionMetadata,
101        /// Description of the oneshot ingestion.
102        request: OneshotIngestionRequest,
103    },
104    /// Render a sink dataflow.
105    RunSinkDataflow(
106        GlobalId,
107        StorageSinkDesc<CollectionMetadata, mz_repr::Timestamp>,
108    ),
109    /// Drop all state and operators for a dataflow. This is a vec because some
110    /// dataflows have their state spread over multiple IDs (i.e. sources that
111    /// spawn subsources); this means that actions taken in response to this
112    /// command should be permissive about missing state.
113    DropDataflow(Vec<GlobalId>),
114
115    /// Update the configuration for rendering dataflows.
116    UpdateConfiguration {
117        /// The new configuration parameters.
118        storage_parameters: StorageParameters,
119    },
120    /// For moving statistics updates to worker 0.
121    StatisticsUpdate {
122        /// Local statistics, with their epochs.
123        sources: Vec<(usize, SourceStatisticsRecord)>,
124        /// Local statistics, with their epochs.
125        sinks: Vec<(usize, SinkStatisticsRecord)>,
126    },
127}
128
129/// A sender broadcasting [`InternalStorageCommand`]s to all workers.
130#[derive(Clone)]
131pub struct InternalCommandSender {
132    tx: mpsc::Sender<InternalStorageCommand>,
133    activator: Rc<RefCell<Option<Activator>>>,
134}
135
136impl InternalCommandSender {
137    /// Creates a sender from externally provided parts, for hosts that
138    /// route internal commands through their own sequencing channel instead of
139    /// `setup_command_sequencer`.
140    pub fn from_parts(
141        tx: mpsc::Sender<InternalStorageCommand>,
142        activator: Rc<RefCell<Option<Activator>>>,
143    ) -> Self {
144        Self { tx, activator }
145    }
146
147    /// Broadcasts the given command to all workers.
148    pub fn send(&self, cmd: InternalStorageCommand) {
149        if self.tx.send(cmd).is_err() {
150            panic!("internal command channel disconnected");
151        }
152
153        self.activator.borrow().as_ref().map(|a| a.activate());
154    }
155}
156
157/// A receiver for [`InternalStorageCommand`]s broadcasted by workers.
158pub struct InternalCommandReceiver {
159    rx: mpsc::Receiver<InternalStorageCommand>,
160}
161
162impl InternalCommandReceiver {
163    /// Returns the next available command, if any.
164    ///
165    /// This returns `None` when there are currently no commands but there might be commands again
166    /// in the future.
167    pub fn try_recv(&self) -> Option<InternalStorageCommand> {
168        match self.rx.try_recv() {
169            Ok(cmd) => Some(cmd),
170            Err(mpsc::TryRecvError::Empty) => None,
171            Err(mpsc::TryRecvError::Disconnected) => {
172                panic!("internal command channel disconnected")
173            }
174        }
175    }
176}
177
178pub(crate) fn setup_command_sequencer<'w>(
179    timely_worker: &'w mut TimelyWorker,
180) -> (InternalCommandSender, InternalCommandReceiver) {
181    let (input_tx, input_rx) = mpsc::channel();
182    let (output_tx, output_rx) = mpsc::channel();
183    let activator = Rc::new(RefCell::new(None));
184
185    timely_worker.dataflow_named::<(), _, _>("command_sequencer", {
186        let activator = Rc::clone(&activator);
187        move |scope| {
188            let scope = scope.with_label();
189            // Create a stream of commands received from `input_rx`.
190            //
191            // The output commands are tagged by worker ID and command index, allowing downstream
192            // operators to ensure their correct relative order.
193            let stream = source(scope, "command_sequencer::source", |cap, info| {
194                *activator.borrow_mut() = Some(scope.activator_for(info.address));
195
196                let worker_id = scope.index();
197                let mut cmd_index = 0;
198                let mut capability = Some(cap);
199
200                move |output| {
201                    let Some(cap) = capability.clone() else {
202                        return;
203                    };
204
205                    let mut session = output.session(&cap);
206                    loop {
207                        match input_rx.try_recv() {
208                            Ok(command) => {
209                                let cmd = IndexedCommand {
210                                    index: cmd_index,
211                                    command,
212                                };
213                                session.give((worker_id, cmd));
214                                cmd_index += 1;
215                            }
216                            Err(mpsc::TryRecvError::Empty) => break,
217                            Err(mpsc::TryRecvError::Disconnected) => {
218                                // Drop our capability to shut down.
219                                capability = None;
220                                break;
221                            }
222                        }
223                    }
224                }
225            });
226
227            // Sequence all commands through a single worker to establish a unique order.
228            //
229            // The output commands are tagged by a command index, allowing downstream operators to
230            // ensure their correct relative order.
231            let stream = stream.unary_frontier(
232                Exchange::new(|_| 0),
233                "command_sequencer::sequencer",
234                |cap, _info| {
235                    let mut cmd_index = 0;
236                    let mut capability = Some(cap);
237
238                    // For each worker, keep an ordered list of pending commands, as well as the
239                    // current index of the next command.
240                    let mut pending_commands = vec![(BTreeSet::new(), 0); scope.peers()];
241
242                    move |(input, frontier), output| {
243                        let Some(cap) = capability.clone() else {
244                            return;
245                        };
246
247                        input.for_each(|_time, data| {
248                            for (worker_id, cmd) in data.drain(..) {
249                                pending_commands[worker_id].0.insert(cmd);
250                            }
251                        });
252
253                        let mut session = output.session(&cap);
254                        for (commands, next_idx) in &mut pending_commands {
255                            while commands.first().is_some_and(|c| c.index == *next_idx) {
256                                let mut cmd = commands.pop_first().unwrap();
257                                cmd.index = cmd_index;
258                                session.give(cmd);
259
260                                *next_idx += 1;
261                                cmd_index += 1;
262                            }
263                        }
264
265                        let _ = session;
266
267                        if frontier.is_empty() {
268                            // Drop our capability to shut down.
269                            capability = None;
270                        }
271                    }
272                },
273            );
274
275            // Broadcast the ordered commands to all workers.
276            let stream = stream.broadcast();
277
278            // Sink the stream back into `output_tx`.
279            stream.sink(Pipeline, "command_sequencer::sink", {
280                // Keep an ordered list of pending commands, as well as the current index of the
281                // next command.
282                let mut pending_commands = BTreeSet::new();
283                let mut next_idx = 0;
284
285                move |(input, _frontier)| {
286                    input.for_each(|_time, data| {
287                        pending_commands.extend(data.drain(..));
288                    });
289
290                    while pending_commands
291                        .first()
292                        .is_some_and(|c| c.index == next_idx)
293                    {
294                        let cmd = pending_commands.pop_first().unwrap();
295                        let _ = output_tx.send(cmd.command);
296                        next_idx += 1;
297                    }
298                }
299            });
300        }
301    });
302
303    let tx = InternalCommandSender {
304        tx: input_tx,
305        activator,
306    };
307    let rx = InternalCommandReceiver { rx: output_rx };
308
309    (tx, rx)
310}
311
312// An [`InternalStorageCommand`] tagged with an index.
313//
314// This is a `(u64, InternalStorageCommand)` in spirit, but implements `Ord` (which
315// `InternalStorageCommand` doesn't) by looking only at the index.
316#[derive(Clone, Debug, Serialize, Deserialize)]
317struct IndexedCommand {
318    index: u64,
319    command: InternalStorageCommand,
320}
321
322impl PartialEq for IndexedCommand {
323    fn eq(&self, other: &Self) -> bool {
324        self.cmp(other).is_eq()
325    }
326}
327
328impl Eq for IndexedCommand {}
329
330impl PartialOrd for IndexedCommand {
331    fn partial_cmp(&self, other: &Self) -> Option<std::cmp::Ordering> {
332        Some(self.cmp(other))
333    }
334}
335
336impl Ord for IndexedCommand {
337    fn cmp(&self, other: &Self) -> std::cmp::Ordering {
338        self.index.cmp(&other.index)
339    }
340}