Skip to main content

meta_srv/handler/
collect_stats_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 std::cmp::Ordering;
16use std::sync::Arc;
17
18use api::v1::meta::{HeartbeatRequest, Role};
19use common_meta::datanode::{DatanodeStatKey, DatanodeStatValue, Stat};
20use common_meta::instruction::CacheIdent;
21use common_meta::key::node_address::{NodeAddressKey, NodeAddressValue};
22use common_meta::key::{MetadataKey, MetadataValue};
23use common_meta::peer::Peer;
24use common_meta::rpc::store::PutRequest;
25use common_telemetry::{error, info, warn};
26use dashmap::DashMap;
27use snafu::ResultExt;
28use tokio::sync::Mutex;
29
30use crate::error::{self, Result};
31use crate::handler::{HandleControl, HeartbeatAccumulator, HeartbeatHandler};
32use crate::metasrv::Context;
33
34#[derive(Debug, Default)]
35struct EpochStats {
36    stats: Vec<Stat>,
37    epoch: Option<u64>,
38}
39
40impl EpochStats {
41    #[inline]
42    fn drain_all(&mut self) -> Vec<Stat> {
43        self.stats.drain(..).collect()
44    }
45
46    #[inline]
47    fn clear_stats(&mut self) {
48        self.stats.clear();
49    }
50
51    #[inline]
52    fn push_stat(&mut self, stat: Stat) {
53        self.stats.push(stat);
54    }
55
56    #[inline]
57    fn len(&self) -> usize {
58        self.stats.len()
59    }
60
61    #[inline]
62    fn epoch(&self) -> Option<u64> {
63        self.epoch
64    }
65
66    #[inline]
67    fn set_epoch(&mut self, epoch: u64) {
68        self.epoch = Some(epoch);
69    }
70}
71
72const DEFAULT_FLUSH_STATS_FACTOR: usize = 3;
73
74pub struct CollectStatsHandler {
75    stats_cache: DashMap<DatanodeStatKey, Arc<Mutex<EpochStats>>>,
76    flush_stats_factor: usize,
77}
78
79impl Default for CollectStatsHandler {
80    fn default() -> Self {
81        Self::new(None)
82    }
83}
84
85impl CollectStatsHandler {
86    pub fn new(flush_stats_factor: Option<usize>) -> Self {
87        Self {
88            flush_stats_factor: flush_stats_factor.unwrap_or(DEFAULT_FLUSH_STATS_FACTOR),
89            stats_cache: DashMap::default(),
90        }
91    }
92}
93
94#[async_trait::async_trait]
95impl HeartbeatHandler for CollectStatsHandler {
96    fn is_acceptable(&self, role: Role) -> bool {
97        role == Role::Datanode
98    }
99
100    async fn handle(
101        &self,
102        _req: &HeartbeatRequest,
103        ctx: &mut Context,
104        acc: &mut HeartbeatAccumulator,
105    ) -> Result<HandleControl> {
106        let Some(current_stat) = acc.stat.take() else {
107            return Ok(HandleControl::Continue);
108        };
109
110        let key = current_stat.stat_key();
111        let state = {
112            let entry = self
113                .stats_cache
114                .entry(key)
115                .or_insert_with(|| Arc::new(Mutex::new(EpochStats::default())));
116            Arc::clone(entry.value())
117        };
118        let mut epoch_stats = state.lock().await;
119
120        let key: Vec<u8> = key.into();
121
122        let refresh = if let Some(epoch) = epoch_stats.epoch() {
123            match current_stat.node_epoch.cmp(&epoch) {
124                Ordering::Greater => {
125                    // This node may have been redeployed.
126                    epoch_stats.clear_stats();
127                    epoch_stats.set_epoch(current_stat.node_epoch);
128                    epoch_stats.push_stat(current_stat);
129                    true
130                }
131                Ordering::Equal => {
132                    epoch_stats.push_stat(current_stat);
133                    false
134                }
135                Ordering::Less => {
136                    warn!("Ignore stale heartbeat: {:?}", current_stat);
137                    false
138                }
139            }
140        } else {
141            epoch_stats.set_epoch(current_stat.node_epoch);
142            epoch_stats.push_stat(current_stat);
143            // If the epoch is empty, it indicates that the current node sending the heartbeat
144            // for the first time to the current meta leader, so it is necessary to save
145            // the data to the KV store as soon as possible.
146            true
147        };
148
149        // Need to refresh the [datanode -> address] mapping
150        if refresh {
151            // Safety: `epoch_stats.stats` is not empty
152            let last = epoch_stats.stats.last().unwrap();
153            rewrite_node_address(ctx, last).await;
154        }
155
156        if !refresh && epoch_stats.len() < self.flush_stats_factor {
157            return Ok(HandleControl::Continue);
158        }
159
160        let value: Vec<u8> = DatanodeStatValue {
161            stats: epoch_stats.drain_all(),
162        }
163        .try_into()
164        .context(error::InvalidDatanodeStatFormatSnafu {})?;
165        let put = PutRequest {
166            key,
167            value,
168            prev_kv: false,
169        };
170
171        let _ = ctx
172            .in_memory
173            .put(put)
174            .await
175            .context(error::KvBackendSnafu)?;
176
177        Ok(HandleControl::Continue)
178    }
179}
180
181async fn rewrite_node_address(ctx: &Context, stat: &Stat) {
182    let peer = Peer {
183        id: stat.id,
184        addr: stat.addr.clone(),
185    };
186    let key = NodeAddressKey::with_datanode(peer.id).to_bytes();
187    if let Ok(value) = NodeAddressValue::new(peer.clone()).try_as_raw_value() {
188        let put = PutRequest {
189            key,
190            value,
191            prev_kv: false,
192        };
193
194        match ctx.leader_cached_kv_backend.put(put).await {
195            Ok(_) => {
196                info!(
197                    "Successfully updated datanode `NodeAddressValue`: {:?}",
198                    peer
199                );
200                // broadcast invalidating cache
201                let cache_idents = stat
202                    .table_ids()
203                    .into_iter()
204                    .map(CacheIdent::TableId)
205                    .collect::<Vec<_>>();
206                if let Err(e) = ctx
207                    .cache_invalidator
208                    .invalidate(&Default::default(), &cache_idents)
209                    .await
210                {
211                    error!(e; "Failed to invalidate {} `NodeAddressKey` cache, peer: {:?}", cache_idents.len(), peer);
212                }
213            }
214            Err(e) => {
215                error!(e; "Failed to update datanode `NodeAddressValue`: {:?}", peer);
216            }
217        }
218    } else {
219        warn!(
220            "Failed to serialize datanode `NodeAddressValue`: {:?}",
221            peer
222        );
223    }
224}
225
226#[cfg(test)]
227mod tests {
228    use std::any::Any;
229    use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering as AtomicOrdering};
230    use std::sync::{Arc, Mutex as StdMutex, mpsc};
231    use std::thread;
232    use std::time::Duration;
233
234    use common_meta::datanode::DatanodeStatKey;
235    use common_meta::error::{Error as MetaError, Result as MetaResult};
236    use common_meta::kv_backend::{KvBackend, KvBackendRef, ResettableKvBackend, TxnService};
237    use common_meta::rpc::store::{
238        BatchDeleteRequest, BatchDeleteResponse, BatchGetRequest, BatchGetResponse,
239        BatchPutRequest, BatchPutResponse, DeleteRangeRequest, DeleteRangeResponse, PutResponse,
240        RangeRequest, RangeResponse,
241    };
242    use tokio::sync::Semaphore;
243    use tokio::time::{sleep, timeout};
244
245    use super::*;
246    use crate::handler::test_utils::TestEnv;
247    use crate::service::store::cached_kv::LeaderCachedKvBackend;
248
249    struct ControlledKvBackend {
250        recorded_puts: StdMutex<Vec<PutRequest>>,
251        put_entered: Semaphore,
252        put_release: Semaphore,
253        block_next_put: AtomicBool,
254        delay_puts: AtomicBool,
255        active_puts: AtomicUsize,
256        max_active_puts: AtomicUsize,
257    }
258
259    impl ControlledKvBackend {
260        fn new() -> Self {
261            Self {
262                recorded_puts: StdMutex::new(Vec::new()),
263                put_entered: Semaphore::new(0),
264                put_release: Semaphore::new(0),
265                block_next_put: AtomicBool::new(false),
266                delay_puts: AtomicBool::new(false),
267                active_puts: AtomicUsize::new(0),
268                max_active_puts: AtomicUsize::new(0),
269            }
270        }
271
272        fn block_next_put(&self) {
273            self.block_next_put.store(true, AtomicOrdering::Relaxed);
274        }
275
276        async fn wait_for_blocked_put(&self) {
277            self.put_entered.acquire().await.unwrap().forget();
278        }
279
280        fn release_one_put(&self) {
281            self.put_release.add_permits(1);
282        }
283
284        fn set_delay_puts(&self, delay: bool) {
285            self.delay_puts.store(delay, AtomicOrdering::Relaxed);
286        }
287
288        fn clear_recorded_puts(&self) {
289            self.recorded_puts.lock().unwrap().clear();
290        }
291
292        fn recorded_puts(&self) -> Vec<PutRequest> {
293            self.recorded_puts.lock().unwrap().clone()
294        }
295
296        fn max_active_puts(&self) -> usize {
297            self.max_active_puts.load(AtomicOrdering::Relaxed)
298        }
299
300        fn start_put(&self, req: &PutRequest) -> PutGuard<'_> {
301            self.recorded_puts.lock().unwrap().push(req.clone());
302            let active = self.active_puts.fetch_add(1, AtomicOrdering::Relaxed) + 1;
303            self.max_active_puts
304                .fetch_max(active, AtomicOrdering::Relaxed);
305            PutGuard { backend: self }
306        }
307    }
308
309    struct PutGuard<'a> {
310        backend: &'a ControlledKvBackend,
311    }
312
313    impl Drop for PutGuard<'_> {
314        fn drop(&mut self) {
315            self.backend
316                .active_puts
317                .fetch_sub(1, AtomicOrdering::Relaxed);
318        }
319    }
320
321    #[async_trait::async_trait]
322    impl TxnService for ControlledKvBackend {
323        type Error = MetaError;
324    }
325
326    #[async_trait::async_trait]
327    impl KvBackend for ControlledKvBackend {
328        fn name(&self) -> &str {
329            "controlled"
330        }
331
332        fn as_any(&self) -> &dyn Any {
333            self
334        }
335
336        async fn range(&self, _req: RangeRequest) -> MetaResult<RangeResponse> {
337            unimplemented!()
338        }
339
340        async fn put(&self, req: PutRequest) -> MetaResult<PutResponse> {
341            let _guard = self.start_put(&req);
342            if self.block_next_put.swap(false, AtomicOrdering::Relaxed) {
343                self.put_entered.add_permits(1);
344                self.put_release.acquire().await.unwrap().forget();
345            }
346            if self.delay_puts.load(AtomicOrdering::Relaxed) {
347                sleep(Duration::from_millis(10)).await;
348            }
349            Ok(PutResponse::default())
350        }
351
352        async fn batch_put(&self, _req: BatchPutRequest) -> MetaResult<BatchPutResponse> {
353            unimplemented!()
354        }
355
356        async fn batch_get(&self, _req: BatchGetRequest) -> MetaResult<BatchGetResponse> {
357            unimplemented!()
358        }
359
360        async fn delete_range(&self, _req: DeleteRangeRequest) -> MetaResult<DeleteRangeResponse> {
361            unimplemented!()
362        }
363
364        async fn batch_delete(&self, _req: BatchDeleteRequest) -> MetaResult<BatchDeleteResponse> {
365            unimplemented!()
366        }
367    }
368
369    impl ResettableKvBackend for ControlledKvBackend {
370        fn reset(&self) {
371            self.clear_recorded_puts();
372        }
373
374        fn as_kv_backend_ref(self: Arc<Self>) -> KvBackendRef {
375            self
376        }
377    }
378
379    fn stat(node_id: u64, epoch: u64, marker: u64, addr: &str) -> Stat {
380        Stat {
381            timestamp_millis: marker as i64,
382            id: node_id,
383            addr: addr.to_string(),
384            region_num: marker,
385            node_epoch: epoch,
386            ..Default::default()
387        }
388    }
389
390    async fn handle_stat(
391        handler: Arc<CollectStatsHandler>,
392        mut ctx: Context,
393        stat: Stat,
394    ) -> Result<HandleControl> {
395        let mut acc = HeartbeatAccumulator {
396            stat: Some(stat),
397            ..Default::default()
398        };
399        handler
400            .handle(&HeartbeatRequest::default(), &mut ctx, &mut acc)
401            .await
402    }
403
404    fn use_controlled_address_backend(ctx: &mut Context) -> Arc<ControlledKvBackend> {
405        let backend = Arc::new(ControlledKvBackend::new());
406        ctx.leader_cached_kv_backend =
407            Arc::new(LeaderCachedKvBackend::with_always_leader(backend.clone()));
408        backend
409    }
410
411    #[tokio::test]
412    async fn test_handle_datanode_stats() {
413        let env = TestEnv::new();
414        let ctx = env.ctx();
415
416        let handler = CollectStatsHandler::default();
417        handle_request_many_times(ctx.clone(), &handler, 1).await;
418
419        let key = DatanodeStatKey { node_id: 101 };
420        let key: Vec<u8> = key.into();
421        let res = ctx.in_memory.get(&key).await.unwrap();
422        let kv = res.unwrap();
423        let key: DatanodeStatKey = kv.key.clone().try_into().unwrap();
424        assert_eq!(101, key.node_id);
425        let val: DatanodeStatValue = kv.value.try_into().unwrap();
426        // first new stat must be set in kv store immediately
427        assert_eq!(1, val.stats.len());
428        assert_eq!(1, val.stats[0].region_num);
429
430        handle_request_many_times(ctx.clone(), &handler, 10).await;
431
432        let key: Vec<u8> = key.into();
433        let res = ctx.in_memory.get(&key).await.unwrap();
434        let kv = res.unwrap();
435        let val: DatanodeStatValue = kv.value.try_into().unwrap();
436        assert_eq!(handler.flush_stats_factor, val.stats.len());
437    }
438
439    #[test]
440    fn test_same_datanode_wait_keeps_current_thread_runtime_responsive() {
441        let (backend_tx, backend_rx) = mpsc::sync_channel(1);
442        let (timer_tx, timer_rx) = mpsc::sync_channel(1);
443        let (done_tx, done_rx) = mpsc::sync_channel(1);
444
445        let worker = thread::spawn(move || {
446            tokio::runtime::Builder::new_current_thread()
447                .enable_time()
448                .build()
449                .unwrap()
450                .block_on(async move {
451                    let env = TestEnv::new();
452                    let mut ctx = env.ctx();
453                    let address_backend = use_controlled_address_backend(&mut ctx);
454                    address_backend.block_next_put();
455                    backend_tx.send(address_backend.clone()).unwrap();
456
457                    let handler = Arc::new(CollectStatsHandler::default());
458                    let first = tokio::spawn(handle_stat(
459                        handler.clone(),
460                        ctx.clone(),
461                        stat(101, 1, 1, "dn-101-v1"),
462                    ));
463                    address_backend.wait_for_blocked_put().await;
464
465                    let second =
466                        tokio::spawn(handle_stat(handler, ctx, stat(101, 1, 2, "dn-101-v1")));
467                    tokio::spawn(async move {
468                        sleep(Duration::from_millis(10)).await;
469                        timer_tx.send(()).unwrap();
470                    });
471                    tokio::task::yield_now().await;
472
473                    first.await.unwrap().unwrap();
474                    second.await.unwrap().unwrap();
475                    done_tx.send(()).unwrap();
476                });
477        });
478
479        let address_backend = backend_rx.recv_timeout(Duration::from_secs(1)).unwrap();
480        let timer_result = timer_rx.recv_timeout(Duration::from_secs(1));
481        address_backend.release_one_put();
482        timer_result.expect("the current-thread runtime must remain responsive");
483        done_rx
484            .recv_timeout(Duration::from_secs(1))
485            .expect("both heartbeat handlers must complete after releasing the write");
486        worker.join().unwrap();
487    }
488
489    #[tokio::test]
490    async fn test_concurrent_flush_persists_every_stat_once() {
491        let env = TestEnv::new();
492        let mut ctx = env.ctx();
493        let stats_backend = Arc::new(ControlledKvBackend::new());
494        ctx.in_memory = stats_backend.clone();
495        stats_backend.set_delay_puts(true);
496
497        let flush_stats_factor = 3;
498        let handler = Arc::new(CollectStatsHandler::new(Some(flush_stats_factor)));
499        handle_stat(handler.clone(), ctx.clone(), stat(101, 1, 0, "dn-101"))
500            .await
501            .unwrap();
502        stats_backend.clear_recorded_puts();
503
504        let mut tasks = Vec::with_capacity(2 * flush_stats_factor);
505        for marker in 1..=(2 * flush_stats_factor) {
506            tasks.push(tokio::spawn(handle_stat(
507                handler.clone(),
508                ctx.clone(),
509                stat(101, 1, marker as u64, "dn-101"),
510            )));
511        }
512        for task in tasks {
513            timeout(Duration::from_secs(1), task)
514                .await
515                .unwrap()
516                .unwrap()
517                .unwrap();
518        }
519
520        let puts = stats_backend.recorded_puts();
521        assert_eq!(2, puts.len());
522        let mut markers = puts
523            .into_iter()
524            .flat_map(|put| {
525                let value: DatanodeStatValue = put.value.try_into().unwrap();
526                value
527                    .stats
528                    .into_iter()
529                    .map(|stat| stat.region_num)
530                    .collect::<Vec<_>>()
531            })
532            .collect::<Vec<_>>();
533        markers.sort_unstable();
534        assert_eq!((1..=6).collect::<Vec<_>>(), markers);
535        assert_eq!(1, stats_backend.max_active_puts());
536    }
537
538    async fn handle_request_many_times(
539        mut ctx: Context,
540        handler: &CollectStatsHandler,
541        loop_times: i32,
542    ) {
543        let req = HeartbeatRequest::default();
544        for i in 1..=loop_times {
545            let mut acc = HeartbeatAccumulator {
546                stat: Some(Stat {
547                    id: 101,
548                    region_num: i as _,
549                    ..Default::default()
550                }),
551                ..Default::default()
552            };
553            handler.handle(&req, &mut ctx, &mut acc).await.unwrap();
554        }
555    }
556}