1#[cfg(test)]
16mod tests;
17
18use std::collections::{HashMap, HashSet};
19use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
20use std::sync::{Arc, Mutex as StdMutex};
21
22use api::v1::meta::heartbeat_request::NodeWorkloads;
23use api::v1::meta::{FrontendWorkloads, HeartbeatRequest, HeartbeatResponse, NodeInfo, Peer};
24use async_trait::async_trait;
25use common_error::ext::BoxedError;
26use common_meta::cache_invalidator::CacheInvalidatorRef;
27use common_meta::datanode::EnvVars;
28use common_meta::heartbeat::handler::invalidate_table_cache::InvalidateCacheHandler;
29use common_meta::heartbeat::handler::parse_mailbox_message::ParseMailboxMessageHandler;
30use common_meta::heartbeat::handler::suspend::SuspendHandler;
31use common_meta::heartbeat::handler::{
32 HandlerGroupExecutor, HeartbeatResponseHandlerContext, HeartbeatResponseHandlerExecutorRef,
33 HeartbeatResponseHandlerRef,
34};
35use common_meta::heartbeat::mailbox::{HeartbeatMailbox, MailboxRef, OutgoingMessage};
36use common_meta::heartbeat::utils::outgoing_message_to_mailbox_message;
37use common_stat::ResourceStatRef;
38use common_telemetry::{debug, error, info, warn};
39use meta_client::client::heartbeat::HeartbeatConfig;
40use meta_client::client::{HeartbeatSender, HeartbeatStream, MetaClient};
41use servers::addrs;
42use snafu::ResultExt;
43use tokio::sync::mpsc::Receiver;
44use tokio::sync::{Mutex, mpsc};
45use tokio::time::{Duration, Instant};
46use tokio_util::sync::CancellationToken;
47
48use crate::error;
49use crate::error::Result;
50use crate::frontend::FrontendOptions;
51use crate::metrics::{HEARTBEAT_RECV_COUNT, HEARTBEAT_SENT_COUNT};
52
53pub type FrontendHeartbeatExtensionResult<T> = std::result::Result<T, BoxedError>;
55
56#[async_trait]
65pub trait FrontendHeartbeatExtension: Send + Sync {
66 fn name(&self) -> &str;
68
69 async fn request_extensions(
71 &self,
72 ) -> FrontendHeartbeatExtensionResult<HashMap<String, Vec<u8>>> {
73 Ok(HashMap::new())
74 }
75
76 fn response_handler(&self) -> Option<HeartbeatResponseHandlerRef> {
81 None
82 }
83
84 async fn connected(&self, _generation: u64) -> FrontendHeartbeatExtensionResult<()> {
86 Ok(())
87 }
88
89 async fn shutdown(&self) -> FrontendHeartbeatExtensionResult<()> {
91 Ok(())
92 }
93}
94
95#[derive(Clone, Default)]
97pub struct FrontendHeartbeatExtensions {
98 inner: Arc<StdMutex<FrontendHeartbeatExtensionsInner>>,
99}
100
101#[derive(Default)]
102struct FrontendHeartbeatExtensionsInner {
103 names: HashSet<String>,
104 extensions: Vec<Arc<dyn FrontendHeartbeatExtension>>,
105}
106
107impl FrontendHeartbeatExtensions {
108 pub fn register(&self, extension: Arc<dyn FrontendHeartbeatExtension>) -> bool {
112 let mut inner = self.inner.lock().unwrap();
113 let name = extension.name().to_string();
114 if !inner.names.insert(name) {
115 return false;
116 }
117 inner.extensions.push(extension);
118 true
119 }
120
121 pub fn extensions(&self) -> Vec<Arc<dyn FrontendHeartbeatExtension>> {
123 self.inner.lock().unwrap().extensions.clone()
124 }
125
126 pub fn response_handlers(&self) -> Vec<HeartbeatResponseHandlerRef> {
128 self.extensions()
129 .into_iter()
130 .filter_map(|extension| extension.response_handler())
131 .collect()
132 }
133
134 pub fn len(&self) -> usize {
136 self.inner.lock().unwrap().extensions.len()
137 }
138
139 pub fn is_empty(&self) -> bool {
141 self.len() == 0
142 }
143}
144
145pub fn heartbeat_response_handler_executor(
150 extensions: &FrontendHeartbeatExtensions,
151 suspend_state: Arc<AtomicBool>,
152 cache_invalidator: CacheInvalidatorRef,
153) -> HeartbeatResponseHandlerExecutorRef {
154 let mut handlers: Vec<HeartbeatResponseHandlerRef> = vec![Arc::new(ParseMailboxMessageHandler)];
155 handlers.extend(extensions.response_handlers().into_iter().map(|handler| {
156 Arc::new(IsolatedHeartbeatResponseHandler(handler)) as HeartbeatResponseHandlerRef
157 }));
158 handlers.extend([
159 Arc::new(SuspendHandler::new(suspend_state)) as HeartbeatResponseHandlerRef,
160 Arc::new(InvalidateCacheHandler::new(cache_invalidator)),
161 ]);
162 Arc::new(HandlerGroupExecutor::new(handlers))
163}
164
165struct IsolatedHeartbeatResponseHandler(HeartbeatResponseHandlerRef);
166
167#[async_trait]
168impl common_meta::heartbeat::handler::HeartbeatResponseHandler
169 for IsolatedHeartbeatResponseHandler
170{
171 fn is_acceptable(&self, _ctx: &HeartbeatResponseHandlerContext) -> bool {
172 true
173 }
174
175 async fn handle(
176 &self,
177 ctx: &mut HeartbeatResponseHandlerContext,
178 ) -> common_meta::error::Result<common_meta::heartbeat::handler::HandleControl> {
179 use common_meta::heartbeat::handler::HandleControl;
180
181 if self.0.is_acceptable(ctx)
182 && let Err(error) = self.0.handle(ctx).await
183 {
184 error!(error; "Heartbeat extension response handler failed");
185 }
186 Ok(HandleControl::Continue)
187 }
188}
189
190#[async_trait]
191trait HeartbeatConnector: Send + Sync {
192 async fn connect(&self) -> Result<HeartbeatConnection>;
193}
194
195#[async_trait]
196trait HeartbeatRequestSender: Send + Sync {
197 async fn send(&self, request: HeartbeatRequest) -> FrontendHeartbeatExtensionResult<()>;
198}
199
200#[async_trait]
201trait HeartbeatResponseStream: Send {
202 async fn message(&mut self) -> FrontendHeartbeatExtensionResult<Option<HeartbeatResponse>>;
203}
204
205struct HeartbeatConnection {
206 sender: Arc<dyn HeartbeatRequestSender>,
207 stream: Box<dyn HeartbeatResponseStream>,
208 config: HeartbeatConfig,
209}
210
211struct MetaHeartbeatConnector {
212 client: Arc<MetaClient>,
213}
214
215#[async_trait]
216impl HeartbeatConnector for MetaHeartbeatConnector {
217 async fn connect(&self) -> Result<HeartbeatConnection> {
218 let (sender, stream, config) = self
219 .client
220 .heartbeat()
221 .await
222 .context(error::CreateMetaHeartbeatStreamSnafu)?;
223 Ok(HeartbeatConnection {
224 sender: Arc::new(MetaHeartbeatSender(sender)),
225 stream: Box::new(MetaHeartbeatStream(stream)),
226 config,
227 })
228 }
229}
230
231struct MetaHeartbeatSender(HeartbeatSender);
232
233#[async_trait]
234impl HeartbeatRequestSender for MetaHeartbeatSender {
235 async fn send(&self, request: HeartbeatRequest) -> FrontendHeartbeatExtensionResult<()> {
236 self.0.send(request).await.map_err(BoxedError::new)
237 }
238}
239
240struct MetaHeartbeatStream(HeartbeatStream);
241
242#[async_trait]
243impl HeartbeatResponseStream for MetaHeartbeatStream {
244 async fn message(&mut self) -> FrontendHeartbeatExtensionResult<Option<HeartbeatResponse>> {
245 self.0.message().await.map_err(BoxedError::new)
246 }
247}
248
249#[derive(Clone)]
250struct HeartbeatRunner {
251 peer_addr: String,
252 connector: Arc<dyn HeartbeatConnector>,
253 resp_handler_executor: HeartbeatResponseHandlerExecutorRef,
254 start_time_ms: u64,
255 resource_stat: ResourceStatRef,
256 env_vars: EnvVars,
257 extensions: FrontendHeartbeatExtensions,
258 cancellation: CancellationToken,
259 generation: Arc<AtomicU64>,
260}
261
262impl HeartbeatRunner {
263 async fn connect(&self) -> Result<Option<HeartbeatConnection>> {
264 tokio::select! {
265 _ = self.cancellation.cancelled() => Ok(None),
266 connection = self.connector.connect() => connection.map(Some),
267 }
268 }
269
270 fn next_generation(&self) -> u64 {
271 self.generation
272 .fetch_add(1, Ordering::AcqRel)
273 .wrapping_add(1)
274 }
275
276 async fn notify_connected(&self, generation: u64) -> bool {
277 for extension in self.extensions.extensions() {
278 let result = tokio::select! {
279 _ = self.cancellation.cancelled() => return false,
280 result = extension.connected(generation) => result,
281 };
282 if let Err(error) = result {
283 error!(error; "Heartbeat extension '{}' failed its connected callback", extension.name());
284 }
285 }
286 true
287 }
288
289 async fn shutdown_extensions(&self) {
290 for extension in self.extensions.extensions() {
291 if let Err(error) = extension.shutdown().await {
292 error!(error; "Failed to shut down heartbeat extension '{}'", extension.name());
293 }
294 }
295 }
296
297 async fn run(self, mut connection: HeartbeatConnection) {
298 loop {
299 let retry_interval = connection.config.retry_interval;
300 if self.run_connection(connection).await == ConnectionEnd::Shutdown {
301 return;
302 }
303
304 loop {
305 if !self.wait_retry(retry_interval).await {
306 return;
307 }
308 info!("Try to re-establish the heartbeat connection to metasrv.");
309
310 match self.connect().await {
311 Ok(Some(next)) => {
312 let generation = self.next_generation();
313 if !self.notify_connected(generation).await {
314 return;
315 }
316 connection = next;
317 break;
318 }
319 Ok(None) => return,
320 Err(error) => {
321 error!(error; "Failed to re-establish heartbeat connection to metasrv");
322 }
323 }
324 }
325 }
326 }
327
328 async fn wait_retry(&self, retry_interval: Duration) -> bool {
329 tokio::select! {
330 _ = self.cancellation.cancelled() => false,
331 _ = tokio::time::sleep(retry_interval) => true,
332 }
333 }
334
335 async fn run_connection(&self, connection: HeartbeatConnection) -> ConnectionEnd {
336 let (outgoing_tx, outgoing_rx) = mpsc::channel(16);
337 let mailbox = Arc::new(HeartbeatMailbox::new(outgoing_tx));
338 let report =
339 self.report_heartbeats(connection.sender, outgoing_rx, connection.config.interval);
340 let responses = self.handle_responses(connection.stream, mailbox);
341 tokio::pin!(report);
342 tokio::pin!(responses);
343
344 tokio::select! {
345 _ = self.cancellation.cancelled() => ConnectionEnd::Shutdown,
346 end = &mut report => end,
347 end = &mut responses => end,
348 }
349 }
350
351 async fn handle_responses(
352 &self,
353 mut stream: Box<dyn HeartbeatResponseStream>,
354 mailbox: MailboxRef,
355 ) -> ConnectionEnd {
356 loop {
357 let response = tokio::select! {
358 _ = self.cancellation.cancelled() => return ConnectionEnd::Shutdown,
359 response = stream.message() => response,
360 };
361 match response {
362 Ok(Some(response)) => {
363 debug!("Receiving heartbeat response: {:?}", response);
364 if let Some(message) = &response.mailbox_message {
365 info!("Received mailbox message: {message:?}");
366 }
367 let context = HeartbeatResponseHandlerContext::new(mailbox.clone(), response);
368 let result = tokio::select! {
369 _ = self.cancellation.cancelled() => return ConnectionEnd::Shutdown,
370 result = self.handle_response(context) => result,
371 };
372 if let Err(error) = result {
373 error!(error; "Error while handling heartbeat response");
374 HEARTBEAT_RECV_COUNT
375 .with_label_values(&["processing_error"])
376 .inc();
377 } else {
378 HEARTBEAT_RECV_COUNT.with_label_values(&["success"]).inc();
379 }
380 }
381 Ok(None) => {
382 warn!("Heartbeat response stream closed");
383 return ConnectionEnd::Reconnect;
384 }
385 Err(error) => {
386 HEARTBEAT_RECV_COUNT.with_label_values(&["error"]).inc();
387 error!(error; "Occur error while reading heartbeat response");
388 return ConnectionEnd::Reconnect;
389 }
390 }
391 }
392 }
393
394 async fn report_heartbeats(
395 &self,
396 sender: Arc<dyn HeartbeatRequestSender>,
397 mut outgoing_rx: Receiver<OutgoingMessage>,
398 report_interval: Duration,
399 ) -> ConnectionEnd {
400 let total_cpu_millicores = self.resource_stat.get_total_cpu_millicores();
401 let total_memory_bytes = self.resource_stat.get_total_memory_bytes();
402 let mut extensions = HashMap::new();
403 self.env_vars.into_extensions(&mut extensions);
404 let heartbeat_request = HeartbeatRequest {
405 peer: Some(Peer {
406 id: 0,
408 addr: self.peer_addr.clone(),
409 }),
410 info: Self::build_node_info(
411 self.start_time_ms,
412 total_cpu_millicores,
413 total_memory_bytes,
414 ),
415 node_workloads: Some(NodeWorkloads::Frontend(FrontendWorkloads { types: vec![] })),
416 extensions,
417 ..Default::default()
418 };
419 let sleep = tokio::time::sleep(Duration::ZERO);
420 tokio::pin!(sleep);
421
422 loop {
423 let request = tokio::select! {
424 _ = self.cancellation.cancelled() => return ConnectionEnd::Shutdown,
425 message = outgoing_rx.recv() => {
426 if let Some(message) = message {
427 Self::new_heartbeat_request(&heartbeat_request, Some(message), 0, 0)
428 } else {
429 warn!("Sender has been dropped, exiting the heartbeat loop");
430 return ConnectionEnd::Reconnect;
431 }
432 }
433 _ = &mut sleep => {
434 sleep.as_mut().reset(Instant::now() + report_interval);
435 Self::new_heartbeat_request(
436 &heartbeat_request,
437 None,
438 self.resource_stat.get_cpu_usage_millicores(),
439 self.resource_stat.get_memory_usage_bytes(),
440 )
441 }
442 };
443
444 if let Some(mut request) = request {
445 if !self.add_request_extensions(&mut request).await {
446 return ConnectionEnd::Shutdown;
447 }
448 debug!(
449 "Sending a heartbeat request to metasrv, content: {:?}",
450 request
451 );
452 let result = tokio::select! {
453 _ = self.cancellation.cancelled() => return ConnectionEnd::Shutdown,
454 result = sender.send(request) => result,
455 };
456 if let Err(error) = result {
457 error!(error; "Failed to send heartbeat to metasrv");
458 return ConnectionEnd::Reconnect;
459 }
460 HEARTBEAT_SENT_COUNT.inc();
461 }
462 }
463 }
464
465 async fn add_request_extensions(&self, request: &mut HeartbeatRequest) -> bool {
466 for extension in self.extensions.extensions() {
467 let generated = tokio::select! {
468 _ = self.cancellation.cancelled() => return false,
469 generated = extension.request_extensions() => generated,
470 };
471 let generated = match generated {
472 Ok(generated) => generated,
473 Err(error) => {
474 error!(error; "Heartbeat extension '{}' failed to generate request extensions", extension.name());
475 continue;
476 }
477 };
478
479 if let Some(key) = generated
480 .keys()
481 .find(|key| request.extensions.contains_key(*key))
482 {
483 warn!(
484 "Heartbeat extension '{}' produced conflicting key '{}'; discarding its output",
485 extension.name(),
486 key
487 );
488 continue;
489 }
490 request.extensions.extend(generated);
491 }
492 true
493 }
494
495 fn new_heartbeat_request(
496 heartbeat_request: &HeartbeatRequest,
497 message: Option<OutgoingMessage>,
498 cpu_usage: i64,
499 memory_usage: i64,
500 ) -> Option<HeartbeatRequest> {
501 let mailbox_message = match message.map(outgoing_message_to_mailbox_message) {
502 Some(Ok(message)) => Some(message),
503 Some(Err(error)) => {
504 error!(error; "Failed to encode mailbox messages");
505 return None;
506 }
507 None => None,
508 };
509
510 let mut heartbeat_request = HeartbeatRequest {
511 mailbox_message,
512 ..heartbeat_request.clone()
513 };
514 if let Some(info) = heartbeat_request.info.as_mut() {
515 info.memory_usage_bytes = memory_usage;
516 info.cpu_usage_millicores = cpu_usage;
517 }
518 Some(heartbeat_request)
519 }
520
521 #[allow(deprecated)]
522 fn build_node_info(
523 start_time_ms: u64,
524 total_cpu_millicores: i64,
525 total_memory_bytes: i64,
526 ) -> Option<NodeInfo> {
527 let build_info = common_version::build_info();
528 Some(NodeInfo {
529 version: build_info.version.to_string(),
530 git_commit: build_info.commit_short.to_string(),
531 start_time_ms,
532 total_cpu_millicores,
533 total_memory_bytes,
534 cpu_usage_millicores: 0,
535 memory_usage_bytes: 0,
536 cpus: total_cpu_millicores as u32,
538 memory_bytes: total_memory_bytes as u64,
539 hostname: hostname::get()
540 .unwrap_or_default()
541 .to_string_lossy()
542 .to_string(),
543 })
544 }
545
546 async fn handle_response(&self, context: HeartbeatResponseHandlerContext) -> Result<()> {
547 self.resp_handler_executor
548 .handle(context)
549 .await
550 .context(error::HandleHeartbeatResponseSnafu)
551 }
552}
553
554#[derive(Debug, PartialEq, Eq)]
555enum ConnectionEnd {
556 Reconnect,
557 Shutdown,
558}
559
560#[derive(Clone)]
562pub struct HeartbeatTask {
563 runner: HeartbeatRunner,
564 start_lock: Arc<Mutex<()>>,
565 shutdown_lock: Arc<Mutex<()>>,
566 supervisor: Arc<Mutex<Option<common_runtime::JoinHandle<()>>>>,
567}
568
569impl HeartbeatTask {
570 pub fn new(
571 peer_addr: String,
572 opts: &FrontendOptions,
573 meta_client: Arc<MetaClient>,
574 resp_handler_executor: HeartbeatResponseHandlerExecutorRef,
575 resource_stat: ResourceStatRef,
576 ) -> Self {
577 Self::new_with_connector(
578 peer_addr,
579 opts,
580 Arc::new(MetaHeartbeatConnector {
581 client: meta_client,
582 }),
583 resp_handler_executor,
584 resource_stat,
585 )
586 }
587
588 fn new_with_connector(
589 peer_addr: String,
590 opts: &FrontendOptions,
591 connector: Arc<dyn HeartbeatConnector>,
592 resp_handler_executor: HeartbeatResponseHandlerExecutorRef,
593 resource_stat: ResourceStatRef,
594 ) -> Self {
595 Self {
596 runner: HeartbeatRunner {
597 peer_addr,
598 connector,
599 resp_handler_executor,
600 start_time_ms: common_time::util::current_time_millis() as u64,
601 resource_stat,
602 env_vars: EnvVars::from_config(&opts.heartbeat_env_vars),
603 extensions: FrontendHeartbeatExtensions::default(),
604 cancellation: CancellationToken::new(),
605 generation: Arc::new(AtomicU64::new(0)),
606 },
607 start_lock: Arc::new(Mutex::new(())),
608 shutdown_lock: Arc::new(Mutex::new(())),
609 supervisor: Arc::new(Mutex::new(None)),
610 }
611 }
612
613 pub fn with_extensions(mut self, extensions: FrontendHeartbeatExtensions) -> Self {
615 self.runner.extensions = extensions;
616 self
617 }
618
619 pub async fn start(&self) -> Result<()> {
621 let _start_guard = self.start_lock.lock().await;
622 if self.runner.cancellation.is_cancelled() {
623 return Ok(());
624 }
625
626 let finished = {
627 let mut supervisor = self.supervisor.lock().await;
628 match supervisor.as_ref() {
629 Some(handle) if !handle.is_finished() => return Ok(()),
630 Some(_) => supervisor.take(),
631 None => None,
632 }
633 };
634 if let Some(handle) = finished
635 && let Err(error) = handle.await
636 && !error.is_cancelled()
637 {
638 error!(error; "Heartbeat supervisor join failed");
639 }
640
641 let Some(connection) = self.runner.connect().await? else {
642 return Ok(());
643 };
644 info!(
645 "Heartbeat started with Metasrv config: {}",
646 connection.config
647 );
648
649 let generation = self.runner.next_generation();
650 if !self.runner.notify_connected(generation).await {
651 return Ok(());
652 }
653
654 let runner = self.runner.clone();
655 let handle = common_runtime::spawn_hb(async move {
656 runner.run(connection).await;
657 });
658 *self.supervisor.lock().await = Some(handle);
659 Ok(())
660 }
661
662 pub async fn shutdown(&self) {
664 let _shutdown_guard = self.shutdown_lock.lock().await;
665 if self.runner.cancellation.is_cancelled() {
666 return;
667 }
668 self.runner.cancellation.cancel();
669
670 let _start_guard = self.start_lock.lock().await;
672 let handle = self.supervisor.lock().await.take();
673 if let Some(handle) = handle
674 && let Err(error) = handle.await
675 && !error.is_cancelled()
676 {
677 error!(error; "Heartbeat supervisor join failed");
678 }
679 self.runner.shutdown_extensions().await;
680 }
681
682 #[cfg(test)]
683 fn generation(&self) -> u64 {
684 self.runner.generation.load(Ordering::Acquire)
685 }
686
687 #[cfg(test)]
688 async fn has_supervisor(&self) -> bool {
689 self.supervisor.lock().await.is_some()
690 }
691
692 #[cfg(test)]
693 pub(crate) fn is_shutdown(&self) -> bool {
694 self.runner.cancellation.is_cancelled()
695 }
696}
697
698pub(crate) fn frontend_peer_addr(opts: &FrontendOptions) -> String {
699 if let Some(internal) = &opts.internal_grpc {
703 addrs::resolve_addr(&internal.bind_addr, Some(&internal.server_addr))
704 } else {
705 addrs::resolve_addr(&opts.grpc.bind_addr, Some(&opts.grpc.server_addr))
706 }
707}