mz_compute/server.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 dataflow server.
11
12use std::cell::RefCell;
13use std::collections::{BTreeMap, BTreeSet};
14use std::convert::Infallible;
15use std::fmt::Debug;
16use std::path::PathBuf;
17use std::rc::Rc;
18use std::sync::{Arc, Mutex};
19use std::thread::Thread;
20use std::time::{Duration, Instant};
21
22use anyhow::Error;
23use mz_cluster::client::{ClusterClient, ClusterSpec, GuestClusterClient};
24use mz_cluster_client::client::TimelyConfig;
25use mz_compute_client::protocol::command::ComputeCommand;
26use mz_compute_client::protocol::history::ComputeCommandHistory;
27use mz_compute_client::protocol::response::ComputeResponse;
28use mz_compute_client::service::ComputeClient;
29use mz_ore::halt;
30use mz_ore::metrics::MetricsRegistry;
31use mz_ore::now::NowFn;
32use mz_ore::tracing::TracingHandle;
33use mz_persist_client::cache::PersistClientCache;
34use mz_rocksdb::config::SharedWriteBufferManager;
35use mz_storage::internal_control::{InternalCommandSender, InternalStorageCommand};
36use mz_storage::metrics::StorageMetrics;
37use mz_storage::storage_state::{StorageInstanceContext, StorageState, Worker as StorageWorker};
38use mz_storage_client::client::{StorageClient, StorageCommand, StorageResponse};
39use mz_storage_types::connections::ConnectionContext;
40use mz_txn_wal::operator::TxnsContext;
41use timely::progress::Antichain;
42use timely::worker::Worker as TimelyWorker;
43use tokio::sync::mpsc;
44use tokio::sync::mpsc::error::SendError;
45use tracing::{info, trace, warn};
46use uuid::Uuid;
47
48use crate::command_channel::{self, StorageLaneInput, UnifiedCommand};
49use crate::compute_state::{
50 ActiveComputeState, ComputeState, PeekPermits, PendingPeek, ReportedFrontier,
51};
52use crate::metrics::{ComputeMetrics, WorkerMetrics};
53
54/// Caller-provided configuration for compute.
55#[derive(Clone, Debug)]
56pub struct ComputeInstanceContext {
57 /// A directory that can be used for scratch work.
58 pub scratch_directory: Option<PathBuf>,
59 /// Whether to set core affinity for Timely workers.
60 pub worker_core_affinity: bool,
61 /// Context required to connect to an external sink from compute,
62 /// like the `CopyToS3OneshotSink` compute sink.
63 pub connection_context: ConnectionContext,
64}
65
66/// Which of a process's compute runtimes a given runtime is.
67///
68/// A clusterd process runs a single `Solo` runtime by default. When an interactive runtime is
69/// configured, the process instead runs a `Maintenance` and an `Interactive` runtime side by side.
70/// The named roles share per-process resources (persist cache, metrics registry, log spans). The
71/// role distinguishes them so that only the globals-owning runtime runs the non-idempotent
72/// process-global initializers, and so metric series and log spans do not collide.
73///
74/// `Solo` exists so the single-runtime default stays behaviorally identical to a deployment without
75/// a second runtime: no `role` metric label, and it owns the process globals just as the sole
76/// runtime always has.
77#[derive(Clone, Copy, Debug, PartialEq, Eq)]
78pub enum ComputeRuntimeRole {
79 /// The sole runtime of a single-runtime process. Owns index maintenance and the process-global
80 /// initializers.
81 Solo,
82 /// The maintenance runtime of a two-runtime process. Owns index maintenance and the
83 /// process-global initializers.
84 Maintenance,
85 /// The interactive runtime of a two-runtime process. Shares the process globals owned by
86 /// maintenance and serves reads.
87 ///
88 /// Test-only until the interactive runtime exists to construct it. It is present because the
89 /// `role` label's entire purpose is that two named roles register into one process registry
90 /// without colliding, and nothing else can express that: `Solo` registers the same metric names
91 /// with no `role` label, so prometheus rejects it alongside a named role for differing label
92 /// dimensions rather than treating it as a second series. Verifying non-collision therefore
93 /// needs a second *named* role.
94 ///
95 /// TODO: drop the `cfg` when the interactive runtime lands and constructs this.
96 #[cfg(test)]
97 Interactive,
98}
99
100impl ComputeRuntimeRole {
101 /// The `role` metric/log label for this role, or `None` for `Solo`.
102 ///
103 /// `Solo` omits the label so a single-runtime deployment registers exactly as it did before a
104 /// second runtime existed, keeping exact-match dashboards and alerts unchanged.
105 pub fn label(self) -> Option<&'static str> {
106 match self {
107 ComputeRuntimeRole::Solo => None,
108 ComputeRuntimeRole::Maintenance => Some("maintenance"),
109 #[cfg(test)]
110 ComputeRuntimeRole::Interactive => Some("interactive"),
111 }
112 }
113
114 /// Whether this role runs the non-idempotent, process-global initializers.
115 ///
116 /// `Solo` and `Maintenance` run them. An interactive runtime shares the same process and
117 /// inherits the globals maintenance installs, so re-running them would either double-apply a
118 /// non-idempotent effect or race maintenance.
119 ///
120 /// NOTE: every role a release build can construct owns the globals, so this is constantly true
121 /// outside tests. The distinction becomes load-bearing when the interactive runtime lands.
122 pub fn owns_process_globals(self) -> bool {
123 matches!(
124 self,
125 ComputeRuntimeRole::Solo | ComputeRuntimeRole::Maintenance
126 )
127 }
128}
129
130/// Configures the server with compute-specific metrics.
131#[derive(Clone)]
132struct Config {
133 /// `persist` client cache.
134 pub persist_clients: Arc<PersistClientCache>,
135 /// Context necessary for rendering txn-wal operators.
136 pub txns_ctx: TxnsContext,
137 /// A process-global handle to tracing configuration.
138 pub tracing_handle: Arc<TracingHandle>,
139 /// Metrics exposed by compute replicas.
140 pub metrics: ComputeMetrics,
141 /// Other configuration for compute.
142 pub context: ComputeInstanceContext,
143 /// The process-global metrics registry.
144 pub metrics_registry: MetricsRegistry,
145 /// The number of timely workers per process.
146 pub workers_per_process: usize,
147 /// Bounds how many offloaded peek walks run at once, shared by every worker this server runs.
148 ///
149 /// NOTE: per compute runtime, not global. A process running a maintenance and an interactive
150 /// runtime calls `serve` twice and admits the bound once per call.
151 pub peek_permits: Arc<PeekPermits>,
152 /// Configuration for hosting storage objects on this cluster, if enabled.
153 pub storage_guest: Option<Arc<StorageGuestConfig>>,
154}
155
156/// A per-worker channel delivering storage client connections.
157type StorageClientRx = mpsc::UnboundedReceiver<(
158 Uuid,
159 mpsc::UnboundedReceiver<StorageCommand>,
160 mpsc::UnboundedSender<StorageResponse>,
161)>;
162
163/// Configuration for hosting storage objects on the compute cluster.
164pub struct StorageGuestConfig {
165 /// Per-worker channels delivering storage client connections, indexed by local worker index.
166 client_rxs: Mutex<Vec<Option<StorageClientRx>>>,
167 /// Metrics for storage objects.
168 metrics: StorageMetrics,
169 /// Function to get wall time now.
170 now: NowFn,
171 /// Configuration for source and sink connections.
172 connection_context: ConnectionContext,
173 /// Other configuration for storage instances.
174 instance_context: StorageInstanceContext,
175 /// Shared rocksdb write buffer manager.
176 shared_rocksdb_write_buffer_manager: SharedWriteBufferManager,
177}
178
179/// Initiates a timely dataflow computation, processing compute commands.
180pub async fn serve(
181 timely_config: TimelyConfig,
182 role: ComputeRuntimeRole,
183 metrics_registry: &MetricsRegistry,
184 persist_clients: Arc<PersistClientCache>,
185 txns_ctx: TxnsContext,
186 tracing_handle: Arc<TracingHandle>,
187 context: ComputeInstanceContext,
188) -> Result<impl Fn() -> Box<dyn ComputeClient> + use<>, Error> {
189 let workers_per_process = timely_config.workers;
190 let config = Config {
191 persist_clients,
192 txns_ctx,
193 tracing_handle,
194 metrics: ComputeMetrics::register_with(metrics_registry, role),
195 context,
196 metrics_registry: metrics_registry.clone(),
197 workers_per_process,
198 peek_permits: Arc::new(PeekPermits::new(workers_per_process)),
199 storage_guest: None,
200 };
201
202 let (_worker_threads, client_builder) = serve_inner(config, timely_config).await?;
203 Ok(client_builder)
204}
205
206/// Initiates a timely dataflow computation that processes compute commands and additionally hosts
207/// storage objects, processing storage commands received over a separate client connection.
208///
209/// Returns client builders for both the compute and the storage side.
210pub async fn serve_unified(
211 timely_config: TimelyConfig,
212 role: ComputeRuntimeRole,
213 metrics_registry: &MetricsRegistry,
214 persist_clients: Arc<PersistClientCache>,
215 txns_ctx: TxnsContext,
216 tracing_handle: Arc<TracingHandle>,
217 context: ComputeInstanceContext,
218 now: NowFn,
219 storage_connection_context: ConnectionContext,
220 storage_instance_context: StorageInstanceContext,
221) -> Result<
222 (
223 impl Fn() -> Box<dyn ComputeClient> + use<>,
224 impl Fn() -> Box<dyn StorageClient> + use<>,
225 ),
226 Error,
227> {
228 let workers_per_process = timely_config.workers;
229
230 // Per-worker channels over which storage client connections are delivered.
231 let mut storage_client_txs = Vec::new();
232 let mut storage_client_rxs = Vec::new();
233 for _ in 0..workers_per_process {
234 let (tx, rx) = mpsc::unbounded_channel();
235 storage_client_txs.push(tx);
236 storage_client_rxs.push(Some(rx));
237 }
238
239 let storage_guest = StorageGuestConfig {
240 client_rxs: Mutex::new(storage_client_rxs),
241 metrics: StorageMetrics::register_with(metrics_registry),
242 now,
243 connection_context: storage_connection_context,
244 instance_context: storage_instance_context,
245 shared_rocksdb_write_buffer_manager: Default::default(),
246 };
247
248 let config = Config {
249 persist_clients,
250 txns_ctx,
251 tracing_handle,
252 metrics: ComputeMetrics::register_with(metrics_registry, role),
253 context,
254 metrics_registry: metrics_registry.clone(),
255 workers_per_process,
256 peek_permits: Arc::new(PeekPermits::new(workers_per_process)),
257 storage_guest: Some(Arc::new(storage_guest)),
258 };
259
260 let (worker_threads, compute_client_builder) = serve_inner(config, timely_config).await?;
261
262 let storage_client_txs = Arc::new(storage_client_txs);
263 let storage_client_builder = move || {
264 let client =
265 GuestClusterClient::new(Arc::clone(&storage_client_txs), worker_threads.clone());
266 let client: Box<dyn StorageClient> = Box::new(client);
267 client
268 };
269
270 Ok((compute_client_builder, storage_client_builder))
271}
272
273/// Builds the Timely cluster for the given config and returns its worker threads along with a
274/// builder for compute clients to it.
275async fn serve_inner(
276 config: Config,
277 timely_config: TimelyConfig,
278) -> Result<(Vec<Thread>, impl Fn() -> Box<dyn ComputeClient> + use<>), Error> {
279 mz_timely_util::column_pager::metrics::register(
280 &config.metrics_registry,
281 mz_timely_util::column_pager::tiered_policy(),
282 );
283 mz_timely_util::pool_config::metrics::register(&config.metrics_registry);
284 mz_cluster::client::register_exert_policy_metrics(&config.metrics_registry);
285
286 let tokio_executor = tokio::runtime::Handle::current();
287
288 let timely_container = config.build_cluster(timely_config, tokio_executor).await?;
289 let worker_threads = timely_container.worker_threads();
290 let timely_container = Arc::new(Mutex::new(timely_container));
291
292 let client_builder = move || {
293 let client = ClusterClient::new(Arc::clone(&timely_container));
294 let client: Box<dyn ComputeClient> = Box::new(client);
295 client
296 };
297
298 Ok((worker_threads, client_builder))
299}
300
301/// Error type returned on connection nonce changes.
302///
303/// A nonce change informs workers that subsequent commands come a from a new client connection
304/// and therefore require reconciliation.
305struct NonceChange(Uuid);
306
307/// Endpoint used by workers to receive compute commands.
308///
309/// Observes nonce changes in the command stream and converts them into receive errors.
310struct CommandReceiver {
311 /// The channel supplying commands.
312 inner: command_channel::Receiver,
313 /// The ID of the Timely worker.
314 worker_id: usize,
315 /// The nonce identifying the current cluster protocol incarnation.
316 nonce: Option<Uuid>,
317 /// A stash to enable peeking the next command, used in `try_recv`.
318 stashed_command: Option<ComputeCommand>,
319}
320
321impl CommandReceiver {
322 fn new(inner: command_channel::Receiver, worker_id: usize) -> Self {
323 Self {
324 inner,
325 worker_id,
326 nonce: None,
327 stashed_command: None,
328 }
329 }
330
331 /// Receive the next pending command, if any.
332 ///
333 /// If the next compute command has a different nonce, this method instead returns an `Err`
334 /// containing the new nonce.
335 fn try_recv(&mut self) -> Result<Option<WorkerCommand>, NonceChange> {
336 if let Some(command) = self.stashed_command.take() {
337 return Ok(Some(WorkerCommand::Compute(command)));
338 }
339 let Some(message) = self.inner.try_recv() else {
340 return Ok(None);
341 };
342
343 let (command, nonce) = match message {
344 UnifiedCommand::Compute(command, nonce) => (command, nonce),
345 UnifiedCommand::Storage(command) => {
346 trace!(
347 worker = self.worker_id,
348 ?command,
349 "received storage command"
350 );
351 return Ok(Some(WorkerCommand::Storage(command)));
352 }
353 };
354
355 trace!(worker = self.worker_id, %nonce, ?command, "received command");
356
357 if Some(nonce) == self.nonce {
358 Ok(Some(WorkerCommand::Compute(command)))
359 } else {
360 self.nonce = Some(nonce);
361 self.stashed_command = Some(command);
362 Err(NonceChange(nonce))
363 }
364 }
365}
366
367/// A command dispatched to the worker from the command channel.
368enum WorkerCommand {
369 /// A compute command.
370 Compute(ComputeCommand),
371 /// A storage-internal command, to be dispatched to the storage guest.
372 Storage(InternalStorageCommand),
373}
374
375/// Endpoint used by workers to send sending compute responses.
376///
377/// Tags responses with the current nonce, allowing receivers to filter out responses intended for
378/// previous client connections.
379pub(crate) struct ResponseSender {
380 /// The channel consuming responses.
381 inner: mpsc::UnboundedSender<(ComputeResponse, Uuid)>,
382 /// The ID of the Timely worker.
383 worker_id: usize,
384 /// The nonce identifying the current cluster protocol incarnation.
385 nonce: Option<Uuid>,
386}
387
388impl ResponseSender {
389 /// `pub(crate)` rather than private so the peek tests can build the sender a worker holds.
390 pub(crate) fn new(
391 inner: mpsc::UnboundedSender<(ComputeResponse, Uuid)>,
392 worker_id: usize,
393 ) -> Self {
394 Self {
395 inner,
396 worker_id,
397 nonce: None,
398 }
399 }
400
401 /// Set the cluster protocol nonce.
402 pub(crate) fn set_nonce(&mut self, nonce: Uuid) {
403 self.nonce = Some(nonce);
404 }
405
406 /// Send a compute response.
407 pub fn send(&self, response: ComputeResponse) -> Result<(), SendError<ComputeResponse>> {
408 let nonce = self.nonce.expect("nonce must be initialized");
409
410 trace!(worker = self.worker_id, %nonce, ?response, "sending response");
411 self.inner
412 .send((response, nonce))
413 .map_err(|SendError((resp, _))| SendError(resp))
414 }
415}
416
417/// State maintained for each worker thread.
418///
419/// Much of this state can be viewed as local variables for the worker thread,
420/// holding state that persists across function calls.
421struct Worker<'w> {
422 /// The underlying Timely worker.
423 timely_worker: &'w mut TimelyWorker,
424 /// The channel over which commands are received.
425 command_rx: CommandReceiver,
426 /// The channel over which responses are sent.
427 response_tx: ResponseSender,
428 compute_state: Option<ComputeState>,
429 /// Compute metrics.
430 metrics: WorkerMetrics,
431 /// A process-global cache of (blob_uri, consensus_uri) -> PersistClient.
432 /// This is intentionally shared between workers
433 persist_clients: Arc<PersistClientCache>,
434 /// Context necessary for rendering txn-wal operators.
435 txns_ctx: TxnsContext,
436 /// A process-global handle to tracing configuration.
437 tracing_handle: Arc<TracingHandle>,
438 context: ComputeInstanceContext,
439 /// The process-global metrics registry.
440 metrics_registry: MetricsRegistry,
441 /// The number of timely workers per process.
442 workers_per_process: usize,
443 /// Bounds how many offloaded peek walks run at once, shared by the workers of one `serve`
444 /// call rather than by the process.
445 peek_permits: Arc<PeekPermits>,
446 /// The hosted storage guest, if any.
447 storage: Option<StorageGuest>,
448}
449
450/// Per-worker state for hosting storage objects on the compute cluster.
451struct StorageGuest {
452 /// Channel delivering new storage client connections.
453 client_rx: StorageClientRx,
454 /// The current storage client connection, if any.
455 conn: Option<StorageConn>,
456 /// The hosted storage worker state.
457 storage_state: StorageState,
458 /// The last time storage maintenance ran.
459 last_maintenance: Instant,
460 /// The last time storage statistics were reported.
461 last_stats_time: Instant,
462}
463
464impl StorageGuest {
465 /// The longest the worker may park before the guest's next periodic duty (frontier reporting
466 /// or statistics collection) comes due, or `None` when no duty is pending.
467 ///
468 /// Mirrors the parking of storage's own server loop: the maintenance and statistics intervals
469 /// bound the park. A maintenance deadline in the past does not bound it, because maintenance
470 /// runs on the next wakeup anyway. The initial zero maintenance interval would otherwise turn
471 /// every park into a spin.
472 fn park_cap(&self) -> Option<Duration> {
473 // Periodic duties run only on a reconciled connection. Without one there is no deadline
474 // to meet, and connection and command arrivals unpark the worker.
475 let conn_serving = self
476 .conn
477 .as_ref()
478 .is_some_and(|conn| conn.reconcile_buf.is_none());
479 if !conn_serving {
480 return None;
481 }
482
483 let maintenance_interval = self.storage_state.server_maintenance_interval;
484 let stats_interval = self
485 .storage_state
486 .storage_configuration
487 .parameters
488 .statistics_collection_interval;
489
490 let next_maintenance =
491 (self.last_maintenance + maintenance_interval).checked_duration_since(Instant::now());
492 let next_stats = stats_interval.saturating_sub(self.last_stats_time.elapsed());
493 match next_maintenance {
494 Some(maintenance) => Some(maintenance.min(next_stats)),
495 None => Some(next_stats),
496 }
497 }
498}
499
500/// A storage client connection.
501struct StorageConn {
502 /// The channel over which storage commands are received.
503 command_rx: mpsc::UnboundedReceiver<StorageCommand>,
504 /// The channel over which storage responses are sent.
505 response_tx: mpsc::UnboundedSender<StorageResponse>,
506 /// Commands buffered for reconciliation, until `InitializationComplete` is received.
507 /// `None` once the connection is reconciled and serving.
508 reconcile_buf: Option<Vec<StorageCommand>>,
509}
510
511impl ClusterSpec for Config {
512 type Command = ComputeCommand;
513 type Response = ComputeResponse;
514
515 const NAME: &str = "compute";
516
517 fn run_worker(
518 &self,
519 timely_worker: &mut TimelyWorker,
520 client_rx: mpsc::UnboundedReceiver<(
521 Uuid,
522 mpsc::UnboundedReceiver<ComputeCommand>,
523 mpsc::UnboundedSender<ComputeResponse>,
524 )>,
525 ) {
526 if self.context.worker_core_affinity {
527 set_core_affinity(timely_worker.index());
528 }
529
530 let worker_id = timely_worker.index();
531 let metrics = self.metrics.for_worker(worker_id);
532
533 let local_index = worker_id % self.workers_per_process;
534
535 // Prepare the storage guest's inputs to the command channel, so
536 // storage-internal commands are sequenced through the same lane as compute commands.
537 let mut storage_lane_input = None;
538 let guest_setup = self.storage_guest.as_ref().map(|cfg| {
539 let storage_client_rx = cfg.client_rxs.lock().expect("poisoned")[local_index]
540 .take()
541 .expect("each worker takes its storage client_rx exactly once");
542 let (internal_tx, internal_rx) = std::sync::mpsc::channel();
543 let activator_slot = Rc::new(RefCell::new(None));
544 storage_lane_input = Some(StorageLaneInput {
545 rx: internal_rx,
546 activator_slot: Rc::clone(&activator_slot),
547 });
548 let internal_cmd_tx = InternalCommandSender::from_parts(internal_tx, activator_slot);
549 (Arc::clone(cfg), storage_client_rx, internal_cmd_tx)
550 });
551
552 // Create the command channel that broadcasts commands from worker 0 to other workers. We
553 // reuse this channel between client connections, to avoid bugs where different workers end
554 // up creating incompatible sides of the channel dataflow after reconnects.
555 // See database-issues#8964.
556 let (cmd_tx, cmd_rx) = command_channel::render(timely_worker, storage_lane_input);
557 let (resp_tx, resp_rx) = mpsc::unbounded_channel();
558
559 spawn_channel_adapter(client_rx, cmd_tx, resp_rx, worker_id);
560
561 // Create the storage guest state.
562 let storage = guest_setup.map(|(cfg, storage_client_rx, internal_cmd_tx)| {
563 let storage_state = StorageState::new_guest(
564 timely_worker.index(),
565 timely_worker.peers(),
566 internal_cmd_tx,
567 // The host dispatches internal commands from the unified command channel, so
568 // the guest reads no receiver of its own.
569 None,
570 cfg.metrics.clone(),
571 cfg.now.clone(),
572 cfg.connection_context.clone(),
573 cfg.instance_context.clone(),
574 Arc::clone(&self.persist_clients),
575 self.txns_ctx.clone(),
576 Arc::clone(&self.tracing_handle),
577 cfg.shared_rocksdb_write_buffer_manager.clone(),
578 );
579
580 StorageGuest {
581 client_rx: storage_client_rx,
582 conn: None,
583 storage_state,
584 last_maintenance: Instant::now(),
585 last_stats_time: Instant::now(),
586 }
587 });
588
589 Worker {
590 timely_worker,
591 command_rx: CommandReceiver::new(cmd_rx, worker_id),
592 response_tx: ResponseSender::new(resp_tx, worker_id),
593 metrics,
594 context: self.context.clone(),
595 persist_clients: Arc::clone(&self.persist_clients),
596 txns_ctx: self.txns_ctx.clone(),
597 compute_state: None,
598 tracing_handle: Arc::clone(&self.tracing_handle),
599 metrics_registry: self.metrics_registry.clone(),
600 workers_per_process: self.workers_per_process,
601 peek_permits: Arc::clone(&self.peek_permits),
602 storage,
603 }
604 .run()
605 }
606}
607
608/// Set the current thread's core affinity, based on the given `worker_id`.
609#[cfg(not(target_os = "macos"))]
610fn set_core_affinity(worker_id: usize) {
611 use tracing::error;
612
613 let Some(mut core_ids) = core_affinity::get_core_ids() else {
614 error!(worker_id, "unable to get core IDs for setting affinity");
615 return;
616 };
617
618 // The `get_core_ids` docs don't say anything about a guaranteed order of the returned Vec,
619 // so sort it just to be safe.
620 core_ids.sort_unstable_by_key(|i| i.id);
621
622 // On multi-process replicas `worker_id` might be greater than the number of available cores.
623 // However, we assume that we always have at least as many cores as there are local workers.
624 // Violating this assumption is safe but might lead to degraded performance due to skew in core
625 // utilization.
626 let idx = worker_id % core_ids.len();
627 let core_id = core_ids[idx];
628
629 if core_affinity::set_for_current(core_id) {
630 info!(
631 worker_id,
632 core_id = core_id.id,
633 "set core affinity for worker"
634 );
635 } else {
636 error!(
637 worker_id,
638 core_id = core_id.id,
639 "failed to set core affinity for worker"
640 )
641 }
642}
643
644/// Set the current thread's core affinity, based on the given `worker_id`.
645#[cfg(target_os = "macos")]
646fn set_core_affinity(_worker_id: usize) {
647 // Setting core affinity is known to not work on Apple Silicon:
648 // https://github.com/Elzair/core_affinity_rs/issues/22
649 info!("setting core affinity is not supported on macOS");
650}
651
652impl<'w> Worker<'w> {
653 /// Runs a compute worker.
654 pub fn run(&mut self) {
655 // The command receiver is initialized without an nonce, so receiving the first command
656 // always triggers a nonce change.
657 let NonceChange(nonce) = self.recv_command().expect_err("change to first nonce");
658 self.set_nonce(nonce);
659
660 loop {
661 let Err(NonceChange(nonce)) = self.run_client();
662 self.set_nonce(nonce);
663 }
664 }
665
666 fn set_nonce(&mut self, nonce: Uuid) {
667 self.response_tx.set_nonce(nonce);
668 }
669
670 /// Handles commands for a client connection, returns when the nonce changes.
671 fn run_client(&mut self) -> Result<Infallible, NonceChange> {
672 self.reconcile()?;
673
674 // The last time we did periodic maintenance.
675 let mut last_maintenance = Instant::now();
676
677 // Commence normal operation.
678 loop {
679 // Get the maintenance interval, default to zero if we don't have a compute state.
680 let maintenance_interval = self
681 .compute_state
682 .as_ref()
683 .map_or(Duration::ZERO, |state| state.server_maintenance_interval);
684
685 let now = Instant::now();
686 // Determine if we need to perform maintenance, which is true if `maintenance_interval`
687 // time has passed since the last maintenance.
688 let sleep_duration;
689 if now >= last_maintenance + maintenance_interval {
690 last_maintenance = now;
691 sleep_duration = None;
692
693 // Report frontier information back the coordinator.
694 if let Some(mut compute_state) = self.activate_compute() {
695 compute_state.compute_state.traces.maintenance();
696 compute_state.report_frontiers();
697 compute_state.report_metrics();
698 compute_state.check_expiration();
699 }
700
701 self.metrics.record_shared_row_metrics();
702 } else {
703 // We didn't perform maintenance, sleep until the next maintenance interval.
704 let next_maintenance = last_maintenance + maintenance_interval;
705 sleep_duration = Some(next_maintenance.saturating_duration_since(now))
706 };
707
708 // Do not sleep while a peek waits for its turn. Only the sweep below gives it one,
709 // and nothing else leaves an activation behind to end the park.
710 let sleep_duration = match &self.compute_state {
711 Some(state) if state.peeks_awaiting_turn() => Some(Duration::ZERO),
712 _ => sleep_duration,
713 };
714
715 // With a storage guest, cap the park duration so the guest's periodic duties run on
716 // time.
717 let sleep_duration = match self.storage.as_ref().and_then(StorageGuest::park_cap) {
718 Some(cap) => Some(sleep_duration.map_or(cap, |d| d.min(cap))),
719 None => sleep_duration,
720 };
721
722 // Step the timely worker, recording the time taken.
723 let timer = self.metrics.timely_step_duration_seconds.start_timer();
724 if self.storage_guest_busy() {
725 self.timely_worker.step();
726 } else {
727 self.timely_worker.step_or_park(sleep_duration);
728 }
729 timer.observe_duration();
730
731 self.handle_pending_commands()?;
732
733 self.process_storage_guest();
734
735 if let Some(mut compute_state) = self.activate_compute() {
736 compute_state.process_peeks();
737 compute_state.process_subscribes();
738 compute_state.process_copy_tos();
739 }
740 }
741 }
742
743 fn handle_pending_commands(&mut self) -> Result<(), NonceChange> {
744 while let Some(cmd) = self.command_rx.try_recv()? {
745 match cmd {
746 WorkerCommand::Compute(cmd) => self.handle_command(cmd),
747 WorkerCommand::Storage(cmd) => self.handle_storage_internal_command(cmd),
748 }
749 }
750 Ok(())
751 }
752
753 /// Whether the storage guest has pending work that forbids parking.
754 ///
755 /// It is critical that we allow Timely to park iff there are no pending commands or async
756 /// worker responses, since those are delivered by other threads that only unpark us once, at
757 /// send time.
758 fn storage_guest_busy(&self) -> bool {
759 self.storage.as_ref().is_some_and(|guest| {
760 !guest.client_rx.is_empty()
761 || guest
762 .conn
763 .as_ref()
764 .is_some_and(|conn| !conn.command_rx.is_empty())
765 || !guest.storage_state.async_worker.is_empty()
766 })
767 }
768
769 /// Dispatch a storage-internal command from the command channel to
770 /// the storage guest. This is where all storage dataflow rendering happens.
771 fn handle_storage_internal_command(&mut self, cmd: InternalStorageCommand) {
772 let mut guest = self
773 .storage
774 .take()
775 .expect("the command channel carries storage commands only when a guest is hosted");
776
777 let mut worker = StorageWorker {
778 timely_worker: &mut *self.timely_worker,
779 client_rx: guest.client_rx,
780 storage_state: guest.storage_state,
781 };
782 worker.handle_internal_storage_command(cmd);
783
784 let StorageWorker {
785 timely_worker: _,
786 client_rx,
787 storage_state,
788 } = worker;
789 guest.client_rx = client_rx;
790 guest.storage_state = storage_state;
791 self.storage = Some(guest);
792 }
793
794 /// Process the storage guest's per-iteration duties: accept client
795 /// connections, handle external storage commands (buffering for reconciliation until
796 /// `InitializationComplete`), forward async worker responses, and report frontiers, dropped
797 /// collections, status updates, and statistics.
798 fn process_storage_guest(&mut self) {
799 let Some(mut guest) = self.storage.take() else {
800 return;
801 };
802
803 // Accept new client connections, replacing any current one. Every new connection starts
804 // with a reconciliation.
805 loop {
806 use tokio::sync::mpsc::error::TryRecvError;
807 match guest.client_rx.try_recv() {
808 Ok((_nonce, command_rx, response_tx)) => {
809 guest.conn = Some(StorageConn {
810 command_rx,
811 response_tx,
812 reconcile_buf: Some(Vec::new()),
813 });
814 }
815 Err(TryRecvError::Empty | TryRecvError::Disconnected) => break,
816 }
817 }
818
819 let mut worker = StorageWorker {
820 timely_worker: &mut *self.timely_worker,
821 client_rx: guest.client_rx,
822 storage_state: guest.storage_state,
823 };
824
825 // Handle responses from the async worker. Only worker 0 does async processing, so only
826 // worker 0 ever receives any. This must run with or without a client connection:
827 // storage-internal commands keep flowing through the command channel while the
828 // controller is away and can issue async work, and an undrained response queue keeps
829 // `storage_guest_busy` true, turning every park into a spin.
830 while let Ok(response) = worker.storage_state.async_worker.try_recv() {
831 worker.handle_async_worker_response(response);
832 }
833
834 let Some(mut conn) = guest.conn.take() else {
835 let StorageWorker {
836 timely_worker: _,
837 client_rx,
838 storage_state,
839 } = worker;
840 guest.client_rx = client_rx;
841 guest.storage_state = storage_state;
842 self.storage = Some(guest);
843 return;
844 };
845
846 // Drain external storage commands.
847 let mut disconnected = false;
848 loop {
849 use tokio::sync::mpsc::error::TryRecvError;
850 match conn.command_rx.try_recv() {
851 Ok(StorageCommand::InitializationComplete) if conn.reconcile_buf.is_some() => {
852 let commands = conn.reconcile_buf.take().expect("checked above");
853 worker.reconcile_commands(commands);
854 }
855 Ok(cmd) => match &mut conn.reconcile_buf {
856 Some(buf) => buf.push(cmd),
857 None => worker.storage_state.handle_storage_command(cmd),
858 },
859 Err(TryRecvError::Empty) => break,
860 Err(TryRecvError::Disconnected) => {
861 disconnected = true;
862 break;
863 }
864 }
865 }
866
867 // Response-producing duties run only on a reconciled connection.
868 if conn.reconcile_buf.is_none() {
869 let maintenance_interval = worker.storage_state.server_maintenance_interval;
870 let now = Instant::now();
871 if now >= guest.last_maintenance + maintenance_interval {
872 guest.last_maintenance = now;
873 worker.report_frontier_progress(&conn.response_tx);
874 }
875
876 for id in std::mem::take(&mut worker.storage_state.dropped_ids) {
877 worker.send_storage_response(&conn.response_tx, StorageResponse::DroppedId(id));
878 }
879
880 worker.process_oneshot_ingestions(&conn.response_tx);
881 worker.report_status_updates(&conn.response_tx);
882
883 let stats_interval = worker
884 .storage_state
885 .storage_configuration
886 .parameters
887 .statistics_collection_interval;
888 if guest.last_stats_time.elapsed() >= stats_interval {
889 worker.report_storage_statistics(&conn.response_tx);
890 guest.last_stats_time = Instant::now();
891 }
892 }
893
894 let StorageWorker {
895 timely_worker: _,
896 client_rx,
897 storage_state,
898 } = worker;
899 guest.client_rx = client_rx;
900 guest.storage_state = storage_state;
901 guest.conn = (!disconnected).then_some(conn);
902 self.storage = Some(guest);
903 }
904
905 fn handle_command(&mut self, cmd: ComputeCommand) {
906 if matches!(&cmd, ComputeCommand::CreateInstance(_)) {
907 self.compute_state = Some(ComputeState::new(
908 Arc::clone(&self.persist_clients),
909 self.txns_ctx.clone(),
910 self.metrics.clone(),
911 Arc::clone(&self.tracing_handle),
912 self.context.clone(),
913 self.metrics_registry.clone(),
914 self.workers_per_process,
915 Arc::clone(&self.peek_permits),
916 ));
917 }
918 self.activate_compute().unwrap().handle_compute_command(cmd);
919 }
920
921 fn activate_compute(&mut self) -> Option<ActiveComputeState<'_>> {
922 if let Some(compute_state) = &mut self.compute_state {
923 Some(ActiveComputeState {
924 timely_worker: &mut *self.timely_worker,
925 compute_state,
926 response_tx: &mut self.response_tx,
927 })
928 } else {
929 None
930 }
931 }
932
933 /// Receive the next compute command.
934 ///
935 /// This method blocks if no command is currently available, but takes care to step the Timely
936 /// worker while doing so.
937 fn recv_command(&mut self) -> Result<ComputeCommand, NonceChange> {
938 loop {
939 if let Some(cmd) = self.command_rx.try_recv()? {
940 match cmd {
941 WorkerCommand::Compute(cmd) => return Ok(cmd),
942 // Storage-internal commands are dispatched even while
943 // waiting for compute commands (e.g. during compute reconciliation), so
944 // storage dataflow construction keeps its lane position on all workers.
945 WorkerCommand::Storage(cmd) => {
946 self.handle_storage_internal_command(cmd);
947 continue;
948 }
949 }
950 }
951
952 // Keep serving the storage guest while blocked on compute
953 // commands, and avoid unbounded parks that would stall its periodic duties.
954 self.process_storage_guest();
955 let park_cap = self.storage.as_ref().and_then(StorageGuest::park_cap);
956
957 let start = Instant::now();
958 if self.storage_guest_busy() {
959 self.timely_worker.step();
960 } else if let Some(cap) = park_cap {
961 self.timely_worker.step_or_park(Some(cap));
962 } else {
963 self.timely_worker.step_or_park(None);
964 }
965 self.metrics
966 .timely_step_duration_seconds
967 .observe(start.elapsed().as_secs_f64());
968 }
969 }
970
971 /// Extract commands until `InitializationComplete`, and make the worker reflect those commands.
972 ///
973 /// This method is meant to be a function of the commands received thus far (as recorded in the
974 /// compute state command history) and the new commands from `command_rx`. It should not be a
975 /// function of other characteristics, like whether the worker has managed to respond to a peek
976 /// or not. Some effort goes in to narrowing our view to only the existing commands we can be sure
977 /// are live at all other workers.
978 ///
979 /// The methodology here is to drain `command_rx` until an `InitializationComplete`, at which point
980 /// the prior commands are "reconciled" in. Reconciliation takes each goal dataflow and looks for an
981 /// existing "compatible" dataflow (per `compatible()`) it can repurpose, with some additional tests
982 /// to be sure that we can cut over from one to the other (no additional compaction, no tails/sinks).
983 /// With any connections established, old orphaned dataflows are allow to compact away, and any new
984 /// dataflows are created from scratch. "Kept" dataflows are allowed to compact up to any new `as_of`.
985 ///
986 /// Some additional tidying happens, cleaning up pending peeks, reported frontiers, and creating a new
987 /// subscribe response buffer. We will need to be vigilant with future modifications to `ComputeState` to
988 /// line up changes there with clean resets here.
989 fn reconcile(&mut self) -> Result<(), NonceChange> {
990 // To initialize the connection, we want to drain all commands until we receive a
991 // `ComputeCommand::InitializationComplete` command to form a target command state.
992 let mut new_commands = Vec::new();
993 loop {
994 match self.recv_command()? {
995 ComputeCommand::InitializationComplete => break,
996 command => new_commands.push(command),
997 }
998 }
999
1000 // Commands we will need to apply before entering normal service.
1001 // These commands may include dropping existing dataflows, compacting existing dataflows,
1002 // and creating new dataflows, in addition to standard peek and compaction commands.
1003 // The result should be the same as if dropping all dataflows and running `new_commands`.
1004 let mut todo_commands = Vec::new();
1005 // We only have a compute history if we are in an initialized state
1006 // (i.e. after a `CreateInstance`).
1007 // If this is not the case, just copy `new_commands` into `todo_commands`.
1008 if let Some(compute_state) = &mut self.compute_state {
1009 // Reduce the installed commands.
1010 // Importantly, act as if all peeks may have been retired (as we cannot know otherwise).
1011 compute_state.command_history.discard_peeks();
1012 compute_state.command_history.reduce();
1013
1014 // At this point, we need to sort out which of the *certainly installed* dataflows are
1015 // suitable replacements for the requested dataflows. A dataflow is "certainly installed"
1016 // as of a frontier if its compaction allows it to go no further. We ignore peeks for this
1017 // reasoning, as we cannot be certain that peeks still exist at any other worker.
1018
1019 // Having reduced our installed command history retaining no peeks (above), we should be able
1020 // to use track down installed dataflows we can use as surrogates for requested dataflows (which
1021 // have retained all of their peeks, creating a more demanding `as_of` requirement).
1022 // NB: installed dataflows may still be allowed to further compact, and we should double check
1023 // this before being too confident. It should be rare without peeks, but could happen with e.g.
1024 // multiple outputs of a dataflow.
1025
1026 // The values with which a prior `CreateInstance` was called, if it was.
1027 let mut old_instance_config = None;
1028 // Index dataflows by `export_ids().collect()`, as this is a precondition for their compatibility.
1029 let mut old_dataflows = BTreeMap::default();
1030 // Maintain allowed compaction, in case installed identifiers may have been allowed to compact.
1031 let mut old_frontiers = BTreeMap::default();
1032 for command in compute_state.command_history.iter() {
1033 match command {
1034 ComputeCommand::CreateInstance(config) => {
1035 old_instance_config = Some(config);
1036 }
1037 ComputeCommand::CreateDataflow(dataflow) => {
1038 let export_ids = dataflow.export_ids().collect::<BTreeSet<_>>();
1039 old_dataflows.insert(export_ids, dataflow);
1040 }
1041 ComputeCommand::AllowCompaction { id, frontier } => {
1042 old_frontiers.insert(id, frontier);
1043 }
1044 _ => {
1045 // Nothing to do in these cases.
1046 }
1047 }
1048 }
1049
1050 // Compaction commands that can be applied to existing dataflows.
1051 let mut old_compaction = BTreeMap::default();
1052 // Exported identifiers from dataflows we retain.
1053 let mut retain_ids = BTreeSet::default();
1054
1055 // Traverse new commands, sorting out what remediation we can do.
1056 for command in new_commands.iter() {
1057 match command {
1058 ComputeCommand::CreateDataflow(dataflow) => {
1059 // Attempt to find an existing match for the dataflow.
1060 let as_of = dataflow.as_of.as_ref().unwrap();
1061 let export_ids = dataflow.export_ids().collect::<BTreeSet<_>>();
1062
1063 if let Some(old_dataflow) = old_dataflows.get(&export_ids) {
1064 let compatible = old_dataflow.compatible_with(dataflow);
1065 let uncompacted = !export_ids
1066 .iter()
1067 .flat_map(|id| old_frontiers.get(id))
1068 .any(|frontier| {
1069 !timely::PartialOrder::less_equal(
1070 *frontier,
1071 dataflow.as_of.as_ref().unwrap(),
1072 )
1073 });
1074
1075 // We cannot reconcile subscribe and copy-to sinks at the moment,
1076 // because the response buffer is shared, and to a first approximation
1077 // must be completely reformed.
1078 let subscribe_free = dataflow.subscribe_ids().next().is_none();
1079 let copy_to_free = dataflow.copy_to_ids().next().is_none();
1080
1081 // If we have replaced any dependency of this dataflow, we need to
1082 // replace this dataflow, to make it use the replacement.
1083 let dependencies_retained = dataflow
1084 .imported_index_ids()
1085 .all(|id| retain_ids.contains(&id));
1086
1087 if compatible
1088 && uncompacted
1089 && subscribe_free
1090 && copy_to_free
1091 && dependencies_retained
1092 {
1093 // Match found; remove the match from the deletion queue,
1094 // and compact its outputs to the dataflow's `as_of`.
1095 old_dataflows.remove(&export_ids);
1096 for id in export_ids.iter() {
1097 old_compaction.insert(*id, as_of.clone());
1098 }
1099 retain_ids.extend(export_ids);
1100 } else {
1101 warn!(
1102 ?export_ids,
1103 ?compatible,
1104 ?uncompacted,
1105 ?subscribe_free,
1106 ?copy_to_free,
1107 ?dependencies_retained,
1108 old_as_of = ?old_dataflow.as_of,
1109 new_as_of = ?as_of,
1110 "dataflow reconciliation failed",
1111 );
1112
1113 // Dump the full dataflow plans if they are incompatible, to
1114 // simplify debugging hard-to-reproduce reconciliation failures.
1115 if !compatible {
1116 warn!(
1117 old = ?old_dataflow,
1118 new = ?dataflow,
1119 "incompatible dataflows in reconciliation",
1120 );
1121 }
1122
1123 todo_commands
1124 .push(ComputeCommand::CreateDataflow(dataflow.clone()));
1125 }
1126
1127 compute_state.metrics.record_dataflow_reconciliation(
1128 compatible,
1129 uncompacted,
1130 subscribe_free,
1131 copy_to_free,
1132 dependencies_retained,
1133 );
1134 } else {
1135 todo_commands.push(ComputeCommand::CreateDataflow(dataflow.clone()));
1136 }
1137 }
1138 ComputeCommand::CreateInstance(config) => {
1139 // Cluster creation should not be performed again!
1140 if old_instance_config.map_or(false, |old| !old.compatible_with(config)) {
1141 halt!(
1142 "new instance configuration not compatible with existing instance configuration:\n{:?}\nvs\n{:?}",
1143 config,
1144 old_instance_config,
1145 );
1146 }
1147 }
1148 // All other commands we apply as requested.
1149 command => {
1150 todo_commands.push(command.clone());
1151 }
1152 }
1153 }
1154
1155 // Issue compaction commands first to reclaim resources.
1156 for (_, dataflow) in old_dataflows.iter() {
1157 for id in dataflow.export_ids() {
1158 // We want to drop anything that has not yet been dropped,
1159 // and nothing that has already been dropped.
1160 if old_frontiers.get(&id) != Some(&&Antichain::new()) {
1161 old_compaction.insert(id, Antichain::new());
1162 }
1163 }
1164 }
1165 for (&id, frontier) in &old_compaction {
1166 let frontier = frontier.clone();
1167 todo_commands.insert(0, ComputeCommand::AllowCompaction { id, frontier });
1168 }
1169
1170 // Clean up worker-local state.
1171 //
1172 // Various aspects of `ComputeState` need to be either uninstalled, or return to a blank slate.
1173 // All dropped dataflows should clean up after themselves, as we plan to install new dataflows
1174 // re-using the same identifiers.
1175 // All re-used dataflows should roll back any believed communicated information (e.g. frontiers)
1176 // so that they recommunicate that information as if from scratch.
1177
1178 // Remove all peeks, whether they have started or are still awaiting a turn.
1179 let queued = std::mem::take(&mut compute_state.queued_peeks);
1180 let pending = std::mem::take(&mut compute_state.pending_peeks);
1181 for peek in queued.into_iter().map(PendingPeek::Index).chain(pending) {
1182 // Log dropping the peek request.
1183 if let Some(logger) = compute_state.compute_logger.as_mut() {
1184 logger.log(&peek.as_log_event(false));
1185 }
1186 }
1187
1188 for (&id, collection) in compute_state.collections.iter_mut() {
1189 // Adjust reported frontiers:
1190 // * For dataflows we continue to use, reset to ensure we report something not
1191 // before the new `as_of` next.
1192 // * For dataflows we drop, set to the empty frontier, to ensure we don't report
1193 // anything for them.
1194 let retained = retain_ids.contains(&id);
1195 let compaction = old_compaction.remove(&id);
1196 let new_reported_frontier = match (retained, compaction) {
1197 (true, Some(new_as_of)) => ReportedFrontier::NotReported { lower: new_as_of },
1198 (true, None) => {
1199 unreachable!("retained dataflows are compacted to the new as_of")
1200 }
1201 (false, Some(new_frontier)) => {
1202 assert!(new_frontier.is_empty());
1203 ReportedFrontier::Reported(new_frontier)
1204 }
1205 (false, None) => {
1206 // Logging dataflows are implicitly retained and don't have a new as_of.
1207 // Reset them to the minimal frontier.
1208 ReportedFrontier::new()
1209 }
1210 };
1211
1212 collection.reset_reported_frontiers(new_reported_frontier);
1213
1214 // Sink tokens should be retained for retained dataflows, and dropped for dropped
1215 // dataflows.
1216 //
1217 // Dropping the tokens of active subscribe and copy-tos makes them place
1218 // `DroppedAt` responses into the respective response buffer. We drop those buffers
1219 // in the next step, which ensures that we don't send out `DroppedAt` responses for
1220 // subscribe/copy-tos dropped during reconciliation.
1221 if !retained {
1222 collection.sink_token = None;
1223 }
1224 }
1225
1226 // We must drop the response buffers as they are global across all subscribe/copy-tos.
1227 // If they were broken out by `GlobalId` then we could drop only the response buffers
1228 // of dataflows we drop.
1229 compute_state.subscribe_response_buffer = Rc::new(RefCell::new(Vec::new()));
1230 compute_state.copy_to_response_buffer = Rc::new(RefCell::new(Vec::new()));
1231
1232 // The controller expects the logging collections to be readable from the minimum time
1233 // initially. We cannot recreate the logging arrangements without restarting the
1234 // instance, but we can pad the compacted times with empty data. Doing so is sound
1235 // because logging collections from different replica incarnations are considered
1236 // distinct TVCs, so the controller doesn't expect any historical consistency from
1237 // these collections when it reconnects to a replica.
1238 //
1239 // TODO(database-issues#8152): Consider resolving this with controller-side reconciliation instead.
1240 if let Some(config) = old_instance_config {
1241 for id in config.logging.index_logs.values() {
1242 let trace = compute_state
1243 .traces
1244 .remove(id)
1245 .expect("logging trace exists");
1246 let padded = trace.into_padded();
1247 compute_state.traces.set(*id, padded);
1248 }
1249 }
1250 } else {
1251 todo_commands.clone_from(&new_commands);
1252 }
1253
1254 // Execute the commands to bring us to `new_commands`.
1255 for command in todo_commands.into_iter() {
1256 self.handle_command(command);
1257 }
1258
1259 // Overwrite `self.command_history` to reflect `new_commands`.
1260 // It is possible that there still isn't a compute state yet.
1261 if let Some(compute_state) = &mut self.compute_state {
1262 let mut command_history = ComputeCommandHistory::new(self.metrics.for_history());
1263 for command in new_commands.iter() {
1264 command_history.push(command.clone());
1265 }
1266 compute_state.command_history = command_history;
1267 }
1268 Ok(())
1269 }
1270}
1271
1272/// Spawn a task to bridge between [`ClusterClient`] and [`Worker`] channels.
1273///
1274/// The [`Worker`] expects a pair of persistent channels, with punctuation marking reconnects,
1275/// while the [`ClusterClient`] provides a new pair of channels on each reconnect.
1276fn spawn_channel_adapter(
1277 mut client_rx: mpsc::UnboundedReceiver<(
1278 Uuid,
1279 mpsc::UnboundedReceiver<ComputeCommand>,
1280 mpsc::UnboundedSender<ComputeResponse>,
1281 )>,
1282 command_tx: command_channel::Sender,
1283 mut response_rx: mpsc::UnboundedReceiver<(ComputeResponse, Uuid)>,
1284 worker_id: usize,
1285) {
1286 mz_ore::task::spawn(
1287 || format!("compute-channel-adapter-{worker_id}"),
1288 async move {
1289 // To make workers aware of the individual client connections, we tag forwarded
1290 // commands with the client nonce. Additionally, we use the nonce to filter out
1291 // responses with a different nonce, which are intended for different client
1292 // connections.
1293 //
1294 // It's possible that we receive responses with nonces from the past but also from the
1295 // future: Worker 0 might have received a new nonce before us and broadcasted it to our
1296 // Timely cluster. When we receive a response with a future nonce, we need to wait with
1297 // forwarding it until we have received the same nonce from a client connection.
1298 //
1299 // Nonces are not ordered so we don't know whether a response nonce is from the past or
1300 // the future. We thus assume that every response with an unknown nonce might be from
1301 // the future and stash them all. Every time we reconnect, we immediately send all
1302 // stashed responses with a matching nonce. Every time we receive a new response with a
1303 // nonce that matches our current one, we can discard the entire response stash as we
1304 // know that all stashed responses must be from the past.
1305 let mut stashed_responses = BTreeMap::<Uuid, Vec<ComputeResponse>>::new();
1306
1307 while let Some((nonce, mut command_rx, response_tx)) = client_rx.recv().await {
1308 // Send stashed responses for this client.
1309 if let Some(resps) = stashed_responses.remove(&nonce) {
1310 for resp in resps {
1311 let _ = response_tx.send(resp);
1312 }
1313 }
1314
1315 // Wait for a new response while forwarding received commands.
1316 let mut serve_rx_channels = async || loop {
1317 tokio::select! {
1318 msg = command_rx.recv() => match msg {
1319 Some(cmd) => command_tx.send((cmd, nonce)),
1320 None => return Err(()),
1321 },
1322 msg = response_rx.recv() => {
1323 return Ok(msg.expect("worker connected"));
1324 }
1325 }
1326 };
1327
1328 // Serve this connection until we see any of the channels disconnect.
1329 loop {
1330 let Ok((resp, resp_nonce)) = serve_rx_channels().await else {
1331 break;
1332 };
1333
1334 if resp_nonce == nonce {
1335 // Response for the current connection; forward it.
1336 stashed_responses.clear();
1337 if response_tx.send(resp).is_err() {
1338 break;
1339 }
1340 } else {
1341 // Response for a past or future connection; stash it.
1342 let stash = stashed_responses.entry(resp_nonce).or_default();
1343 stash.push(resp);
1344 }
1345 }
1346 }
1347 },
1348 );
1349}