Skip to main content

meta_srv/handler/
collect_leader_region_handler.rs

1// Copyright 2023 Greptime Team
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 api::v1::meta::{HeartbeatRequest, Role};
16use common_meta::region_registry::{LeaderRegion, LeaderRegionManifestInfo};
17use store_api::region_engine::RegionRole;
18
19use crate::error::Result;
20use crate::handler::{HandleControl, HeartbeatAccumulator, HeartbeatHandler};
21use crate::metasrv::Context;
22
23pub struct CollectLeaderRegionHandler;
24
25#[async_trait::async_trait]
26impl HeartbeatHandler for CollectLeaderRegionHandler {
27    fn is_acceptable(&self, role: Role) -> bool {
28        role == Role::Datanode
29    }
30
31    async fn handle(
32        &self,
33        _req: &HeartbeatRequest,
34        ctx: &mut Context,
35        acc: &mut HeartbeatAccumulator,
36    ) -> Result<HandleControl> {
37        let Some(current_stat) = acc.stat.as_ref() else {
38            return Ok(HandleControl::Continue);
39        };
40
41        let mut key_values = Vec::with_capacity(current_stat.region_stats.len());
42        for stat in current_stat.region_stats.iter() {
43            if !matches!(stat.role, RegionRole::Leader | RegionRole::StagingLeader) {
44                continue;
45            }
46
47            let manifest = LeaderRegionManifestInfo::from_region_stat(stat);
48            let value = LeaderRegion {
49                datanode_id: current_stat.id,
50                manifest,
51            };
52            key_values.push((stat.id, value));
53        }
54        ctx.leader_region_registry.batch_put(key_values);
55
56        Ok(HandleControl::Continue)
57    }
58}
59
60#[cfg(test)]
61mod tests {
62    use common_meta::datanode::{RegionManifestInfo, RegionStat, Stat};
63    use common_meta::region_registry::LeaderRegionManifestInfo;
64    use store_api::region_engine::RegionRole;
65    use store_api::storage::RegionId;
66
67    use super::*;
68    use crate::handler::test_utils::TestEnv;
69
70    fn new_region_stat(id: RegionId, manifest_version: u64, role: RegionRole) -> RegionStat {
71        RegionStat {
72            id,
73            region_manifest: RegionManifestInfo::Mito {
74                manifest_version,
75                flushed_entry_id: 0,
76                file_removed_cnt: 0,
77            },
78            rcus: 0,
79            wcus: 0,
80            approximate_bytes: 0,
81            engine: "mito".to_string(),
82            role,
83            num_rows: 0,
84            memtable_size: 0,
85            manifest_size: 0,
86            sst_size: 0,
87            sst_num: 0,
88            index_size: 0,
89            data_topic_latest_entry_id: 0,
90            metadata_topic_latest_entry_id: 0,
91            min_timestamp: None,
92            max_timestamp: None,
93            written_bytes: 0,
94            query_cpu_time: 0,
95            query_scanned_bytes: 0,
96        }
97    }
98
99    #[tokio::test]
100    async fn test_handle_collect_leader_region() {
101        let env = TestEnv::new();
102        let mut ctx = env.ctx();
103
104        let mut acc = HeartbeatAccumulator {
105            stat: Some(Stat {
106                id: 1,
107                region_stats: vec![
108                    new_region_stat(RegionId::new(1, 1), 1, RegionRole::Leader),
109                    new_region_stat(RegionId::new(1, 2), 2, RegionRole::Follower),
110                ],
111                addr: "127.0.0.1:0000".to_string(),
112                region_num: 2,
113                ..Default::default()
114            }),
115            ..Default::default()
116        };
117
118        let handler = CollectLeaderRegionHandler;
119        let control = handler
120            .handle(&HeartbeatRequest::default(), &mut ctx, &mut acc)
121            .await
122            .unwrap();
123
124        assert_eq!(control, HandleControl::Continue);
125        let regions = ctx
126            .leader_region_registry
127            .batch_get(vec![RegionId::new(1, 1), RegionId::new(1, 2)].into_iter());
128        assert_eq!(regions.len(), 1);
129        assert_eq!(
130            regions.get(&RegionId::new(1, 1)),
131            Some(&LeaderRegion {
132                datanode_id: 1,
133                manifest: LeaderRegionManifestInfo::Mito {
134                    manifest_version: 1,
135                    flushed_entry_id: 0,
136                    topic_latest_entry_id: 0,
137                },
138            })
139        );
140
141        // New heartbeat with new manifest version
142        acc.stat = Some(Stat {
143            id: 1,
144            region_stats: vec![new_region_stat(RegionId::new(1, 1), 2, RegionRole::Leader)],
145            timestamp_millis: 0,
146            addr: "127.0.0.1:0000".to_string(),
147            region_num: 1,
148            node_epoch: 0,
149            ..Default::default()
150        });
151        let control = handler
152            .handle(&HeartbeatRequest::default(), &mut ctx, &mut acc)
153            .await
154            .unwrap();
155
156        assert_eq!(control, HandleControl::Continue);
157        let regions = ctx
158            .leader_region_registry
159            .batch_get(vec![RegionId::new(1, 1)].into_iter());
160        assert_eq!(regions.len(), 1);
161        assert_eq!(
162            regions.get(&RegionId::new(1, 1)),
163            Some(&LeaderRegion {
164                datanode_id: 1,
165                manifest: LeaderRegionManifestInfo::Mito {
166                    manifest_version: 2,
167                    flushed_entry_id: 0,
168                    topic_latest_entry_id: 0,
169                },
170            })
171        );
172
173        // New heartbeat with old manifest version
174        acc.stat = Some(Stat {
175            id: 1,
176            region_stats: vec![new_region_stat(RegionId::new(1, 1), 1, RegionRole::Leader)],
177            timestamp_millis: 0,
178            addr: "127.0.0.1:0000".to_string(),
179            region_num: 1,
180            node_epoch: 0,
181            ..Default::default()
182        });
183        let control = handler
184            .handle(&HeartbeatRequest::default(), &mut ctx, &mut acc)
185            .await
186            .unwrap();
187
188        assert_eq!(control, HandleControl::Continue);
189        let regions = ctx
190            .leader_region_registry
191            .batch_get(vec![RegionId::new(1, 1)].into_iter());
192        assert_eq!(regions.len(), 1);
193        assert_eq!(
194            regions.get(&RegionId::new(1, 1)),
195            // The manifest version is not updated
196            Some(&LeaderRegion {
197                datanode_id: 1,
198                manifest: LeaderRegionManifestInfo::Mito {
199                    manifest_version: 2,
200                    flushed_entry_id: 0,
201                    topic_latest_entry_id: 0,
202                },
203            })
204        );
205    }
206}