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