1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
use std::collections::HashMap;
use std::sync::{Arc, Mutex};
use anyhow::anyhow;
use mz_build_info::BuildInfo;
use mz_persist_client::PersistConfig;
use timely::communication::initialize::WorkerGuards;
use tokio::sync::mpsc;
use mz_ore::metrics::MetricsRegistry;
use mz_ore::now::NowFn;
use mz_persist_client::cache::PersistClientCache;
use mz_service::local::LocalClient;
use crate::protocol::client::StorageClient;
use crate::sink::SinkBaseMetrics;
use crate::source::metrics::SourceBaseMetrics;
use crate::storage_state::{StorageState, Worker};
use crate::types::connections::ConnectionContext;
use crate::DecodeMetrics;
pub struct Config {
pub build_info: &'static BuildInfo,
pub workers: usize,
pub timely_config: timely::Config,
pub now: NowFn,
pub metrics_registry: MetricsRegistry,
pub connection_context: ConnectionContext,
}
pub struct Server {
_worker_guards: WorkerGuards<()>,
}
pub fn serve(
config: Config,
) -> Result<(Server, impl Fn() -> Box<dyn StorageClient>), anyhow::Error> {
assert!(config.workers > 0);
let source_metrics = SourceBaseMetrics::register_with(&config.metrics_registry);
let sink_metrics = SinkBaseMetrics::register_with(&config.metrics_registry);
let decode_metrics = DecodeMetrics::register_with(&config.metrics_registry);
let metrics_bundle = (source_metrics, sink_metrics, decode_metrics);
let (client_txs, client_rxs): (Vec<_>, Vec<_>) = (0..config.workers)
.map(|_| crossbeam_channel::unbounded())
.unzip();
let client_rxs: Mutex<Vec<_>> = Mutex::new(client_rxs.into_iter().map(Some).collect());
let tokio_executor = tokio::runtime::Handle::current();
let now = config.now;
let persist_clients = PersistClientCache::new(
PersistConfig::new(config.build_info, now.clone()),
&config.metrics_registry,
);
let persist_clients = Arc::new(tokio::sync::Mutex::new(persist_clients));
let worker_guards = timely::execute::execute(config.timely_config, move |timely_worker| {
let timely_worker_index = timely_worker.index();
let timely_worker_peers = timely_worker.peers();
let _tokio_guard = tokio_executor.enter();
let client_rx = client_rxs.lock().unwrap()[timely_worker_index % config.workers]
.take()
.unwrap();
let (source_metrics, sink_metrics, decode_metrics) = metrics_bundle.clone();
let persist_clients = Arc::clone(&persist_clients);
Worker {
timely_worker,
client_rx,
storage_state: StorageState {
source_uppers: HashMap::new(),
source_tokens: HashMap::new(),
decode_metrics,
reported_frontiers: HashMap::new(),
ingestions: HashMap::new(),
exports: HashMap::new(),
now: now.clone(),
source_metrics,
sink_metrics,
timely_worker_index,
timely_worker_peers,
connection_context: config.connection_context.clone(),
persist_clients,
sink_tokens: HashMap::new(),
sink_write_frontiers: HashMap::new(),
},
}
.run()
})
.map_err(|e| anyhow!("{}", e))?;
let worker_threads = worker_guards
.guards()
.iter()
.map(|g| g.thread().clone())
.collect::<Vec<_>>();
let client_builder = move || {
let (command_txs, command_rxs): (Vec<_>, Vec<_>) = (0..config.workers)
.map(|_| crossbeam_channel::unbounded())
.unzip();
let (response_txs, response_rxs): (Vec<_>, Vec<_>) = (0..config.workers)
.map(|_| mpsc::unbounded_channel())
.unzip();
for (client_tx, channels) in client_txs
.iter()
.zip(command_rxs.into_iter().zip(response_txs))
{
client_tx
.send(channels)
.expect("worker should not drop first");
}
let client =
LocalClient::new_partitioned(response_rxs, command_txs, worker_threads.clone());
Box::new(client) as Box<dyn StorageClient>
};
let server = Server {
_worker_guards: worker_guards,
};
Ok((server, client_builder))
}