1use std::collections::HashMap;
16
17use anyhow::{Context, anyhow};
18use risingwave_common::secret::{LocalSecretManager, SecretEncryption};
19use risingwave_hummock_sdk::FrontendHummockVersion;
20use risingwave_meta::MetaResult;
21use risingwave_meta::controller::catalog::Catalog;
22use risingwave_meta::controller::fragment::FragmentServingInfo;
23use risingwave_meta::manager::MetadataManager;
24use risingwave_meta_model::{FragmentId, WorkerId};
25use risingwave_pb::backup_service::MetaBackupManifestId;
26use risingwave_pb::catalog::{Secret, Table};
27use risingwave_pb::common::worker_node::State::Running;
28use risingwave_pb::common::{ClusterResource, WorkerNode, WorkerType};
29use risingwave_pb::hummock::WriteLimits;
30use risingwave_pb::meta::meta_snapshot::SnapshotVersion;
31use risingwave_pb::meta::notification_service_server::NotificationService;
32use risingwave_pb::meta::{
33 FragmentWorkerSlotMapping, GetSessionParamsResponse, MetaSnapshot, PbTableCacheRefillPolicies,
34 SubscribeRequest, SubscribeType,
35};
36use risingwave_pb::user::UserInfo;
37use tokio::sync::mpsc;
38use tokio_stream::wrappers::UnboundedReceiverStream;
39use tonic::{Request, Response, Status};
40
41use crate::backup_restore::BackupManagerRef;
42use crate::hummock::HummockManagerRef;
43use crate::manager::{MetaSrvEnv, Notification, NotificationVersion, WorkerKey};
44use crate::serving::ServingVnodeMappingRef;
45use crate::table_refill::build_hummock_table_refill_runtime_config;
46
47pub struct NotificationServiceImpl {
48 env: MetaSrvEnv,
49
50 metadata_manager: MetadataManager,
51 hummock_manager: HummockManagerRef,
52 backup_manager: BackupManagerRef,
53 serving_vnode_mapping: ServingVnodeMappingRef,
54}
55
56impl NotificationServiceImpl {
57 pub async fn new(
58 env: MetaSrvEnv,
59 metadata_manager: MetadataManager,
60 hummock_manager: HummockManagerRef,
61 backup_manager: BackupManagerRef,
62 serving_vnode_mapping: ServingVnodeMappingRef,
63 ) -> MetaResult<Self> {
64 let service = Self {
65 env,
66 metadata_manager,
67 hummock_manager,
68 backup_manager,
69 serving_vnode_mapping,
70 };
71 let (secrets, _catalog_version) = service.get_decrypted_secret_snapshot().await?;
72 LocalSecretManager::global().init_secrets(secrets);
73 Ok(service)
74 }
75
76 async fn get_catalog_snapshot(
77 &self,
78 ) -> MetaResult<(Catalog, Vec<UserInfo>, NotificationVersion)> {
79 let catalog_guard = self
80 .metadata_manager
81 .catalog_controller
82 .get_inner_read_guard()
83 .await;
84 let (
85 (
86 databases,
87 schemas,
88 tables,
89 sources,
90 sinks,
91 subscriptions,
92 indexes,
93 views,
94 functions,
95 connections,
96 secrets,
97 ),
98 users,
99 ) = catalog_guard.snapshot().await?;
100 let notification_version = self.env.notification_manager().current_version().await;
101 Ok((
102 (
103 databases,
104 schemas,
105 tables,
106 sources,
107 sinks,
108 subscriptions,
109 indexes,
110 views,
111 functions,
112 connections,
113 secrets,
114 ),
115 users,
116 notification_version,
117 ))
118 }
119
120 async fn get_decrypted_secret_snapshot(
122 &self,
123 ) -> MetaResult<(Vec<Secret>, NotificationVersion)> {
124 let catalog_guard = self
125 .metadata_manager
126 .catalog_controller
127 .get_inner_read_guard()
128 .await;
129 let secrets = catalog_guard.list_secrets().await?;
130 let notification_version = self.env.notification_manager().current_version().await;
131
132 let decrypted_secrets = self.decrypt_secrets(secrets)?;
133
134 Ok((decrypted_secrets, notification_version))
135 }
136
137 fn decrypt_secrets(&self, secrets: Vec<Secret>) -> MetaResult<Vec<Secret>> {
138 if secrets.is_empty() {
140 return Ok(vec![]);
141 }
142 let secret_store_private_key = self
143 .env
144 .opts
145 .secret_store_private_key
146 .clone()
147 .ok_or_else(|| anyhow!("secret_store_private_key is not configured"))?;
148 let mut decrypted_secrets = Vec::with_capacity(secrets.len());
149 for mut secret in secrets {
150 let encrypted_secret = SecretEncryption::deserialize(secret.get_value())
151 .context(format!("failed to deserialize secret {}", secret.name))?;
152 let decrypted_secret = encrypted_secret
153 .decrypt(secret_store_private_key.as_slice())
154 .context(format!("failed to decrypt secret {}", secret.name))?;
155 secret.value = decrypted_secret;
156 decrypted_secrets.push(secret);
157 }
158 Ok(decrypted_secrets)
159 }
160
161 async fn get_worker_slot_mapping_snapshot(
162 &self,
163 ) -> MetaResult<(Vec<FragmentWorkerSlotMapping>, NotificationVersion)> {
164 let mappings = self
165 .metadata_manager
166 .catalog_controller
167 .get_worker_slot_mappings();
168
169 let notification_version = self.env.notification_manager().current_version().await;
170 Ok((mappings, notification_version))
171 }
172
173 fn get_serving_vnode_mappings(&self) -> Vec<FragmentWorkerSlotMapping> {
174 self.serving_vnode_mapping
175 .all()
176 .iter()
177 .map(|(&fragment_id, mapping)| FragmentWorkerSlotMapping {
178 fragment_id,
179 mapping: Some(mapping.to_protobuf()),
180 })
181 .collect()
182 }
183
184 async fn get_worker_node_snapshot(&self) -> MetaResult<(Vec<WorkerNode>, NotificationVersion)> {
185 let cluster_guard = self
186 .metadata_manager
187 .cluster_controller
188 .get_inner_read_guard()
189 .await;
190 let compute_nodes = cluster_guard
191 .list_workers(Some(WorkerType::ComputeNode.into()), Some(Running.into()))
192 .await?;
193 let frontends = cluster_guard
194 .list_workers(Some(WorkerType::Frontend.into()), Some(Running.into()))
195 .await?;
196 let worker_nodes = compute_nodes.into_iter().chain(frontends).collect();
197 let notification_version = self.env.notification_manager().current_version().await;
198 Ok((worker_nodes, notification_version))
199 }
200
201 async fn get_tables_snapshot(&self) -> MetaResult<(Vec<Table>, NotificationVersion)> {
202 let catalog_guard = self
203 .metadata_manager
204 .catalog_controller
205 .get_inner_read_guard()
206 .await;
207 let mut tables = catalog_guard.list_all_state_tables().await?;
208 tables.extend(catalog_guard.dropped_tables.values().cloned());
209 let notification_version = self.env.notification_manager().current_version().await;
210 Ok((tables, notification_version))
211 }
212
213 async fn get_hummock_catalog_snapshot(
214 &self,
215 ) -> MetaResult<(
216 Vec<Table>,
217 PbTableCacheRefillPolicies,
218 HashMap<FragmentId, FragmentServingInfo>,
219 NotificationVersion,
220 )> {
221 let catalog_guard = self
222 .metadata_manager
223 .catalog_controller
224 .get_inner_read_guard()
225 .await;
226 let mut tables = catalog_guard.list_all_state_tables().await?;
227 tables.extend(catalog_guard.dropped_tables.values().cloned());
228 let table_cache_refill_policies =
229 catalog_guard.table_cache_refill_policies_snapshot().await?;
230 let fragment_serving_infos = catalog_guard.fragment_serving_infos().await?;
231 let notification_version = self.env.notification_manager().current_version().await;
232 Ok((
233 tables,
234 table_cache_refill_policies,
235 fragment_serving_infos,
236 notification_version,
237 ))
238 }
239
240 async fn get_cluster_resource(&self) -> MetaResult<ClusterResource> {
242 self.metadata_manager
243 .cluster_controller
244 .cluster_resource()
245 .await
246 }
247
248 async fn compactor_subscribe(&self) -> MetaResult<MetaSnapshot> {
249 let (tables, catalog_version) = self.get_tables_snapshot().await?;
250 let cluster_resource = self.get_cluster_resource().await?;
251
252 Ok(MetaSnapshot {
253 tables,
254 version: Some(SnapshotVersion {
255 catalog_version,
256 ..Default::default()
257 }),
258 cluster_resource: Some(cluster_resource),
259 ..Default::default()
260 })
261 }
262
263 async fn frontend_subscribe(&self) -> MetaResult<MetaSnapshot> {
264 let (
265 (
266 databases,
267 schemas,
268 tables,
269 sources,
270 sinks,
271 subscriptions,
272 indexes,
273 views,
274 functions,
275 connections,
276 secrets,
277 ),
278 users,
279 catalog_version,
280 ) = self.get_catalog_snapshot().await?;
281 let object_dependencies = self
282 .metadata_manager
283 .catalog_controller
284 .list_created_object_dependencies()
285 .await?;
286
287 let decrypted_secrets = self.decrypt_secrets(secrets)?;
289
290 let (streaming_worker_slot_mappings, streaming_worker_slot_mapping_version) =
291 self.get_worker_slot_mapping_snapshot().await?;
292
293 let streaming_job_count = self.metadata_manager.count_streaming_job().await?;
294 if streaming_job_count > 0 && streaming_worker_slot_mappings.is_empty() {
295 tracing::warn!(
296 streaming_job_count,
297 "frontend subscribe returns empty streaming_worker_slot_mappings while streaming jobs exist; meta may still be recovering"
298 );
299 }
300
301 let serving_worker_slot_mappings = self.get_serving_vnode_mappings();
302
303 let (nodes, worker_node_version) = self.get_worker_node_snapshot().await?;
304
305 let hummock_version = self
306 .hummock_manager
307 .on_current_version(|version| {
308 FrontendHummockVersion::from_version(version).to_protobuf()
309 })
310 .await;
311
312 let session_params = self
313 .env
314 .session_params_manager_impl_ref()
315 .get_params()
316 .await;
317
318 let session_params = Some(GetSessionParamsResponse {
319 params: serde_json::to_string(&session_params)
320 .context("failed to encode session params")?,
321 });
322
323 let cluster_resource = self.get_cluster_resource().await?;
324
325 Ok(MetaSnapshot {
326 databases,
327 schemas,
328 sources,
329 sinks,
330 tables,
331 indexes,
332 views,
333 subscriptions,
334 functions,
335 connections,
336 secrets: decrypted_secrets,
337 users,
338 nodes,
339 hummock_version: Some(hummock_version),
340 version: Some(SnapshotVersion {
341 catalog_version,
342 worker_node_version,
343 streaming_worker_slot_mapping_version,
344 }),
345 serving_worker_slot_mappings,
346 streaming_worker_slot_mappings,
347 session_params,
348 object_dependencies,
349 cluster_resource: Some(cluster_resource),
350 ..Default::default()
351 })
352 }
353
354 async fn hummock_subscribe(&self, worker_id: WorkerId) -> MetaResult<MetaSnapshot> {
355 let (tables, table_cache_refill_policies, fragment_serving_infos, catalog_version) =
356 self.get_hummock_catalog_snapshot().await?;
357 let hummock_version = self
358 .hummock_manager
359 .on_current_version(|version| version.into())
360 .await;
361 let hummock_write_limits = self.hummock_manager.write_limits().await;
362 let meta_backup_manifest_id = self.backup_manager.manifest().await.manifest_id;
363 let cluster_resource = self.get_cluster_resource().await?;
364 let table_refill_runtime_config = build_hummock_table_refill_runtime_config(
365 &self.serving_vnode_mapping,
366 worker_id,
367 table_cache_refill_policies,
368 &fragment_serving_infos,
369 catalog_version,
370 );
371
372 Ok(MetaSnapshot {
373 tables,
374 hummock_version: Some(hummock_version),
375 version: Some(SnapshotVersion {
376 catalog_version,
377 ..Default::default()
378 }),
379 meta_backup_manifest_id: Some(MetaBackupManifestId {
380 id: meta_backup_manifest_id,
381 }),
382 hummock_write_limits: Some(WriteLimits {
383 write_limits: hummock_write_limits,
384 }),
385 cluster_resource: Some(cluster_resource),
386 table_refill_runtime_config: Some(table_refill_runtime_config),
387 ..Default::default()
388 })
389 }
390
391 async fn compute_subscribe(&self) -> MetaResult<MetaSnapshot> {
392 let (secrets, catalog_version) = self.get_decrypted_secret_snapshot().await?;
393 let cluster_resource = self.get_cluster_resource().await?;
394
395 Ok(MetaSnapshot {
396 secrets,
397 version: Some(SnapshotVersion {
398 catalog_version,
399 ..Default::default()
400 }),
401 cluster_resource: Some(cluster_resource),
402 ..Default::default()
403 })
404 }
405}
406
407#[async_trait::async_trait]
408impl NotificationService for NotificationServiceImpl {
409 type SubscribeStream = UnboundedReceiverStream<Notification>;
410
411 async fn subscribe(
412 &self,
413 request: Request<SubscribeRequest>,
414 ) -> Result<Response<Self::SubscribeStream>, Status> {
415 let req = request.into_inner();
416 let host_address = req.get_host()?.clone();
417 let subscribe_type = req.get_subscribe_type()?;
418
419 let worker_key = WorkerKey(host_address);
420
421 let (tx, rx) = mpsc::unbounded_channel();
422 self.env
423 .notification_manager()
424 .insert_sender(subscribe_type, worker_key, tx.clone());
425
426 let meta_snapshot = match subscribe_type {
427 SubscribeType::Compactor => self.compactor_subscribe().await?,
428 SubscribeType::Frontend => self.frontend_subscribe().await?,
429 SubscribeType::Hummock => {
430 self.hummock_manager
431 .pin_version(req.get_worker_id())
432 .await?;
433 self.hummock_subscribe(req.get_worker_id()).await?
434 }
435 SubscribeType::Compute => self.compute_subscribe().await?,
436 SubscribeType::Unspecified => unreachable!(),
437 };
438
439 self.env
440 .notification_manager()
441 .notify_snapshot(tx, meta_snapshot);
442
443 Ok(Response::new(UnboundedReceiverStream::new(rx)))
444 }
445}