Skip to main content

mz_cluster/
client.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//! An interactive cluster server.
11
12use std::fmt;
13use std::sync::atomic::{AtomicU64, Ordering};
14use std::sync::{Arc, Mutex, OnceLock};
15use std::thread::Thread;
16
17use anyhow::{Error, anyhow};
18use async_trait::async_trait;
19use differential_dataflow::trace::ExertionLogic;
20use futures::future;
21use mz_cluster_client::client::{TimelyConfig, TryIntoProtocolNonce};
22use mz_ore::metric;
23use mz_ore::metrics::{ComputedUIntGauge, MetricsRegistry};
24use mz_service::client::{GenericClient, Partitionable, Partitioned};
25use mz_service::local::LocalClient;
26use timely::WorkerConfig;
27use timely::communication::allocator::zero_copy::bytes_slab::BytesRefill;
28use timely::communication::initialize::WorkerGuards;
29use timely::execute::execute_from;
30use timely::worker::Worker as TimelyWorker;
31use tokio::runtime::Handle;
32use tokio::sync::mpsc;
33use tracing::{info, info_span};
34use uuid::Uuid;
35
36use crate::communication::initialize_networking;
37
38type PartitionedClient<C, R> = Partitioned<LocalClient<C, R>, C, R>;
39
40/// A client managing access to the local portion of a Timely cluster
41pub struct ClusterClient<C>
42where
43    C: ClusterSpec,
44    (C::Command, C::Response): Partitionable<C::Command, C::Response>,
45{
46    /// The actual client to talk to the cluster
47    inner: Option<PartitionedClient<C::Command, C::Response>>,
48    /// The running timely instance
49    timely_container: Arc<Mutex<TimelyContainer<C>>>,
50}
51
52/// Metadata about timely workers in this process.
53pub struct TimelyContainer<C: ClusterSpec> {
54    /// Channels over which to send endpoints for wiring up a new Client
55    client_txs: Vec<
56        mpsc::UnboundedSender<(
57            Uuid,
58            mpsc::UnboundedReceiver<C::Command>,
59            mpsc::UnboundedSender<C::Response>,
60        )>,
61    >,
62    /// Thread guards that keep worker threads alive
63    worker_guards: WorkerGuards<()>,
64}
65
66impl<C: ClusterSpec> TimelyContainer<C> {
67    /// The threads of the Timely workers in this process.
68    pub fn worker_threads(&self) -> Vec<Thread> {
69        self.worker_guards
70            .guards()
71            .iter()
72            .map(|h| h.thread().clone())
73            .collect()
74    }
75}
76
77/// A client to a secondary ("guest") command stream served by workers of
78/// an existing Timely cluster. Like [`ClusterClient`], but the per-worker client channels are
79/// provided externally instead of coming from a [`TimelyContainer`] built for this command type.
80pub struct GuestClusterClient<Cmd, Resp>
81where
82    (Cmd, Resp): Partitionable<Cmd, Resp>,
83{
84    /// Per-worker channels over which to send endpoints for wiring up a new client.
85    client_txs: Arc<
86        Vec<
87            mpsc::UnboundedSender<(
88                Uuid,
89                mpsc::UnboundedReceiver<Cmd>,
90                mpsc::UnboundedSender<Resp>,
91            )>,
92        >,
93    >,
94    /// The worker threads, for unparking on send.
95    worker_threads: Vec<Thread>,
96    /// The actual client to talk to the cluster.
97    inner: Option<PartitionedClient<Cmd, Resp>>,
98}
99
100impl<Cmd, Resp> GuestClusterClient<Cmd, Resp>
101where
102    Cmd: fmt::Debug + Send + TryIntoProtocolNonce,
103    Resp: fmt::Debug + Send,
104    (Cmd, Resp): Partitionable<Cmd, Resp>,
105{
106    /// Create a new `GuestClusterClient`.
107    pub fn new(
108        client_txs: Arc<
109            Vec<
110                mpsc::UnboundedSender<(
111                    Uuid,
112                    mpsc::UnboundedReceiver<Cmd>,
113                    mpsc::UnboundedSender<Resp>,
114                )>,
115            >,
116        >,
117        worker_threads: Vec<Thread>,
118    ) -> Self {
119        Self {
120            client_txs,
121            worker_threads,
122            inner: None,
123        }
124    }
125
126    fn connect(&mut self, nonce: Uuid) {
127        let mut command_txs = Vec::new();
128        let mut response_rxs = Vec::new();
129        for client_tx in self.client_txs.iter() {
130            let (cmd_tx, cmd_rx) = mpsc::unbounded_channel();
131            let (resp_tx, resp_rx) = mpsc::unbounded_channel();
132
133            client_tx
134                .send((nonce, cmd_rx, resp_tx))
135                .expect("worker not dropped");
136
137            command_txs.push(cmd_tx);
138            response_rxs.push(resp_rx);
139        }
140
141        self.inner = Some(LocalClient::new_partitioned(
142            response_rxs,
143            command_txs,
144            self.worker_threads.clone(),
145        ));
146    }
147}
148
149impl<Cmd, Resp> fmt::Debug for GuestClusterClient<Cmd, Resp>
150where
151    (Cmd, Resp): Partitionable<Cmd, Resp>,
152{
153    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
154        f.debug_struct("GuestClusterClient").finish_non_exhaustive()
155    }
156}
157
158#[async_trait]
159impl<Cmd, Resp> GenericClient<Cmd, Resp> for GuestClusterClient<Cmd, Resp>
160where
161    Cmd: fmt::Debug + Send + TryIntoProtocolNonce,
162    Resp: fmt::Debug + Send,
163    (Cmd, Resp): Partitionable<Cmd, Resp>,
164{
165    async fn send(&mut self, cmd: Cmd) -> Result<(), Error> {
166        match cmd.try_into_protocol_nonce() {
167            Ok(nonce) => {
168                self.connect(nonce);
169                Ok(())
170            }
171            Err(cmd) => self.inner.as_mut().expect("initialized").send(cmd).await,
172        }
173    }
174
175    /// # Cancel safety
176    ///
177    /// This method is cancel safe, see [`ClusterClient::recv`].
178    async fn recv(&mut self) -> Result<Option<Resp>, Error> {
179        if let Some(client) = self.inner.as_mut() {
180            client.recv().await
181        } else {
182            future::pending().await
183        }
184    }
185}
186
187impl<C> ClusterClient<C>
188where
189    C: ClusterSpec,
190    (C::Command, C::Response): Partitionable<C::Command, C::Response>,
191{
192    /// Create a new `ClusterClient`.
193    pub fn new(timely_container: Arc<Mutex<TimelyContainer<C>>>) -> Self {
194        Self {
195            timely_container,
196            inner: None,
197        }
198    }
199
200    /// Connect to the Timely cluster with the given client nonce.
201    fn connect(&mut self, nonce: Uuid) -> Result<(), Error> {
202        let timely = self.timely_container.lock().expect("poisoned");
203
204        let mut command_txs = Vec::new();
205        let mut response_rxs = Vec::new();
206        for client_tx in &timely.client_txs {
207            let (cmd_tx, cmd_rx) = mpsc::unbounded_channel();
208            let (resp_tx, resp_rx) = mpsc::unbounded_channel();
209
210            client_tx
211                .send((nonce, cmd_rx, resp_tx))
212                .expect("worker not dropped");
213
214            command_txs.push(cmd_tx);
215            response_rxs.push(resp_rx);
216        }
217
218        self.inner = Some(LocalClient::new_partitioned(
219            response_rxs,
220            command_txs,
221            timely.worker_threads(),
222        ));
223        Ok(())
224    }
225}
226
227impl<C> fmt::Debug for ClusterClient<C>
228where
229    C: ClusterSpec,
230    (C::Command, C::Response): Partitionable<C::Command, C::Response>,
231{
232    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
233        f.debug_struct("ClusterClient")
234            .field("inner", &self.inner)
235            .finish_non_exhaustive()
236    }
237}
238
239#[async_trait]
240impl<C> GenericClient<C::Command, C::Response> for ClusterClient<C>
241where
242    C: ClusterSpec,
243    (C::Command, C::Response): Partitionable<C::Command, C::Response>,
244{
245    async fn send(&mut self, cmd: C::Command) -> Result<(), Error> {
246        // Changing this debug statement requires changing the replica-isolation test
247        tracing::debug!("ClusterClient send={:?}", &cmd);
248
249        match cmd.try_into_protocol_nonce() {
250            Ok(nonce) => self.connect(nonce),
251            Err(cmd) => self.inner.as_mut().expect("initialized").send(cmd).await,
252        }
253    }
254
255    /// # Cancel safety
256    ///
257    /// This method is cancel safe. If `recv` is used as the event in a [`tokio::select!`]
258    /// statement and some other branch completes first, it is guaranteed that no messages were
259    /// received by this client.
260    async fn recv(&mut self) -> Result<Option<C::Response>, Error> {
261        if let Some(client) = self.inner.as_mut() {
262            // `Partitioned::recv` is documented as cancel safe.
263            client.recv().await
264        } else {
265            future::pending().await
266        }
267    }
268}
269
270/// Specification for a Timely cluster to which a [`ClusterClient`] connects.
271///
272/// This trait is used to make the [`ClusterClient`] generic over the compute and storage cluster
273/// implementations.
274#[async_trait]
275pub trait ClusterSpec: Clone + Send + Sync + 'static {
276    /// The cluster command type.
277    type Command: fmt::Debug + Send + TryIntoProtocolNonce;
278    /// The cluster response type.
279    type Response: fmt::Debug + Send;
280
281    /// The name of this cluster ("compute" or "storage").
282    const NAME: &str;
283
284    /// Run the given Timely worker.
285    fn run_worker(
286        &self,
287        timely_worker: &mut TimelyWorker,
288        client_rx: mpsc::UnboundedReceiver<(
289            Uuid,
290            mpsc::UnboundedReceiver<Self::Command>,
291            mpsc::UnboundedSender<Self::Response>,
292        )>,
293    );
294
295    /// Build a Timely cluster using the given config.
296    async fn build_cluster(
297        &self,
298        config: TimelyConfig,
299        tokio_executor: Handle,
300    ) -> Result<TimelyContainer<Self>, Error> {
301        info!("Building timely container with config {config:?}");
302        let (client_txs, client_rxs): (Vec<_>, Vec<_>) = (0..config.workers)
303            .map(|_| mpsc::unbounded_channel())
304            .unzip();
305        let client_rxs: Mutex<Vec<_>> = Mutex::new(client_rxs.into_iter().map(Some).collect());
306
307        let refill = if config.enable_zero_copy_lgalloc {
308            BytesRefill {
309                logic: Arc::new(|size| Box::new(alloc::lgalloc_refill(size))),
310                limit: config.zero_copy_limit,
311            }
312        } else {
313            BytesRefill {
314                logic: Arc::new(|size| Box::new(vec![0; size])),
315                limit: config.zero_copy_limit,
316            }
317        };
318
319        let (builders, other) = initialize_networking(
320            config.workers,
321            config.process,
322            config.addresses.clone(),
323            refill,
324            config.enable_zero_copy,
325        )
326        .await?;
327
328        let mut worker_config = WorkerConfig::default();
329
330        // We set a custom exertion logic for proportionality > 0. A proportionality value of 0
331        // means that no arrangement merge effort is exerted and merging occurs only in response to
332        // updates.
333        if config.arrangement_exert_proportionality > 0 {
334            let merge_effort = Some(1000);
335
336            // ExertionLogic defines a function to determine if a spine is sufficiently tidied.
337            // Its arguments are an iterator over the index of a layer, the count of batches in the
338            // layer and the length of batches at the layer. The iterator enumerates layers from the
339            // largest to the smallest layer.
340
341            let arc: ExertionLogic = Arc::new(move |layers| {
342                EXERT_POLICY_CALLS.fetch_add(1, Ordering::Relaxed);
343                let mut prop = config.arrangement_exert_proportionality;
344
345                // Layers are ordered from largest to smallest.
346                // Skip to the largest occupied layer.
347                let layers = layers
348                    .iter()
349                    .copied()
350                    .skip_while(|(_idx, count, _len)| *count == 0);
351
352                let mut first = true;
353                for (_idx, count, len) in layers {
354                    if count > 1 {
355                        // Found an in-progress merge that we should continue.
356                        EXERT_POLICY_MERGE_GRANTS.fetch_add(1, Ordering::Relaxed);
357                        return merge_effort;
358                    }
359
360                    if !first && prop > 0 && len > 0 {
361                        // Found a non-empty batch within `arrangement_exert_proportionality` of
362                        // the largest one.
363                        EXERT_POLICY_CONSOLIDATION_GRANTS.fetch_add(1, Ordering::Relaxed);
364                        return merge_effort;
365                    }
366
367                    first = false;
368                    prop /= 2;
369                }
370
371                None
372            });
373            worker_config.set::<ExertionLogic>("differential/default_exert_logic".to_string(), arc);
374        }
375
376        let spec = self.clone();
377        let worker_guards = execute_from(builders, other, worker_config, move |timely_worker| {
378            let worker_idx = timely_worker.index();
379
380            // Per worker tracing span, lets us identify Timely clusters and workers in the logs.
381            let span = info_span!("timely", name = Self::NAME, worker_id = worker_idx);
382            let _span_guard = span.enter();
383
384            // Every Timely instance in this process names its threads `timely:work-N`, restarting
385            // the index at 0, so storage and compute worker threads collide under the same OS
386            // thread name. Rename to disambiguate them for profilers and `top -H`.
387            mz_ore::process::set_current_thread_name(&format!("{}:{worker_idx}", Self::NAME));
388
389            let _tokio_guard = tokio_executor.enter();
390            let client_rx = client_rxs.lock().unwrap()[worker_idx % config.workers]
391                .take()
392                .unwrap();
393            spec.run_worker(timely_worker, client_rx);
394        })
395        .map_err(|e| anyhow!(e))?;
396
397        Ok(TimelyContainer {
398            client_txs,
399            worker_guards,
400        })
401    }
402}
403
404// Evaluations of the exertion policy installed above, process-wide across
405// every timely runtime and trace implementation. A spine evaluates the policy
406// both when exerted and after each insert, so evaluations and grants count
407// decisions, not applications of effort. With
408// `arrangement_exert_proportionality` at zero no policy is installed and the
409// counts stay zero.
410static EXERT_POLICY_CALLS: AtomicU64 = AtomicU64::new(0);
411static EXERT_POLICY_MERGE_GRANTS: AtomicU64 = AtomicU64::new(0);
412static EXERT_POLICY_CONSOLIDATION_GRANTS: AtomicU64 = AtomicU64::new(0);
413
414/// Register gauges for the arrangement exertion policy's decisions.
415///
416/// The counters are process-wide, so only the first call registers them, and
417/// they cover every timely runtime in the process, not only the caller's.
418pub fn register_exert_policy_metrics(registry: &MetricsRegistry) {
419    static REGISTERED: OnceLock<()> = OnceLock::new();
420    REGISTERED.get_or_init(|| {
421        let _: ComputedUIntGauge = registry.register_computed_gauge(
422            metric!(name: "mz_arrangement_exert_policy_calls_total", help: "Arrangement exertion policy evaluations."),
423            || EXERT_POLICY_CALLS.load(Ordering::Relaxed),
424        );
425        for (reason, grants) in [
426            ("active_merge", &EXERT_POLICY_MERGE_GRANTS),
427            ("consolidation", &EXERT_POLICY_CONSOLIDATION_GRANTS),
428        ] {
429            let _: ComputedUIntGauge = registry.register_computed_gauge(
430                metric!(name: "mz_arrangement_exert_policy_grants_total", help: "Arrangement exertion policy evaluations that returned effort, by reason.", const_labels: {"reason" => reason}),
431                move || grants.load(Ordering::Relaxed),
432            );
433        }
434    });
435}
436
437mod alloc {
438    /// A Timely communication refill function that uses lgalloc.
439    ///
440    /// Returns a handle to lgalloc'ed memory if lgalloc can handle
441    /// the request, otherwise we fall back to a heap allocation. In either case, the handle must
442    /// be dropped to free the memory.
443    pub(crate) fn lgalloc_refill(size: usize) -> LgallocHandle {
444        match lgalloc::allocate::<u8>(size) {
445            Ok((pointer, capacity, handle)) => {
446                let handle = Some(handle);
447                LgallocHandle {
448                    handle,
449                    pointer,
450                    capacity,
451                }
452            }
453            Err(_) => {
454                // Allocate memory
455                let mut alloc = vec![0_u8; size];
456                // Ensure that the length matches the capacity.
457                alloc.shrink_to_fit();
458                // Get a pointer to the allocated memory.
459                let pointer = std::ptr::NonNull::new(alloc.as_mut_ptr()).unwrap();
460                // Forget the vector to avoid dropping it. We'll free the memory in `drop`.
461                std::mem::forget(alloc);
462                LgallocHandle {
463                    handle: None,
464                    pointer,
465                    capacity: size,
466                }
467            }
468        }
469    }
470
471    /// A handle to memory allocated by lgalloc. This can either be memory allocated by lgalloc or
472    /// memory allocated by Vec. If the handle is set, it's lgalloc memory. If the handle is None,
473    /// it's a regular heap allocation.
474    pub(crate) struct LgallocHandle {
475        /// Lgalloc handle, set if the memory was allocated by lgalloc.
476        handle: Option<lgalloc::Handle>,
477        /// Pointer to the allocated memory. Always well-aligned, but can be dangling.
478        pointer: std::ptr::NonNull<u8>,
479        /// Capacity of the allocated memory in bytes.
480        capacity: usize,
481    }
482
483    // SAFETY: `LgallocHandle` exclusively owns its allocation (either an
484    // lgalloc handle or a leaked `Vec`), and the `NonNull<u8>` is just a pointer
485    // into that owned memory. Transferring the handle transfers ownership, so it
486    // is safe to send across threads. Required because timely's `BytesRefill`
487    // produces `Box<dyn DerefMut<Target=[u8]> + Send>`.
488    unsafe impl Send for LgallocHandle {}
489
490    impl std::ops::Deref for LgallocHandle {
491        type Target = [u8];
492        #[inline(always)]
493        fn deref(&self) -> &Self::Target {
494            unsafe { std::slice::from_raw_parts(self.pointer.as_ptr(), self.capacity) }
495        }
496    }
497
498    impl std::ops::DerefMut for LgallocHandle {
499        #[inline(always)]
500        fn deref_mut(&mut self) -> &mut Self::Target {
501            unsafe { std::slice::from_raw_parts_mut(self.pointer.as_ptr(), self.capacity) }
502        }
503    }
504
505    impl Drop for LgallocHandle {
506        fn drop(&mut self) {
507            // If we have a handle, it's lgalloc memory. Otherwise, it's a heap allocation.
508            if let Some(handle) = self.handle.take() {
509                lgalloc::deallocate(handle);
510            } else {
511                unsafe { Vec::from_raw_parts(self.pointer.as_ptr(), 0, self.capacity) };
512            }
513            // Update pointer and capacity such that we don't double-free.
514            self.pointer = std::ptr::NonNull::dangling();
515            self.capacity = 0;
516        }
517    }
518}
519
520#[cfg(test)]
521mod exert_policy_metrics_tests {
522    use super::*;
523
524    /// Registration is guarded process-wide, so under a test runner that
525    /// shares the process another test may have registered first, into its
526    /// own registry. Either way a second call must not panic, and this
527    /// registry holds both families or neither.
528    #[mz_ore::test]
529    fn registering_twice_does_not_panic() {
530        let registry = MetricsRegistry::new();
531        register_exert_policy_metrics(&registry);
532        register_exert_policy_metrics(&registry);
533        let names: Vec<_> = registry
534            .gather()
535            .into_iter()
536            .map(|family| family.name().to_string())
537            .filter(|name| name.starts_with("mz_arrangement_exert_policy"))
538            .collect();
539        assert!(matches!(names.len(), 0 | 2), "{names:?}");
540    }
541}