Skip to main content

mz_compute/
command_channel.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// As of the Change Date specified in that file, in accordance with
7// the Business Source License, use of this software will be governed
8// by the Apache License, Version 2.0.
9
10//! A channel for sequencing commands between all workers of a Timely cluster.
11//!
12//! Compute uses a dataflow to distribute commands between workers. This is to ensure workers
13//! retain a consistent dataflow state across reconnects. If each worker would receive its commands
14//! directly from the controller, there wouldn't be any guarantee that after a reconnect all
15//! workers have seen the same sequence of commands. This is particularly problematic for
16//! `CreateDataflow` commands, since Timely requires that all workers render the same dataflows in
17//! the same order. So the controller instead sends commands only to worker 0, which then
18//! broadcasts them to other workers through the Timely fabric, taking care of the correct
19//! sequencing.
20//!
21//! Commands in the command channel are tagged with a nonce identifying the incarnation of the
22//! compute protocol the command belongs to, allowing workers to recognize client reconnects that
23//! require a reconciliation.
24//!
25//! The channel optionally also carries storage-internal commands, for
26//! clusters that host storage objects alongside compute objects. Both command kinds are sequenced
27//! through a single lane, so all workers observe one consistent interleaving and therefore
28//! construct all dataflows, compute and storage alike, in the same order. Unlike compute commands,
29//! storage-internal commands may be injected from any worker (e.g. by health operators triggering
30//! a suspend-and-restart), so the channel uses a two-hop structure copied from storage's command
31//! sequencer: producers tag commands with a per-producer index, worker 0 fixes one definitive
32//! order and assigns a global index, and receivers restore that order.
33
34use std::cell::RefCell;
35use std::collections::BTreeMap;
36use std::rc::Rc;
37use std::sync::mpsc::{self, TryRecvError};
38use std::sync::{Arc, Mutex};
39
40use itertools::Itertools;
41use mz_compute_client::protocol::command::ComputeCommand;
42use mz_compute_types::dataflows::{BuildDesc, DataflowDescription};
43use mz_ore::cast::CastFrom;
44use mz_storage::internal_control::InternalStorageCommand;
45use mz_timely_util::scope_label::ScopeExt;
46use serde::{Deserialize, Serialize};
47use timely::dataflow::channels::pact::Exchange;
48use timely::dataflow::operators::Operator;
49use timely::dataflow::operators::generic::source;
50use timely::scheduling::{Activator, SyncActivator};
51use timely::worker::Worker as TimelyWorker;
52use uuid::Uuid;
53
54#[cfg(test)]
55mod tests;
56
57/// A command in the unified command lane.
58#[derive(Clone, Debug, Serialize, Deserialize)]
59pub enum UnifiedCommand {
60    /// A compute command, tagged with the client nonce.
61    Compute(ComputeCommand, Uuid),
62    /// A storage-internal command.
63    Storage(InternalStorageCommand),
64}
65
66/// A sender pushing compute commands onto the command channel.
67pub struct Sender {
68    tx: mpsc::Sender<(ComputeCommand, Uuid)>,
69    activator: Arc<Mutex<Option<SyncActivator>>>,
70}
71
72impl Sender {
73    /// Broadcasts the given command to all workers.
74    pub fn send(&self, message: (ComputeCommand, Uuid)) {
75        if self.tx.send(message).is_err() {
76            unreachable!("command channel never shuts down");
77        }
78
79        self.activator
80            .lock()
81            .expect("poisoned")
82            .as_ref()
83            .map(|a| a.activate());
84    }
85}
86
87/// A receiver reading commands from the command channel.
88pub struct Receiver {
89    rx: mpsc::Receiver<UnifiedCommand>,
90}
91
92impl Receiver {
93    /// Returns the next available command, if any.
94    ///
95    /// This returns `None` when there are currently no commands but there might be commands again
96    /// in the future.
97    pub fn try_recv(&self) -> Option<UnifiedCommand> {
98        match self.rx.try_recv() {
99            Ok(msg) => Some(msg),
100            Err(TryRecvError::Empty) => None,
101            Err(TryRecvError::Disconnected) => {
102                unreachable!("command channel never shuts down");
103            }
104        }
105    }
106}
107
108/// Per-worker storage-side inputs to the command channel.
109///
110/// Created by the host before rendering the channel. The sending half of `rx` and the filled
111/// `activator_slot` together back the guest's `InternalCommandSender`.
112pub struct StorageLaneInput {
113    /// Receiver for storage-internal commands injected on this worker.
114    pub rx: mpsc::Receiver<InternalStorageCommand>,
115    /// Slot the channel fills with an activator for the source operator, so sends wake the
116    /// dataflow.
117    pub activator_slot: Rc<RefCell<Option<Activator>>>,
118}
119
120/// Render the command channel dataflow.
121pub fn render(
122    timely_worker: &mut TimelyWorker,
123    storage_input: Option<StorageLaneInput>,
124) -> (Sender, Receiver) {
125    let (input_tx, input_rx) = mpsc::channel();
126    let (output_tx, output_rx) = mpsc::channel();
127    let activator = Arc::new(Mutex::new(None));
128
129    timely_worker.dataflow_named::<(), _, _>("command_channel", {
130        let activator = Arc::clone(&activator);
131        move |scope| {
132            let scope = scope.with_label();
133
134            let peers = scope.peers();
135
136            // Create a stream of commands received from this worker's input queues.
137            //
138            // The output commands are tagged by worker ID and a per-producer command index,
139            // allowing the sequencer to restore their correct relative order.
140            let stream = source(scope, "command_channel::source", |cap, info| {
141                let sync_activator = scope.worker().sync_activator_for(info.address.to_vec());
142                *activator.lock().expect("poisoned") = Some(sync_activator);
143
144                if let Some(input) = &storage_input {
145                    let act = scope.activator_for(info.address);
146                    *input.activator_slot.borrow_mut() = Some(act);
147                }
148
149                let worker_id = scope.index();
150                let mut cmd_index = 0_u64;
151                let mut capability = Some(cap);
152
153                move |output| {
154                    let Some(cap) = &capability else {
155                        return;
156                    };
157
158                    let mut session = output.session(cap);
159
160                    let mut compute_disconnected = false;
161                    loop {
162                        match input_rx.try_recv() {
163                            Ok((cmd, nonce)) if worker_id == 0 => {
164                                session.give((
165                                    worker_id,
166                                    cmd_index,
167                                    UnifiedCommand::Compute(cmd, nonce),
168                                ));
169                                cmd_index += 1;
170                            }
171                            Ok((cmd, _nonce)) => {
172                                // Non-leader workers only receive `UpdateConfiguration` commands
173                                // from the controller and must drop them to not sequence
174                                // duplicates.
175                                assert!(matches!(cmd, ComputeCommand::UpdateConfiguration(_)));
176                            }
177                            Err(TryRecvError::Empty) => break,
178                            Err(TryRecvError::Disconnected) => {
179                                compute_disconnected = true;
180                                break;
181                            }
182                        }
183                    }
184
185                    let mut storage_disconnected = true;
186                    if let Some(input) = &storage_input {
187                        storage_disconnected = false;
188                        loop {
189                            match input.rx.try_recv() {
190                                Ok(cmd) => {
191                                    session.give((
192                                        worker_id,
193                                        cmd_index,
194                                        UnifiedCommand::Storage(cmd),
195                                    ));
196                                    cmd_index += 1;
197                                }
198                                Err(TryRecvError::Empty) => break,
199                                Err(TryRecvError::Disconnected) => {
200                                    storage_disconnected = true;
201                                    break;
202                                }
203                            }
204                        }
205                    }
206
207                    drop(session);
208
209                    // Once every sender is gone no further commands can arrive, so release the
210                    // capability to let the dataflow shut down.
211                    if compute_disconnected && storage_disconnected {
212                        capability = None;
213                    }
214                }
215            });
216
217            // Sequence all commands through a single worker to establish a unique order.
218            //
219            // The output commands are tagged with a global command index and a target worker,
220            // allowing downstream operators to ensure their correct relative order.
221            let stream = stream.unary_frontier(
222                Exchange::new(|_| 0),
223                "command_channel::sequencer",
224                |cap, _info| {
225                    let mut global_index = 0_u64;
226                    let mut capability = Some(cap);
227
228                    // For each producer, keep an ordered list of pending commands, as well as the
229                    // index of the next command.
230                    let mut pending: Vec<(BTreeMap<u64, UnifiedCommand>, u64)> =
231                        vec![(BTreeMap::new(), 0); peers];
232
233                    move |(input, frontier), output| {
234                        let Some(cap) = capability.clone() else {
235                            return;
236                        };
237
238                        input.for_each(|_time, data| {
239                            for (producer, index, cmd) in data.drain(..) {
240                                // An index below the next expected one was already sequenced, so
241                                // it would sit in the pending map forever, wedging this producer.
242                                let next_idx = pending[producer].1;
243                                mz_ore::soft_assert_or_log!(
244                                    index >= next_idx,
245                                    "command index {index} from producer {producer} \
246                                     stepped back behind {next_idx}"
247                                );
248                                let duplicate = pending[producer].0.insert(index, cmd);
249                                mz_ore::soft_assert_or_log!(
250                                    duplicate.is_none(),
251                                    "duplicate command index {index} from producer {producer}"
252                                );
253                            }
254                        });
255
256                        let mut session = output.session(&cap);
257                        for (commands, next_idx) in &mut pending {
258                            while commands
259                                .first_key_value()
260                                .is_some_and(|(i, _)| i == next_idx)
261                            {
262                                let (_, cmd) = commands.pop_first().unwrap();
263                                for (target, part) in split_command(cmd, peers) {
264                                    session.give((target, global_index, part));
265                                }
266
267                                *next_idx += 1;
268                                global_index += 1;
269                            }
270                        }
271
272                        drop(session);
273
274                        if frontier.is_empty() {
275                            // Drop our capability to shut down.
276                            capability = None;
277                        }
278                    }
279                },
280            );
281
282            // Sink the stream back into `output_tx`, restoring the global order.
283            stream.sink(
284                Exchange::new(|(target, _, _)| u64::cast_from(*target)),
285                "command_channel::sink",
286                {
287                    // Pending commands by global index, and the index of the next command.
288                    let mut pending = BTreeMap::new();
289                    let mut next_idx = 0_u64;
290
291                    move |(input, _frontier)| {
292                        input.for_each(|_time, data| {
293                            for (_target, index, cmd) in data.drain(..) {
294                                // An index below the next expected one was already delivered, so
295                                // it would sit in the pending map forever, wedging the channel.
296                                mz_ore::soft_assert_or_log!(
297                                    index >= next_idx,
298                                    "global command index {index} stepped back behind {next_idx}"
299                                );
300                                let duplicate = pending.insert(index, cmd);
301                                mz_ore::soft_assert_or_log!(
302                                    duplicate.is_none(),
303                                    "duplicate global command index {index}"
304                                );
305                            }
306                        });
307
308                        while pending
309                            .first_key_value()
310                            .is_some_and(|(i, _)| *i == next_idx)
311                        {
312                            let (_, cmd) = pending.pop_first().unwrap();
313                            let _ = output_tx.send(cmd);
314                            next_idx += 1;
315                        }
316                    }
317                },
318            );
319        }
320    });
321
322    let tx = Sender {
323        tx: input_tx,
324        activator,
325    };
326    let rx = Receiver { rx: output_rx };
327
328    (tx, rx)
329}
330
331/// Split the given command into one part per target worker.
332///
333/// Compute `CreateDataflow` commands are partitioned among the workers. Every other command is
334/// replicated to all workers.
335fn split_command(
336    command: UnifiedCommand,
337    parts: usize,
338) -> impl Iterator<Item = (usize, UnifiedCommand)> {
339    use itertools::Either;
340
341    let commands = match command {
342        UnifiedCommand::Compute(ComputeCommand::CreateDataflow(dataflow), nonce) => {
343            let dataflow = *dataflow;
344
345            // A list of descriptions of objects for each part to build.
346            let mut builds_parts = vec![Vec::new(); parts];
347            // Partition each build description among `parts`.
348            for build_desc in dataflow.objects_to_build {
349                let build_part = build_desc.plan.partition_among(parts);
350                for (plan, objects_to_build) in
351                    build_part.into_iter().zip_eq(builds_parts.iter_mut())
352                {
353                    objects_to_build.push(BuildDesc {
354                        id: build_desc.id,
355                        plan,
356                    });
357                }
358            }
359
360            // Each list of build descriptions results in a dataflow description.
361            let commands = builds_parts
362                .into_iter()
363                .map(move |objects_to_build| DataflowDescription {
364                    source_imports: dataflow.source_imports.clone(),
365                    index_imports: dataflow.index_imports.clone(),
366                    objects_to_build,
367                    index_exports: dataflow.index_exports.clone(),
368                    sink_exports: dataflow.sink_exports.clone(),
369                    as_of: dataflow.as_of.clone(),
370                    until: dataflow.until.clone(),
371                    debug_name: dataflow.debug_name.clone(),
372                    initial_storage_as_of: dataflow.initial_storage_as_of.clone(),
373                    refresh_schedule: dataflow.refresh_schedule.clone(),
374                    time_dependence: dataflow.time_dependence.clone(),
375                })
376                .map(Box::new)
377                .map(move |dataflow| {
378                    UnifiedCommand::Compute(ComputeCommand::CreateDataflow(dataflow), nonce)
379                });
380            Either::Left(commands)
381        }
382        command => {
383            let commands = std::iter::repeat_n(command, parts);
384            Either::Right(commands)
385        }
386    };
387
388    commands.into_iter().enumerate()
389}