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}