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}