Skip to main content

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}