use std::sync::Arc;
use mz_cluster::client::{ClusterClient, ClusterSpec};
use mz_ore::metrics::MetricsRegistry;
use mz_ore::now::NowFn;
use mz_ore::tracing::TracingHandle;
use mz_persist_client::cache::PersistClientCache;
use mz_rocksdb::config::SharedWriteBufferManager;
use mz_storage_client::client::{StorageClient, StorageCommand, StorageResponse};
use mz_storage_types::connections::ConnectionContext;
use mz_txn_wal::operator::TxnsContext;
use timely::communication::Allocate;
use timely::worker::Worker as TimelyWorker;
use tokio::sync::mpsc;
use crate::metrics::StorageMetrics;
use crate::storage_state::{StorageInstanceContext, Worker};
#[derive(Clone)]
struct Config {
pub persist_clients: Arc<PersistClientCache>,
pub txns_ctx: TxnsContext,
pub tracing_handle: Arc<TracingHandle>,
pub now: NowFn,
pub connection_context: ConnectionContext,
pub instance_context: StorageInstanceContext,
pub metrics: StorageMetrics,
pub shared_rocksdb_write_buffer_manager: SharedWriteBufferManager,
}
pub fn serve(
metrics_registry: &MetricsRegistry,
persist_clients: Arc<PersistClientCache>,
txns_ctx: TxnsContext,
tracing_handle: Arc<TracingHandle>,
now: NowFn,
connection_context: ConnectionContext,
instance_context: StorageInstanceContext,
) -> Result<impl Fn() -> Box<dyn StorageClient>, anyhow::Error> {
let config = Config {
persist_clients,
txns_ctx,
tracing_handle,
now,
connection_context,
instance_context,
metrics: StorageMetrics::register_with(metrics_registry),
shared_rocksdb_write_buffer_manager: Default::default(),
};
let tokio_executor = tokio::runtime::Handle::current();
let timely_container = Arc::new(tokio::sync::Mutex::new(None));
let client_builder = move || {
let client = ClusterClient::new(
Arc::clone(&timely_container),
tokio_executor.clone(),
config.clone(),
);
let client: Box<dyn StorageClient> = Box::new(client);
client
};
Ok(client_builder)
}
impl ClusterSpec for Config {
type Command = StorageCommand;
type Response = StorageResponse;
fn run_worker<A: Allocate + 'static>(
&self,
timely_worker: &mut TimelyWorker<A>,
client_rx: crossbeam_channel::Receiver<(
crossbeam_channel::Receiver<StorageCommand>,
mpsc::UnboundedSender<StorageResponse>,
)>,
) {
Worker::new(
timely_worker,
client_rx,
self.metrics.clone(),
self.now.clone(),
self.connection_context.clone(),
self.instance_context.clone(),
Arc::clone(&self.persist_clients),
self.txns_ctx.clone(),
Arc::clone(&self.tracing_handle),
self.shared_rocksdb_write_buffer_manager.clone(),
)
.run();
}
}