Skip to main content

risingwave_meta_service/
notification_service.rs

1// Copyright 2023 RisingWave Labs
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use 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    /// Get decrypted secret snapshot
121    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        // Skip getting `secret_store_private_key` if there is no secret
139        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    /// Get the total resource of the cluster.
241    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        // Use the plain text secret value for frontend. The secret value will be masked in frontend handle.
288        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}