1use 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 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 true
147 };
148
149 if refresh {
151 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 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 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}