Skip to main content

risingwave_meta_service/
hummock_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, HashSet};
16use std::time::Duration;
17
18use compact_task::PbTaskStatus;
19use futures::StreamExt;
20use itertools::Itertools;
21use risingwave_common::catalog::SYS_CATALOG_START_ID;
22use risingwave_hummock_sdk::key_range::KeyRange;
23use risingwave_hummock_sdk::version::HummockVersionDelta;
24use risingwave_meta::backup_restore::BackupManagerRef;
25use risingwave_meta::manager::MetadataManager;
26use risingwave_meta::manager::iceberg_compaction::IcebergCompactionManagerRef;
27use risingwave_pb::hummock::get_compaction_score_response::PickerInfo;
28use risingwave_pb::hummock::hummock_manager_service_server::HummockManagerService;
29use risingwave_pb::hummock::subscribe_compaction_event_request::Event as RequestEvent;
30use risingwave_pb::hummock::*;
31use risingwave_pb::iceberg_compaction::subscribe_iceberg_compaction_event_request::Event as IcebergRequestEvent;
32use risingwave_pb::iceberg_compaction::{
33    SubscribeIcebergCompactionEventRequest, SubscribeIcebergCompactionEventResponse,
34};
35use tonic::{Request, Response, Status, Streaming};
36
37use crate::RwReceiverStream;
38use crate::hummock::compaction::selector::ManualCompactionOption;
39use crate::hummock::{HummockManagerRef, ManualCompactionTriggerResult};
40
41pub struct HummockServiceImpl {
42    hummock_manager: HummockManagerRef,
43    metadata_manager: MetadataManager,
44    backup_manager: BackupManagerRef,
45    iceberg_compaction_manager: IcebergCompactionManagerRef,
46}
47
48impl HummockServiceImpl {
49    pub fn new(
50        hummock_manager: HummockManagerRef,
51        metadata_manager: MetadataManager,
52        backup_manager: BackupManagerRef,
53        iceberg_compaction_manager: IcebergCompactionManagerRef,
54    ) -> Self {
55        HummockServiceImpl {
56            hummock_manager,
57            metadata_manager,
58            backup_manager,
59            iceberg_compaction_manager,
60        }
61    }
62}
63
64macro_rules! fields_to_kvs {
65    ($struct:ident, $($field:ident),*) => {
66        {
67            let mut kvs = HashMap::default();
68            $(
69                kvs.insert(stringify!($field).to_string(), $struct.$field.to_string());
70            )*
71            kvs
72        }
73    }
74}
75
76#[async_trait::async_trait]
77impl HummockManagerService for HummockServiceImpl {
78    type SubscribeCompactionEventStream = RwReceiverStream<SubscribeCompactionEventResponse>;
79    type SubscribeIcebergCompactionEventStream =
80        RwReceiverStream<SubscribeIcebergCompactionEventResponse>;
81
82    async fn unpin_version_before(
83        &self,
84        request: Request<UnpinVersionBeforeRequest>,
85    ) -> Result<Response<UnpinVersionBeforeResponse>, Status> {
86        let req = request.into_inner();
87        self.hummock_manager
88            .unpin_version_before(req.context_id, req.unpin_version_before)
89            .await?;
90        Ok(Response::new(UnpinVersionBeforeResponse { status: None }))
91    }
92
93    async fn get_current_version(
94        &self,
95        _request: Request<GetCurrentVersionRequest>,
96    ) -> Result<Response<GetCurrentVersionResponse>, Status> {
97        let current_version = self
98            .hummock_manager
99            .on_current_version(|version| version.into())
100            .await;
101        Ok(Response::new(GetCurrentVersionResponse {
102            status: None,
103            current_version: Some(current_version),
104        }))
105    }
106
107    async fn replay_version_delta(
108        &self,
109        request: Request<ReplayVersionDeltaRequest>,
110    ) -> Result<Response<ReplayVersionDeltaResponse>, Status> {
111        #[cfg(any(test, feature = "test"))]
112        {
113            let req = request.into_inner();
114            let (version, compaction_groups) = self
115                .hummock_manager
116                .replay_version_delta(HummockVersionDelta::from_rpc_protobuf(
117                    &req.version_delta.unwrap(),
118                ))
119                .await?;
120            Ok(Response::new(ReplayVersionDeltaResponse {
121                version: Some(version.into()),
122                modified_compaction_groups: compaction_groups,
123            }))
124        }
125        #[cfg(not(any(test, feature = "test")))]
126        {
127            let _ = request;
128            Err(Status::unimplemented(
129                "replay_version_delta is only available in test builds",
130            ))
131        }
132    }
133
134    async fn trigger_compaction_deterministic(
135        &self,
136        request: Request<TriggerCompactionDeterministicRequest>,
137    ) -> Result<Response<TriggerCompactionDeterministicResponse>, Status> {
138        let req = request.into_inner();
139        self.hummock_manager
140            .trigger_compaction_deterministic(req.version_id, req.compaction_groups)
141            .await?;
142        Ok(Response::new(TriggerCompactionDeterministicResponse {}))
143    }
144
145    async fn disable_commit_epoch(
146        &self,
147        _request: Request<DisableCommitEpochRequest>,
148    ) -> Result<Response<DisableCommitEpochResponse>, Status> {
149        let version = self.hummock_manager.disable_commit_epoch().await;
150        Ok(Response::new(DisableCommitEpochResponse {
151            current_version: Some(PbHummockVersion::from(&*version)),
152        }))
153    }
154
155    async fn list_version_deltas(
156        &self,
157        request: Request<ListVersionDeltasRequest>,
158    ) -> Result<Response<ListVersionDeltasResponse>, Status> {
159        let req = request.into_inner();
160        let version_deltas = self
161            .hummock_manager
162            .list_version_deltas(req.start_id, req.num_limit)
163            .await?;
164        let resp = ListVersionDeltasResponse {
165            version_deltas: Some(PbHummockVersionDeltas {
166                version_deltas: version_deltas
167                    .into_iter()
168                    .map(HummockVersionDelta::into)
169                    .collect(),
170            }),
171        };
172        Ok(Response::new(resp))
173    }
174
175    async fn get_new_object_ids(
176        &self,
177        request: Request<GetNewObjectIdsRequest>,
178    ) -> Result<Response<GetNewObjectIdsResponse>, Status> {
179        let object_id_range = self
180            .hummock_manager
181            .get_new_object_ids(request.into_inner().number)
182            .await?;
183        Ok(Response::new(GetNewObjectIdsResponse {
184            status: None,
185            start_id: object_id_range.start_id,
186            end_id: object_id_range.end_id,
187        }))
188    }
189
190    async fn trigger_manual_compaction(
191        &self,
192        request: Request<TriggerManualCompactionRequest>,
193    ) -> Result<Response<TriggerManualCompactionResponse>, Status> {
194        let request = request.into_inner();
195        let compaction_group_id = request.compaction_group_id;
196        let mut option = ManualCompactionOption {
197            level: request.level as usize,
198            target_level: request
199                .target_level
200                .map(|target_level| target_level as usize),
201            sst_ids: request.sst_ids,
202            exclusive: request.exclusive.unwrap_or(false),
203            ..Default::default()
204        };
205
206        // rewrite the key_range
207        match request.key_range {
208            Some(pb_key_range) => {
209                option.key_range = KeyRange {
210                    left: pb_key_range.left.into(),
211                    right: pb_key_range.right.into(),
212                    right_exclusive: pb_key_range.right_exclusive,
213                };
214            }
215
216            None => {
217                option.key_range = KeyRange::default();
218            }
219        }
220
221        // get internal_table_id by metadata_manger
222        if request.table_id.as_raw_id() < SYS_CATALOG_START_ID as u32 {
223            // We need to make sure to use the correct table_id to filter sst
224            let job_id = request.table_id;
225            if let Ok(table_fragment) = self.metadata_manager.get_job_fragments_by_id(job_id).await
226            {
227                option.internal_table_id = HashSet::from_iter(table_fragment.all_table_ids());
228            }
229        }
230
231        assert!(
232            option
233                .internal_table_id
234                .iter()
235                .all(|table_id| table_id.as_raw_id() < SYS_CATALOG_START_ID as u32),
236        );
237
238        tracing::info!(
239            "Try trigger_manual_compaction compaction_group_id {} option {:?}",
240            compaction_group_id,
241            &option
242        );
243
244        let should_retry = match self
245            .hummock_manager
246            .trigger_manual_compaction(compaction_group_id, option)
247            .await?
248        {
249            ManualCompactionTriggerResult::Submitted => false,
250            ManualCompactionTriggerResult::Retry => true,
251        };
252
253        Ok(Response::new(TriggerManualCompactionResponse {
254            status: None,
255            should_retry: Some(should_retry),
256        }))
257    }
258
259    async fn trigger_full_gc(
260        &self,
261        request: Request<TriggerFullGcRequest>,
262    ) -> Result<Response<TriggerFullGcResponse>, Status> {
263        let req = request.into_inner();
264        let backup_manager_2 = self.backup_manager.clone();
265        let hummock_manager_2 = self.hummock_manager.clone();
266        tokio::task::spawn(async move {
267            use thiserror_ext::AsReport;
268            let _ = hummock_manager_2
269                .start_full_gc(
270                    Duration::from_secs(req.sst_retention_time_sec),
271                    req.prefix,
272                    Some(backup_manager_2),
273                )
274                .await
275                .inspect_err(|e| tracing::warn!(error = %e.as_report(), "Failed to start GC."));
276        });
277        Ok(Response::new(TriggerFullGcResponse { status: None }))
278    }
279
280    async fn rise_ctl_get_pinned_versions_summary(
281        &self,
282        _request: Request<RiseCtlGetPinnedVersionsSummaryRequest>,
283    ) -> Result<Response<RiseCtlGetPinnedVersionsSummaryResponse>, Status> {
284        let pinned_versions = self.hummock_manager.list_pinned_version().await;
285        let workers = self
286            .hummock_manager
287            .list_workers(&pinned_versions.iter().map(|v| v.context_id).collect_vec())
288            .await?;
289        Ok(Response::new(RiseCtlGetPinnedVersionsSummaryResponse {
290            summary: Some(PinnedVersionsSummary {
291                pinned_versions,
292                workers,
293            }),
294        }))
295    }
296
297    async fn get_assigned_compact_task_num(
298        &self,
299        _request: Request<GetAssignedCompactTaskNumRequest>,
300    ) -> Result<Response<GetAssignedCompactTaskNumResponse>, Status> {
301        let num_tasks = self.hummock_manager.get_assigned_compact_task_num().await;
302        Ok(Response::new(GetAssignedCompactTaskNumResponse {
303            num_tasks: num_tasks as u32,
304        }))
305    }
306
307    async fn rise_ctl_list_compaction_group(
308        &self,
309        _request: Request<RiseCtlListCompactionGroupRequest>,
310    ) -> Result<Response<RiseCtlListCompactionGroupResponse>, Status> {
311        let compaction_groups = self.hummock_manager.list_compaction_group().await;
312        Ok(Response::new(RiseCtlListCompactionGroupResponse {
313            status: None,
314            compaction_groups,
315        }))
316    }
317
318    async fn rise_ctl_update_compaction_config(
319        &self,
320        request: Request<RiseCtlUpdateCompactionConfigRequest>,
321    ) -> Result<Response<RiseCtlUpdateCompactionConfigResponse>, Status> {
322        let RiseCtlUpdateCompactionConfigRequest {
323            compaction_group_ids,
324            configs,
325        } = request.into_inner();
326        self.hummock_manager
327            .update_compaction_config(
328                compaction_group_ids.as_slice(),
329                configs
330                    .into_iter()
331                    .map(|c| c.mutable_config.unwrap())
332                    .collect::<Vec<_>>()
333                    .as_slice(),
334            )
335            .await?;
336        Ok(Response::new(RiseCtlUpdateCompactionConfigResponse {
337            status: None,
338        }))
339    }
340
341    async fn init_metadata_for_replay(
342        &self,
343        request: Request<InitMetadataForReplayRequest>,
344    ) -> Result<Response<InitMetadataForReplayResponse>, Status> {
345        let InitMetadataForReplayRequest {
346            tables,
347            compaction_groups,
348        } = request.into_inner();
349
350        self.hummock_manager
351            .init_metadata_for_version_replay(tables, compaction_groups)?;
352        Ok(Response::new(InitMetadataForReplayResponse {}))
353    }
354
355    async fn pin_version(
356        &self,
357        request: Request<PinVersionRequest>,
358    ) -> Result<Response<PinVersionResponse>, Status> {
359        let req = request.into_inner();
360        let version = self.hummock_manager.pin_version(req.context_id).await?;
361        Ok(Response::new(PinVersionResponse {
362            pinned_version: Some(PbHummockVersion::from(&*version)),
363        }))
364    }
365
366    async fn split_compaction_group(
367        &self,
368        request: Request<SplitCompactionGroupRequest>,
369    ) -> Result<Response<SplitCompactionGroupResponse>, Status> {
370        let req = request.into_inner();
371        let new_group_id = self
372            .hummock_manager
373            .move_state_tables_to_dedicated_compaction_group(
374                req.group_id,
375                &req.table_ids,
376                if req.partition_vnode_count > 0 {
377                    Some(req.partition_vnode_count)
378                } else {
379                    None
380                },
381            )
382            .await?
383            .0;
384        Ok(Response::new(SplitCompactionGroupResponse { new_group_id }))
385    }
386
387    async fn rise_ctl_pause_version_checkpoint(
388        &self,
389        _request: Request<RiseCtlPauseVersionCheckpointRequest>,
390    ) -> Result<Response<RiseCtlPauseVersionCheckpointResponse>, Status> {
391        self.hummock_manager.pause_version_checkpoint();
392        Ok(Response::new(RiseCtlPauseVersionCheckpointResponse {}))
393    }
394
395    async fn rise_ctl_resume_version_checkpoint(
396        &self,
397        _request: Request<RiseCtlResumeVersionCheckpointRequest>,
398    ) -> Result<Response<RiseCtlResumeVersionCheckpointResponse>, Status> {
399        self.hummock_manager.resume_version_checkpoint();
400        Ok(Response::new(RiseCtlResumeVersionCheckpointResponse {}))
401    }
402
403    async fn rise_ctl_get_checkpoint_version(
404        &self,
405        _request: Request<RiseCtlGetCheckpointVersionRequest>,
406    ) -> Result<Response<RiseCtlGetCheckpointVersionResponse>, Status> {
407        let checkpoint_version = self.hummock_manager.get_checkpoint_version().await;
408        Ok(Response::new(RiseCtlGetCheckpointVersionResponse {
409            checkpoint_version: Some(PbHummockVersion::from(&*checkpoint_version)),
410        }))
411    }
412
413    async fn rise_ctl_list_compaction_status(
414        &self,
415        _request: Request<RiseCtlListCompactionStatusRequest>,
416    ) -> Result<Response<RiseCtlListCompactionStatusResponse>, Status> {
417        let (compaction_statuses, task_assignment) =
418            self.hummock_manager.list_compaction_status().await;
419        let task_progress = self.hummock_manager.compactor_manager.get_progress();
420        Ok(Response::new(RiseCtlListCompactionStatusResponse {
421            compaction_statuses,
422            task_assignment,
423            task_progress,
424        }))
425    }
426
427    async fn subscribe_compaction_event(
428        &self,
429        request: Request<Streaming<SubscribeCompactionEventRequest>>,
430    ) -> Result<Response<Self::SubscribeCompactionEventStream>, tonic::Status> {
431        let mut request_stream: Streaming<SubscribeCompactionEventRequest> = request.into_inner();
432        let register_req = {
433            let req = request_stream.next().await.ok_or_else(|| {
434                Status::invalid_argument("subscribe_compaction_event request is empty")
435            })??;
436
437            match req.event {
438                Some(RequestEvent::Register(register)) => register,
439                _ => {
440                    return Err(Status::invalid_argument(
441                        "the first message must be `Register`",
442                    ));
443                }
444            }
445        };
446
447        let context_id = register_req.context_id;
448
449        // check_context and add_compactor as a whole is not atomic, but compactor_manager will
450        // remove invalid compactor eventually.
451        if !self.hummock_manager.check_context(context_id).await? {
452            return Err(Status::new(
453                tonic::Code::Internal,
454                format!("invalid hummock context {}", context_id),
455            ));
456        }
457        let compactor_manager = self.hummock_manager.compactor_manager.clone();
458
459        let rx: tokio::sync::mpsc::UnboundedReceiver<
460            Result<SubscribeCompactionEventResponse, crate::MetaError>,
461        > = compactor_manager.add_compactor(context_id);
462
463        // register request stream to hummock
464        self.hummock_manager
465            .add_compactor_stream(context_id, request_stream);
466
467        // Trigger compaction on all compaction groups.
468        for cg_id in self.hummock_manager.compaction_group_ids().await {
469            self.hummock_manager
470                .try_send_compaction_request(cg_id, compact_task::TaskType::Dynamic);
471        }
472
473        Ok(Response::new(RwReceiverStream::new(rx)))
474    }
475
476    async fn subscribe_iceberg_compaction_event(
477        &self,
478        request: Request<Streaming<SubscribeIcebergCompactionEventRequest>>,
479    ) -> Result<Response<Self::SubscribeIcebergCompactionEventStream>, tonic::Status> {
480        let mut request_stream: Streaming<SubscribeIcebergCompactionEventRequest> =
481            request.into_inner();
482        let register_req = {
483            let req = request_stream.next().await.ok_or_else(|| {
484                Status::invalid_argument("subscribe_compaction_event request is empty")
485            })??;
486
487            match req.event {
488                Some(IcebergRequestEvent::Register(register)) => register,
489                _ => {
490                    return Err(Status::invalid_argument(
491                        "the first message must be `Register`",
492                    ));
493                }
494            }
495        };
496
497        let context_id = register_req.context_id;
498
499        // check_context and add_compactor as a whole is not atomic, but compactor_manager will
500        // remove invalid compactor eventually.
501        if !self.hummock_manager.check_context(context_id).await? {
502            return Err(Status::new(
503                tonic::Code::Internal,
504                format!("invalid hummock context {}", context_id),
505            ));
506        }
507
508        let rx: tokio::sync::mpsc::UnboundedReceiver<
509            Result<SubscribeIcebergCompactionEventResponse, crate::MetaError>,
510        > = self
511            .iceberg_compaction_manager
512            .iceberg_compactor_manager
513            .add_compactor(context_id);
514
515        self.iceberg_compaction_manager
516            .add_compactor_stream(context_id, request_stream);
517
518        // TODO: Trigger iceberg compaction
519
520        Ok(Response::new(RwReceiverStream::new(rx)))
521    }
522
523    async fn report_compaction_task(
524        &self,
525        _request: Request<ReportCompactionTaskRequest>,
526    ) -> Result<Response<ReportCompactionTaskResponse>, Status> {
527        unreachable!()
528    }
529
530    async fn list_branched_object(
531        &self,
532        _request: Request<ListBranchedObjectRequest>,
533    ) -> Result<Response<ListBranchedObjectResponse>, Status> {
534        let branched_objects = self
535            .hummock_manager
536            .list_branched_objects()
537            .await
538            .into_iter()
539            .flat_map(|(object_id, v)| {
540                v.into_iter()
541                    .map(move |(compaction_group_id, sst_ids)| BranchedObject {
542                        object_id,
543                        sst_id: sst_ids,
544                        compaction_group_id,
545                    })
546            })
547            .collect();
548        Ok(Response::new(ListBranchedObjectResponse {
549            branched_objects,
550        }))
551    }
552
553    async fn list_active_write_limit(
554        &self,
555        _request: Request<ListActiveWriteLimitRequest>,
556    ) -> Result<Response<ListActiveWriteLimitResponse>, Status> {
557        Ok(Response::new(ListActiveWriteLimitResponse {
558            write_limits: self.hummock_manager.write_limits().await,
559        }))
560    }
561
562    async fn list_hummock_meta_config(
563        &self,
564        _request: Request<ListHummockMetaConfigRequest>,
565    ) -> Result<Response<ListHummockMetaConfigResponse>, Status> {
566        let opt = &self.hummock_manager.env.opts;
567        let configs = fields_to_kvs!(
568            opt,
569            vacuum_interval_sec,
570            vacuum_spin_interval_ms,
571            hummock_version_checkpoint_interval_sec,
572            min_delta_log_num_for_hummock_version_checkpoint,
573            min_sst_retention_time_sec,
574            full_gc_interval_sec,
575            periodic_compaction_interval_sec,
576            periodic_space_reclaim_compaction_interval_sec,
577            periodic_ttl_reclaim_compaction_interval_sec,
578            periodic_tombstone_reclaim_compaction_interval_sec,
579            periodic_scheduling_compaction_group_split_interval_sec,
580            enable_compaction_group_normalize,
581            max_normalize_splits_per_round,
582            do_not_config_object_storage_lifecycle,
583            partition_vnode_count,
584            table_high_write_throughput_threshold,
585            table_low_write_throughput_threshold,
586            compaction_task_max_heartbeat_interval_secs,
587            periodic_scheduling_compaction_group_merge_interval_sec
588        );
589        Ok(Response::new(ListHummockMetaConfigResponse { configs }))
590    }
591
592    async fn rise_ctl_rebuild_table_stats(
593        &self,
594        _request: Request<RiseCtlRebuildTableStatsRequest>,
595    ) -> Result<Response<RiseCtlRebuildTableStatsResponse>, Status> {
596        self.hummock_manager.rebuild_table_stats().await?;
597        Ok(Response::new(RiseCtlRebuildTableStatsResponse {}))
598    }
599
600    async fn get_compaction_score(
601        &self,
602        request: Request<GetCompactionScoreRequest>,
603    ) -> Result<Response<GetCompactionScoreResponse>, Status> {
604        let compaction_group_id = request.into_inner().compaction_group_id;
605        let scores = self
606            .hummock_manager
607            .get_compaction_scores(compaction_group_id)
608            .await
609            .into_iter()
610            .map(|s| PickerInfo {
611                score: s.score,
612                select_level: s.select_level as _,
613                target_level: s.target_level as _,
614                picker_type: s.picker_type.to_string(),
615            })
616            .collect();
617        Ok(Response::new(GetCompactionScoreResponse {
618            compaction_group_id,
619            scores,
620        }))
621    }
622
623    async fn list_compact_task_assignment(
624        &self,
625        _request: Request<ListCompactTaskAssignmentRequest>,
626    ) -> Result<Response<ListCompactTaskAssignmentResponse>, Status> {
627        let (_compaction_statuses, task_assignment) =
628            self.hummock_manager.list_compaction_status().await;
629        Ok(Response::new(ListCompactTaskAssignmentResponse {
630            task_assignment,
631        }))
632    }
633
634    async fn list_compact_task_progress(
635        &self,
636        _request: Request<ListCompactTaskProgressRequest>,
637    ) -> Result<Response<ListCompactTaskProgressResponse>, Status> {
638        let task_progress = self.hummock_manager.compactor_manager.get_progress();
639
640        Ok(Response::new(ListCompactTaskProgressResponse {
641            task_progress,
642        }))
643    }
644
645    async fn cancel_compact_task(
646        &self,
647        request: Request<CancelCompactTaskRequest>,
648    ) -> Result<Response<CancelCompactTaskResponse>, Status> {
649        let request = request.into_inner();
650        let ret = self
651            .hummock_manager
652            .cancel_compact_task(
653                request.task_id,
654                PbTaskStatus::try_from(request.task_status).unwrap(),
655            )
656            .await?;
657
658        let response = Response::new(CancelCompactTaskResponse { ret });
659        return Ok(response);
660    }
661
662    async fn get_version_by_epoch(
663        &self,
664        request: Request<GetVersionByEpochRequest>,
665    ) -> Result<Response<GetVersionByEpochResponse>, Status> {
666        let GetVersionByEpochRequest { epoch, table_id } = request.into_inner();
667        let version = self
668            .hummock_manager
669            .epoch_to_version(epoch, table_id)
670            .await?;
671        Ok(Response::new(GetVersionByEpochResponse {
672            version: Some(version.to_protobuf()),
673        }))
674    }
675
676    async fn merge_compaction_group(
677        &self,
678        request: Request<MergeCompactionGroupRequest>,
679    ) -> Result<Response<MergeCompactionGroupResponse>, Status> {
680        let req = request.into_inner();
681        self.hummock_manager
682            .merge_compaction_group(req.left_group_id, req.right_group_id)
683            .await?;
684        Ok(Response::new(MergeCompactionGroupResponse {}))
685    }
686
687    async fn get_table_change_logs(
688        &self,
689        request: Request<GetTableChangeLogsRequest>,
690    ) -> Result<Response<GetTableChangeLogsResponse>, Status> {
691        let GetTableChangeLogsRequest {
692            epoch_only,
693            start_epoch_inclusive,
694            end_epoch_inclusive,
695            table_ids,
696            exclude_empty,
697            limit,
698        } = request.into_inner();
699        let table_change_logs = self
700            .hummock_manager
701            .get_table_change_logs(
702                epoch_only,
703                start_epoch_inclusive,
704                end_epoch_inclusive,
705                table_ids
706                    .map(|t| t.table_ids.into_iter().collect::<HashSet<_>>())
707                    .clone(),
708                exclude_empty,
709                limit,
710            )
711            .await?
712            .into_iter()
713            .map(|(i, l)| (i.as_raw_id(), l.to_protobuf()))
714            .collect();
715        Ok(Response::new(GetTableChangeLogsResponse {
716            table_change_logs,
717        }))
718    }
719}
720
721#[cfg(test)]
722mod tests {
723    use std::collections::HashMap;
724    #[test]
725    fn test_fields_to_kvs() {
726        struct S {
727            foo: u64,
728            bar: String,
729        }
730        let s = S {
731            foo: 15,
732            bar: "foobar".to_owned(),
733        };
734        let kvs: HashMap<String, String> = fields_to_kvs!(s, foo, bar);
735        assert_eq!(kvs.len(), 2);
736        assert_eq!(kvs.get("foo").unwrap(), "15");
737        assert_eq!(kvs.get("bar").unwrap(), "foobar");
738    }
739}