1use 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
40pub struct ClusterClient<C>
42where
43 C: ClusterSpec,
44 (C::Command, C::Response): Partitionable<C::Command, C::Response>,
45{
46 inner: Option<PartitionedClient<C::Command, C::Response>>,
48 timely_container: Arc<Mutex<TimelyContainer<C>>>,
50}
51
52pub struct TimelyContainer<C: ClusterSpec> {
54 client_txs: Vec<
56 mpsc::UnboundedSender<(
57 Uuid,
58 mpsc::UnboundedReceiver<C::Command>,
59 mpsc::UnboundedSender<C::Response>,
60 )>,
61 >,
62 worker_guards: WorkerGuards<()>,
64}
65
66impl<C: ClusterSpec> TimelyContainer<C> {
67 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
77pub struct GuestClusterClient<Cmd, Resp>
81where
82 (Cmd, Resp): Partitionable<Cmd, Resp>,
83{
84 client_txs: Arc<
86 Vec<
87 mpsc::UnboundedSender<(
88 Uuid,
89 mpsc::UnboundedReceiver<Cmd>,
90 mpsc::UnboundedSender<Resp>,
91 )>,
92 >,
93 >,
94 worker_threads: Vec<Thread>,
96 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 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 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 pub fn new(timely_container: Arc<Mutex<TimelyContainer<C>>>) -> Self {
194 Self {
195 timely_container,
196 inner: None,
197 }
198 }
199
200 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 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 async fn recv(&mut self) -> Result<Option<C::Response>, Error> {
261 if let Some(client) = self.inner.as_mut() {
262 client.recv().await
264 } else {
265 future::pending().await
266 }
267 }
268}
269
270#[async_trait]
275pub trait ClusterSpec: Clone + Send + Sync + 'static {
276 type Command: fmt::Debug + Send + TryIntoProtocolNonce;
278 type Response: fmt::Debug + Send;
280
281 const NAME: &str;
283
284 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 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 if config.arrangement_exert_proportionality > 0 {
334 let merge_effort = Some(1000);
335
336 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 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 EXERT_POLICY_MERGE_GRANTS.fetch_add(1, Ordering::Relaxed);
357 return merge_effort;
358 }
359
360 if !first && prop > 0 && len > 0 {
361 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 let span = info_span!("timely", name = Self::NAME, worker_id = worker_idx);
382 let _span_guard = span.enter();
383
384 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
404static 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
414pub 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 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 let mut alloc = vec![0_u8; size];
456 alloc.shrink_to_fit();
458 let pointer = std::ptr::NonNull::new(alloc.as_mut_ptr()).unwrap();
460 std::mem::forget(alloc);
462 LgallocHandle {
463 handle: None,
464 pointer,
465 capacity: size,
466 }
467 }
468 }
469 }
470
471 pub(crate) struct LgallocHandle {
475 handle: Option<lgalloc::Handle>,
477 pointer: std::ptr::NonNull<u8>,
479 capacity: usize,
481 }
482
483 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 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 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 #[mz_ore::test]
529 fn registering_twice_does_not_panic() {
530 let registry = MetricsRegistry::new();
531 register_exert_policy_metrics(®istry);
532 register_exert_policy_metrics(®istry);
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}